The API has five calls: init_transactions() once at startup, then begin_transaction(), produce(), send_offsets_to_transaction() when consuming as well, and commit_transaction() or abort_transaction(). init_transactions() also fences zombies: it bumps the epoch, aborts any open transaction of an older instance with the same ID, and from then on the coordinator rejects that instance.
from confluent_kafka import KafkaException, Producer
conf = {"bootstrap.servers": "localhost:33092", "transactional.id": "booknest-shop-tx"}
old = Producer(conf)
old.init_transactions()
old.begin_transaction()
old.produce("booknest.order-events", '{"type":"order_paid","order_id":8}', "8")
new = Producer(conf) # a restarted instance with the same ID
new.init_transactions() # bumps the epoch, aborts the old transaction
try:
old.commit_transaction()
except KafkaException as e:
err = e.args[0]
print(f"old instance: {err.name()} (fatal={err.fatal()})")
new.begin_transaction()
new.produce("booknest.order-events", '{"type":"order_paid","order_id":8}', "8")
new.commit_transaction()
print("new instance: committed")Output
old instance: _FENCED (fatal=True) new instance: committed
A fenced producer is finished: close it and exit, since its work belongs to the new instance. Give each instance of a job its own stable ID (booknest-refunds-0, -1, ...), never a random one per start, or no fencing happens. Retriable errors inside a transaction call for abort_transaction() and a retry of the batch.