Skip to content

AxonOps — AI-Native Control Plane for Open Source Data Platforms

Kafka Transaction Coordinator

The transaction coordinator enables atomic writes across multiple partitions and exactly-once semantics in Kafka. It manages producer identities, transaction state, and coordinates the two-phase commit protocol.


Kafka was originally designed for high-throughput streaming with at-least-once delivery. However, certain use cases require stronger guarantees:

ProblemWithout TransactionsWith Transactions
Duplicate writes on retryProducer retries may create duplicatesIdempotent writes prevent duplicates
Partial failuresWriting to 3 topics may succeed for 2, fail for 1All-or-nothing atomicity
Read-process-write consistencyConsumer offset committed separately from outputOffset and output committed atomically
Zombie producersCrashed producer restarts, old instance still writingEpoch fencing blocks old instances

Consider a stream processing application that reads from an input topic, processes records, and writes to an output topic:

Without TransactionsWith Transactions1. Read from input2. Process3. Write to output4. Commit offsetIf crash here:Output writtenOffset NOT committed→ Duplicate processing on restart1. Read from input2. Process3. Write to output+ commit offsetATOMICOutput + offset committed atomicallyNo duplicates possible

Use CaseWhy Transactions Help
Financial data processingCannot tolerate duplicates or partial updates
Exactly-once stream processingKafka Streams with exactly_once_v2 uses transactions internally
Multi-topic atomic writesOrder service writes to orders, inventory, notifications atomically
Read-process-write pipelinesETL jobs that must not reprocess or skip records
Event sourcing with projectionsEvent and projection updates must be atomic
Cross-partition aggregationsAggregation results written atomically with source offsets

Example 1: Order Processing

An order service must update multiple topics atomically:

Transaction {
write("orders", order_created_event)
write("inventory", reserve_stock_event)
write("payments", payment_request_event)
write("notifications", order_confirmation_event)
}

If any write fails, all are rolled back. No partial orders.

Example 2: Stream Aggregation

A real-time dashboard aggregates sales by region:

Transaction {
// Read sales events from input
// Aggregate by region
write("sales-by-region", aggregated_results)
commit_offsets("sales-events", consumer_group)
}

If the application crashes after writing results but before committing offsets, restart would not reprocess—the transaction ensures both happen or neither.

Example 3: Change Data Capture (CDC)

Database changes replicated to Kafka must maintain consistency:

Transaction {
write("users", user_updated)
write("user-search-index", search_document)
write("user-cache-invalidation", cache_key)
}

All downstream systems see the change atomically or not at all.


Transactions add overhead and complexity. They are not always the right choice.

ScenarioWhy Transactions Are WrongBetter Approach
High-throughput loggingOverhead can reduce throughputIdempotent producer (no transactions)
Fire-and-forget metricsOccasional duplicates acceptableacks=1 or acks=0
Single-partition writesNo cross-partition atomicity neededIdempotent producer only
Latency-critical pathsTransaction commit adds latencyAt-least-once with idempotent consumers
Independent eventsEvents don’t need atomic groupingSeparate non-transactional writes
MetricNon-TransactionalTransactionalImpact
Throughput100% baseline50-70%30-50% reduction
Latency (p50)~5ms~15-25ms3-5x increase
Latency (p99)~20ms~50-100ms2.5-5x increase
Broker CPUBaselineHigherCoordinator overhead
Do you need atomic writesacross multiple partitions/topics?YesyesnoDo you need exactly-onceread-process-write?YesyesnoUse TransactionsCan you tolerateoccasional duplicates?NonoyesUse TransactionsUse Idempotent ProducerDo you need deduplicationon producer retries?YesyesnoUse Idempotent ProducerStandard Producer(acks=all for durability)
AlternativeWhen to UseTrade-off
Idempotent producerSingle-partition deduplicationNo cross-partition atomicity
Idempotent consumerConsumer handles duplicatesApplication complexity
Outbox patternDatabase + Kafka consistencyRequires CDC or polling
Saga patternLong-running distributed workflowsCompensation logic needed

Kafka transactions provide atomicity—a set of writes either all succeed or all fail. The transaction coordinator is a broker-side component that manages transaction state and coordinates commits.

Transactional ProducerKafka ClusterTransactionCoordinatorPartition LeadersGroupCoordinatortransactional.idProducer ID (PID)Epoch__transaction_statetopic-A-0topic-A-1topic-B-0__consumer_offsetsInitProducerIdAddPartitionsToTxnEndTxnProduce (transactional)TxnOffsetCommitpersist stateWriteTxnMarkersWriteTxnMarkers
ComponentResponsibility
Transaction CoordinatorManages transaction state and coordinates commit/abort
__transaction_stateInternal topic storing transaction metadata
Producer ID (PID)Unique identifier for producer instance
EpochFencing mechanism to prevent zombie producers
Transaction MarkersControl records marking transaction boundaries

