The shop's producer combines idempotence, Java-compatible partitioning, zstd 126 , a short linger, the event type in a header (so consumers can filter without parsing) and a callback that records failures. Records keep the send time as their timestamp: stamping the 2025 event time would let 7-day retention delete them (Retention and Compaction).
"""booknest_producer.py: publish BookNest's order events the way the shop backend does."""
import json
from confluent_kafka import Producer
producer = Producer({
"bootstrap.servers": "localhost:33092", "client.id": "booknest-shop",
"enable.idempotence": True, # acks=all, no duplicates
"partitioner": "murmur2_random", # same partition as Java
"compression.type": "zstd", "linger.ms": 20, "queue.buffering.max.messages": 500_000})
failed, sent = [], 0
for line in open("data/order_events.jsonl", encoding="utf-8"):
event = json.loads(line)
producer.produce("booknest.order-events", line.rstrip("\n"), str(event["order_id"]),
headers={"event-type": event["type"]},
on_delivery=lambda err, msg: err and failed.append((msg.key(), err)))
producer.poll(0) # serve delivery reports as we go
sent += 1
print(f"{sent:,} sent, {len(failed)} failed, {producer.flush(60)} undelivered")Output
390,737 sent, 0 failed, 0 undelivered
On a freshly recreated booknest.order-events it filled the partitions with 130,491, 130,420 and 129,826 events in seconds; Transactions and Exactly-Once and Consumers and Queues read them.