03 / 05

How would you implement a consume-transform-produce pipeline using Kafka transactions?

Difficulty: 7/10
Transactional processing, Exactly-once, Idempotency

Implementing a Consume-Transform-Produce Pipeline with Kafka Transactions

A consume-transform-produce pipeline with Kafka transactions has three moving parts: a transactional producer, a consumer with auto-commit disabled and isolation.level=read_committed, and a loop that ties the consumer offsets to the producer transaction. The pattern is: poll records, begin transaction, transform and send each record to the output topic, send the consumer offsets to the transaction, then commit. If anything fails, abort the transaction, and the consumer will re-read the records because the offsets were not committed. This gives exactly-once semantics for the Kafka-to-Kafka path because the output records and the consumer offsets are committed atomically. The key detail is that the offsets must be sent to the transaction via sendOffsetsToTransaction, not committed separately with consumer.commitSync(). If you commit offsets separately, you reintroduce the window where a crash between the output write and the offset commit causes duplicates.

The mechanism depends on a stable transactional.id. Each producer instance must have a unique transactional.id that does not change across restarts. On restart, the producer calls initTransactions(), which fences any previous instance with the same transactional.id and recovers any in-flight transactions. This is what prevents a zombie producer from continuing to write after a new instance has taken over. The consumer must be configured with enable.auto.commit=false because the offsets are managed by the transaction, not by the consumer. It should also use isolation.level=read_committed so that it does not read uncommitted or aborted records from upstream. The transaction timeout (transaction.timeout.ms on the producer, transaction.max.timeout.ms on the broker) bounds how long a transaction can stay open; if it exceeds the timeout, the coordinator aborts it. This is a safety net against a stuck producer but also means that long-running processing must either fit within the timeout or be broken into smaller transactions.

A common mistake is to use the same transactional.id for multiple concurrent producer instances. This causes fencing and is one of the most frequent production issues with transactions. Another mistake is to process records asynchronously after commitTransaction, which breaks the guarantee because the side effect is not covered by the transaction. The trade-off is between batch size and latency: larger batches mean fewer transactions and higher throughput but longer transaction durations and more memory. Smaller batches mean lower latency but more coordinator round trips. For high-throughput pipelines, you tune max.poll.records and the transaction timeout to balance this. Version note: Kafka Streams exposes this as processing.guarantee=exactly_once_v2 (KIP-447, Kafka 2.5), which is more scalable than the older exactly_once. If you are writing the pipeline manually, the same principles apply. Also note that transactions are scoped to a single Kafka cluster; if your output topic is in a different cluster, transactions do not extend there.

javascript
  1. 1

    Use a transactional producer with a unique, stable transactional.id.

  2. 2

    Disable auto-commit on the consumer and set isolation.level=read_committed.

  3. 3

    Send consumer offsets via sendOffsetsToTransaction, not commitSync, to keep them in the transaction.

  4. 4

    Begin transaction, transform and send records, send offsets, commit; abort on failure.

  5. 5

    On restart, initTransactions fences any previous instance with the same transactional.id.

  6. 6

    transaction.timeout.ms bounds transaction duration; long processing must fit or be split.

  7. 7

    Transactions are scoped to one Kafka cluster and do not cover external side effects.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.