BookNest's refunds service must never pay twice, so a job copies every order_cancelled and order_returned event from the 390,737 order events to booknest.refunds, one transaction per batch of 1,000. Given an argument, it crashes inside that batch after its records reached the broker, the worst moment.
"""refunds_eos.py [CRASH_AT]: copy cancellations and returns to booknest.refunds once."""
import json, os, sys
from confluent_kafka import Consumer, Producer
B, crash_at = "localhost:33092", int(sys.argv[1]) if len(sys.argv) > 1 else 0
consumer = Consumer({"bootstrap.servers": B, "group.id": "booknest-refunds",
"enable.auto.commit": False, "auto.offset.reset": "earliest",
"isolation.level": "read_committed", "session.timeout.ms": 6000})
producer = Producer({"bootstrap.servers": B, "transactional.id": "booknest-refunds-0"})
producer.init_transactions() # fences and aborts a crashed predecessor
consumer.subscribe(["booknest.order-events"])
batches = copied = idle = 0
while idle < 3: # stop after three empty polls once assigned
msgs = consumer.consume(1000, timeout=2)
idle = 0 if msgs else idle + bool(consumer.assignment())
if not msgs:
continue
producer.begin_transaction()
for m in msgs:
event = json.loads(m.value())
if event["type"] in ("order_cancelled", "order_returned"):
producer.produce("booknest.refunds", m.value(), m.key()); copied += 1
producer.send_offsets_to_transaction( # the input position, in the same transaction
consumer.position(consumer.assignment()), consumer.consumer_group_metadata())
batches += 1
if batches == crash_at:
producer.flush() # this batch's refunds reach the broker...
print(f"crash in batch {batches} after copying {copied}", flush=True)
os._exit(1) # ...but the transaction is never committed
producer.commit_transaction()
print(f"{batches} transactions, {copied} refunds copied")listings/l0665_run.sh runs it with a crash in batch 50, waits for the dead member's 6-second session to expire, runs it again to the end, and counts the records and distinct event_ids in booknest.refunds with both isolation levels.
crash in batch 50 after copying 1298 342 transactions, 8550 refunds copied read_uncommitted 9848 records, 9819 distinct events read_committed 9819 records, 9819 distinct events
The order events hold 5,890 cancellations and 3,929 returns: 9,819. The 29 refunds of the crashed batch are in the log but aborted; the restart began at batch 50's first offset and copied them again. A read_committed reader sees each refund exactly once; a read_uncommitted one would pay 29 twice.
Passing consumer_group_metadata() with the offsets (KIP-447, Kafka 2.5 129 ) lets the coordinator fence a zombie by its group generation, so one transactional ID per job instance is enough; older guides demand one per input partition. Each commit costs a few round trips, so batch sizes trade latency for throughput: here the 342 transactions of the restart, with the 8-second wait, took under a minute on the shared host.