Kafka Transactions and Exactly-Once Processing, Explained
Kafka transactions let stream-processing applications publish related records and commit the offsets they consumed as one atomic operation. For developers and platform engineers, this is the foundation of exactly-once processing inside Kafka: after a crash, a batch is either visible with its input progress or retried without exposing partial output. The guarantee has a boundary, however. It does not make an unrelated database write, HTTP request, or email part of the Kafka transaction.
Why Kafka Transactions Are Being Discussed
Event-driven services often read records, transform them, publish results, and then record their progress. A crash between any two of those steps can repeat work or lose output. That matters when a pipeline updates a derived topic, calculates account totals, or moves data between Kafka topics: duplicate results can corrupt downstream state even though the broker retained the original events.
Kafka transactions provide a broker-coordinated way to group output records and consumed offsets. They are useful when both sides of the operation are in Kafka, rather than a universal switch that makes every application effect exactly once. The distinction is especially important in microservice designs where Kafka is only one participant among databases and remote services.
What Are Kafka Transactions?
A Kafka transaction is a group of writes from a transactional producer that become committed or aborted together. A producer can write records to multiple topic partitions in a transaction. A stream processor can also include its input consumer offsets in that same transaction, so output and progress are committed together.
The Apache Kafka project describes Kafka as an event-streaming platform built around producers, consumers, topics, and durable logs. Transactions add atomicity to the producer’s writes; they do not change the underlying partition model or create a global order across topics.
An idempotent producer and a transactional producer solve related but different problems. Idempotence uses producer IDs and sequence numbers to prevent retrying a send from appending a duplicate record. A transaction groups multiple writes and, when used by a consume-transform-produce application, can atomically include offsets. The Kafka documentation on delivery semantics explains the producer guarantees and transaction boundaries.
The Problem: Processing and Offsets Can Disagree
Kafka stores records independently from the progress a consumer group has committed. If an application writes a result and crashes before committing the input offset, the next attempt can process the same input and write the result again. If it commits the offset first and crashes before writing the result, that input may never produce an output.
Disabling automatic offset commits and choosing a careful order helps, but does not remove the crash window between separate broker operations. Writing to an external database makes the gap larger: Kafka cannot atomically commit its offset and an unrelated SQL transaction just because both calls happen in one function.
The transaction design described in KIP-98 adds a coordinator and commit markers so participating Kafka records can be interpreted consistently. The application still has to use the producer and consumer APIs in a way that keeps input progress inside the same transaction as the output.
How Kafka Transactions Work
The application configures a transactional ID and initializes its producer. Kafka associates that ID with a producer epoch, which allows a newer producer instance to fence off an older, still-running instance using the same logical identity. The producer starts a transaction, sends records, adds offsets if it is processing a consumer group, and then commits or aborts.
At a high level, a consume-transform-produce batch follows this sequence:
- The consumer polls records and keeps automatic offset commits disabled.
- The producer begins a transaction and sends the derived records.
- The application adds the next offsets for the consumed partitions to the transaction.
- The producer commits. Kafka makes the records and offsets visible as one successful unit.
- If processing fails before commit, the producer aborts. A later attempt can read the uncommitted input again.
Transaction metadata is coordinated by Kafka and commit or abort markers are written to the affected logs. Consumers configured with read_committed return only committed transactional records; they do not expose aborted output. This isolation level can also make a consumer wait behind an open transaction in a partition, because it must not pass the last stable offset while a transaction there is unresolved.
| Mechanism | What it makes safe | Atomic boundary | Typical use |
|---|---|---|---|
| Idempotent producer | Retries of individual sends | A producer’s append operations | Avoiding duplicate records from network retries |
| Producer transaction | A set of output records | Records written by one transaction | Publishing a consistent multi-record result |
| Transaction plus offsets | Kafka input progress and Kafka output | Output records and consumer-group offsets | Kafka-to-Kafka processing that can be retried |
| External idempotency or outbox | Effects involving a database or remote service | Defined by the external system and application | Coordinating Kafka with systems outside Kafka |
Kafka’s transaction guarantee is not global ordering. Ordering remains per partition, and a transaction’s records can span partitions without creating a single total order between them. Producers, brokers, and consumers must also use compatible transactional settings for the intended read-committed behavior.
Key Components and Concepts
- Idempotent producer: Use
enable.idempotence=trueto make producer retries safe from duplicates caused by retrying the same send. Modern clients enable this by default when compatible settings are used, but making it explicit can clarify an application’s intent. Idempotence alone does not atomically connect output records with consumer offsets. - Transactional ID: Set
transactional.idto a stable identity for a logical producer instance. It must be unique among concurrently active instances. Reusing the same ID when restarting that logical instance enables Kafka to resolve old transactional state and fence an obsolete producer; using one ID for every replica causes them to fence one another. - Transaction coordinator: The broker-side coordinator tracks transaction state and coordinates commit or abort markers. This work adds network round trips and broker state, so transactions are not free throughput-wise.
- Consumer isolation:
read_committedhides aborted transactional output.read_uncommittedmay return records from transactions that later abort, so applications that depend on transactional visibility should configure consumers accordingly. - Offset metadata:
sendOffsetsToTransactionadds the consumer group’s next offsets to the producer’s current transaction. The application must turn off auto-commit and pass the right offsets and group metadata for the records it processed. - Kafka Streams processing guarantee: Kafka Streams can manage the transaction and state-store coordination itself. Its
processing.guarantee=exactly_once_v2setting is the supported option for exactly-once stream processing; it is not the same as a generic guarantee for arbitrary application code.
Real-World Use Cases
Transactions fit a stream processor that reads orders from one topic, validates or enriches them, and writes results to another topic. If the process crashes during a batch, downstream consumers see either the committed batch with its offsets or no committed output from that attempt.
They are also useful when one event produces multiple Kafka records that must agree, such as a state update and a corresponding audit event. Kafka Streams applications can use transactions for stateful operations and changelog updates under the configured processing guarantee.
Transactions are not a replacement for every reliability pattern. If a service changes a relational database and publishes an event, the transactional outbox pattern can atomically save the database change and event intent. A consumer that updates an external system still needs an idempotency key, deduplication record, or another system-specific recovery design.
Practical Guide: Configure a Transactional Processor
Use a Kafka distribution or managed cluster with transaction support, and create the input and output topics before running the application. A local development topic can use one replica:
bin/kafka-topics.sh --create \
--topic orders.raw \
--bootstrap-server localhost:9092 \
--partitions 3 \
--replication-factor 1
For production, set replication and minimum in-sync replicas to match the cluster’s failure and durability requirements rather than copying the single-node example.
The following Java example uses the Apache Kafka client library. It commits the next offset for each input partition together with the output batch. It is a single logical processor; a multi-instance deployment needs a distinct, stable transactional ID per live processor instance.
import java.time.Duration;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Properties;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.TopicPartition;
public class OrderTransformer {
public static void main(String[] args) {
Properties consumerProperties = new Properties();
consumerProperties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
consumerProperties.put(ConsumerConfig.GROUP_ID_CONFIG, "order-transformer");
consumerProperties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
consumerProperties.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
consumerProperties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
"org.apache.kafka.common.serialization.StringDeserializer");
consumerProperties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
"org.apache.kafka.common.serialization.StringDeserializer");
Properties producerProperties = new Properties();
producerProperties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
producerProperties.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-transformer-0");
producerProperties.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
producerProperties.put(ProducerConfig.ACKS_CONFIG, "all");
producerProperties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
"org.apache.kafka.common.serialization.StringSerializer");
producerProperties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
"org.apache.kafka.common.serialization.StringSerializer");
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProperties);
KafkaProducer<String, String> producer = new KafkaProducer<>(producerProperties)) {
consumer.subscribe(List.of("orders.raw"));
producer.initTransactions();
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
if (records.isEmpty()) {
continue;
}
producer.beginTransaction();
try {
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
for (ConsumerRecord<String, String> record : records) {
producer.send(new ProducerRecord<>(
"orders.validated", record.key(), record.value()));
TopicPartition partition =
new TopicPartition(record.topic(), record.partition());
offsets.put(partition, new OffsetAndMetadata(record.offset() + 1));
}
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction();
} catch (RuntimeException error) {
producer.abortTransaction();
throw error;
}
}
}
}
}
In real deployments, choose a transaction timeout that covers normal processing but does not leave failed work open indefinitely. Handle producer fencing and authorization or configuration errors as fatal conditions that require operator or instance recovery; do not blindly retry them forever. Monitor transaction aborts, producer errors, consumer lag, and processing latency. The Kafka consumer group guide explains how group membership and committed offsets behave during rebalances.
For a Kafka Streams application, configure the processing guarantee directly:
processing.guarantee=exactly_once_v2
Use it with the appropriate application ID and broker permissions. Verify behavior by restarting the processor during a test workload, then confirm the read-committed output and consumer-group offsets do not contain a partially committed batch. The Apache Kafka event streaming guide covers topics, partitions, and consumer groups needed to interpret those checks.
Common Misconceptions
“Idempotence means the whole workflow runs exactly once.”
Idempotence addresses retry duplicates for producer sends. A separate offset commit or database write can still fail independently. Transactions extend the boundary to a set of Kafka records and, when explicitly included, consumer offsets.
“Exactly once means no code ever runs twice.”
A processor can execute the same input again after an aborted attempt. The guarantee is about committed Kafka-visible results and offsets, not how many times application code was invoked. Processing must still tolerate retries, and external calls need their own deduplication or idempotency.
“A Kafka transaction includes my database automatically.”
It does not. Kafka’s coordinator cannot commit a relational database transaction, payment request, or email delivery. Use an outbox, idempotent consumer, or carefully designed saga when the workflow crosses system boundaries.
Related Articles
- Start with event streaming with Apache Kafka for the topic, partition, and broker model.
- See Kafka consumer groups and partition rebalancing for offset progress and membership changes.
- Compare Kafka transactions with the transactional outbox pattern for database-to-broker delivery.
- Review event-driven microservices for the wider integration patterns around asynchronous services.
Changelog
- 2026-10-11: First publication.