Every producer is assigned a unique 64-bit Producer ID. For idempotent producers, this enables deduplication. For transactional producers, it identifies the transaction owner.

ProducerTransactionProducerProducerTransactionCoordinatorTransactionCoordinatorNon-Transactional (Idempotent)InitProducerId(transactional_id=null)generate new PIDPID=1000, epoch=0PID generated each timeNo persistenceTransactionalInitProducerId(transactional_id="order-service-1")lookup or create PIDfor transactional.idPID=2000, epoch=0PID persisted in__transaction_state

The epoch is a 16-bit counter that increments each time a producer with the same transactional.id initializes. This enables “zombie fencing”—preventing old producers from causing inconsistencies.

Producer ATransactionProducer A.PartitionProducer A(epoch=0)Producer A(epoch=0)TransactionCoordinatorTransactionCoordinatorProducer A'(epoch=1)Producer A'(epoch=1)PartitionLeaderPartitionLeaderOriginal producerInitProducerId("order-svc")PID=100, epoch=0Produce(PID=100, epoch=0)ackCrash / Network partitionInitProducerId("order-svc")increment epochPID=100, epoch=1New producer instancewith same transactional.idProduce(PID=100, epoch=0)ProducerFencedExceptionOld producer is fencedmust shut down
ScenarioBehavior
Producer restartNew epoch assigned; old instance fenced
Network partitionFirst to re-initialize gets new epoch
Horizontal scalingEach instance needs unique transactional.id
Zombie producerRejected with ProducerFencedException

Each transactional.id maps to a specific coordinator via hashing:

coordinator = hash(transactional.id) % num_partitions(__transaction_state)

The broker hosting that partition of __transaction_state is the coordinator.

__transaction_state topic (50 partitions)Partition 0(Broker 1)Partition 1(Broker 2)Partition 2(Broker 3)...Partition 49(Broker 1)transactional.id = 'order-svc-1'transactional.id = 'payment-svc-1'Coordinator is the leader of theassigned __transaction_state partitionhash % 50 = 1hash % 50 = 2

The __transaction_state topic stores:

FieldDescription
transactional.idProducer’s transaction identifier
producer_idAssigned PID
producer_epochCurrent epoch
transaction_stateCurrent state (Empty, Ongoing, etc.)
topic_partitionsPartitions participating in transaction
transaction_timeout_msTimeout for this transaction
transaction_start_timeWhen transaction began
ConfigurationDefaultDescription
transaction.state.log.replication.factor3Replication factor for __transaction_state
transaction.state.log.num.partitions50Partitions in __transaction_state
transaction.state.log.min.isr2Minimum ISR for transaction log
transaction.state.log.segment.bytes104857600Segment size

EmptyNo active transactionOngoingTransaction in progressPrepareCommitCommit initiatedPrepareAbortAbort initiatedCompleteCommitCommit markers writtenCompleteAbortAbort markers writtenDeadtransactional.id expiredInitProducerIdAddPartitionsToTxnAddPartitionsToTxnProduceAddOffsetsToTxnEndTxn(COMMIT)EndTxn(ABORT)timeoutmarkers writtenmarkers writtencleanupcleanuptransactional.id.expiration.ms
StateDescription
EmptyProducer initialized; no active transaction
OngoingTransaction active; partitions being written
PrepareCommitCommit requested; writing markers
PrepareAbortAbort requested; writing markers
CompleteCommitAll commit markers written
CompleteAbortAll abort markers written
Deadtransactional.id expired due to inactivity

Kafka transactions use a variant of two-phase commit to ensure atomicity across partitions.

ProducerTransactiontopic-A Leadertopic-B LeaderProducerProducerTransactionCoordinatorTransactionCoordinatortopic-A Leadertopic-A Leadertopic-B Leadertopic-B LeaderbeginTransaction()Local onlyAddPartitionsToTxn([topic-A-0])record partitionOKProduce(txn records)ackAddPartitionsToTxn([topic-B-0])record partitionOKProduce(txn records)ackTransaction state: OngoingPartitions: [topic-A-0, topic-B-0]
ProducerTransactiontopic-A Leadertopic-B LeaderProducerProducerTransactionCoordinatorTransactionCoordinatortopic-A Leadertopic-A Leadertopic-B Leadertopic-B LeaderEndTxn(COMMIT)state = PrepareCommitpersist to __transaction_stateWriteTxnMarkers(COMMIT, PID, epoch)write COMMIT markerackWriteTxnMarkers(COMMIT, PID, epoch)write COMMIT markerackstate = CompleteCommitCOMMIT successTransaction markers writtenRecords now visible to consumerswith isolation.level=read_committed

Transaction markers are special control records written to each partition:

Marker TypeMeaning
COMMITTransaction committed; records are valid
ABORTTransaction aborted; records should be ignored

