A transactional producer has a stable transactional.id that you choose. The broker that leads the partition of the internal __transaction_state topic for that ID's hash is its transaction coordinator; it maps the ID to a producer ID and an epoch, and logs each transaction's state. When the producer commits, the coordinator writes a COMMIT or ABORT control record (marker) into every partition the transaction touched. Since transaction version 2 (KIP-890, active on this cluster, Kafka 4.x and KRaft), the epoch also rises with every transaction, which lets brokers reject stray writes from earlier ones. kafka-transactions.sh shows a transaction in flight.
import subprocess
from confluent_kafka import Producer
p = Producer({"bootstrap.servers": "localhost:33092", "transactional.id": "booknest-shop-tx"})
p.init_transactions() # find the coordinator, get a producer ID and epoch
p.begin_transaction()
p.produce("booknest.order-events", '{"type":"order_paid","order_id":7}', "7")
p.flush() # sent, but not yet committed
out = subprocess.run(["docker", "exec", "l3-kafka", "/opt/kafka/bin/kafka-transactions.sh",
"--bootstrap-server", "l3-kafka:9092", "describe", "--transactional-id",
"booknest-shop-tx"], capture_output=True, text=True)
for name, value in zip(*(line.split() for line in out.stdout.splitlines())):
print(f"{name:<30} {value}")
p.commit_transaction()CoordinatorId 1 TransactionalId booknest-shop-tx ProducerId 6 ProducerEpoch 4 TransactionState Ongoing TransactionTimeoutMs 60000 CurrentTransactionStartTimeMs 1790911611480 TransactionDurationMs 4300 TopicPartitions booknest.order-events-0
The epoch is 4 because earlier runs had used this ID: every new instance, and with transaction version 2 every finished transaction, raises it. A transaction left open longer than transaction.timeout.ms (60 s) is aborted by the coordinator. kafka-transactions.sh can also list transactions, find hanging ones and abort them by hand.