Exactly-Once in Practice

Exactly-Once BookNest Order Events in Practice

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: a consume-transform-produce job with transactionsPython
"""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.

Output of 36
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.