Markers contain:

  • Producer ID
  • Producer epoch
  • Coordinator epoch
  • Control type (COMMIT/ABORT)

When using read-process-write patterns, consumer offsets can be committed as part of the transaction:

Producer.ConsumerTransactionGroupOutput TopicProducer/ConsumerProducer/ConsumerTransactionCoordinatorTransactionCoordinatorGroupCoordinatorGroupCoordinatorOutput TopicOutput TopicbeginTransaction()poll() from input topicAddPartitionsToTxn([output-0])Produce(processed records)AddOffsetsToTxn(group.id)record group coordinatorOKTxnOffsetCommit(offsets)write pending offsetsOKEndTxn(COMMIT)WriteTxnMarkers(COMMIT)WriteTxnMarkers(COMMIT)Offsets now committedConsumer won't re-readthese records
isolation.levelBehavior
read_uncommittedRead all records including aborted transactions
read_committedRead only committed records; aborted filtered

The LSO is the offset below which all transactions are complete:

Partition Log0committed1committed2txn-A3committed4txn-A5committedLSO = 2 (txn-A still ongoing) read_committed consumers see: 0, 1read_uncommitted consumers see: 0, 1, 2, 3, 4, 5 After txn-A commits: LSO = 6read_committed now sees all records

ScenarioCoordinator Action
Producer crashes before EndTxnTransaction times out; coordinator aborts
Producer crashes during EndTxnNew producer instance completes or aborts
Network partitionTransaction times out if no progress
ProducerProducerOld CoordinatorOld CoordinatorNew CoordinatorNew CoordinatorProducerProducerOld Coordinator(Broker 1)Old Coordinator(Broker 1)New Coordinator(Broker 2)New Coordinator(Broker 2)EndTxn(COMMIT)Broker 1 failsleader election for __transaction_statebecome coordinatorload transaction statefrom __transaction_statealt[Transaction in PrepareCommit]complete commit(write markers)[Transaction in Ongoing]wait for produceror timeoutretry EndTxn(COMMIT)OK (already committed)or complete commit

If a transaction exceeds transaction.timeout.ms, the coordinator aborts it:

ConfigurationDefaultDescription
transaction.timeout.ms60000 (1 min)Max transaction duration
transactional.id.expiration.ms604800000 (7 days)Expiration for inactive transactional.id

FeatureIdempotentTransactional
DeduplicationWithin single sessionAcross sessions
ScopeSingle partitionMultiple partitions
AtomicityPer-recordMulti-record
Configurationenable.idempotence=truetransactional.id required
OverheadLowHigher (coordinator RPCs)
Use CaseRecommendation
Prevent duplicate writesIdempotent producer
Atomic multi-partition writesTransactional producer
Read-process-write exactly-onceTransactional producer
High-throughput, no atomicity neededIdempotent producer

OperationOverhead
InitProducerIdOne-time per producer start
AddPartitionsToTxnPer new partition in transaction
EndTxnTwo-phase commit across partitions
Transaction markersAdditional records in each partition

Larger transactions amortize overhead:

PatternTransactions/secThroughput
1 record per transactionLowLow
100 records per transactionMediumMedium
1000+ records per transactionHighHigh
ConfigurationDefaultTuning Guidance
transaction.timeout.ms60000Increase for long-running transactions
max.block.ms60000Time to wait for transaction coordinator
delivery.timeout.ms120000Must exceed transaction.timeout.ms

Kafka’s EOS combines idempotent producers, transactions, and transactional consumers:

Exactly-Once ComponentsIdempotent ProducerTransactionsTransactional ConsumerPID + sequence numbersDeduplication per partitionAtomic multi-partition writesConsumer offset + output atomicread_committed isolationFilter aborted recordsTogether: exactly-once processingRecord processed exactly oncefrom input to outputenablescompletes

Kafka Streams uses transactions internally for exactly-once:

processing.guaranteeBehavior
at_least_onceNo transactions; duplicates possible on failure
exactly_once_v2Transactions per task; atomic state + output

FeatureMinimum Version
Idempotent producer0.11.0
Transactions0.11.0
exactly_once_v2 (Streams)2.5.0
Transaction protocol improvements2.5.0

ConfigurationDefaultDescription
enable.idempotencetrueEnable idempotent producer
transactional.idnullTransaction identifier (enables transactions)
transaction.timeout.ms60000Transaction timeout
max.in.flight.requests.per.connection5Must be ≤5 for idempotence
ConfigurationDefaultDescription
transaction.state.log.replication.factor3__transaction_state replication
transaction.state.log.num.partitions50__transaction_state partitions
transaction.state.log.min.isr2Minimum ISR
transactional.id.expiration.ms604800000Expiration for inactive IDs
ConfigurationDefaultDescription
isolation.levelread_uncommittedread_committed for transactional
enable.auto.committrueSet to false for transactional