Preventing Duplicate Events in Kafka with Idempotent and Transactional Producers
Learn how Kafka’s idempotent and transactional producer APIs eliminate duplicate records caused by retries, and how to configure them safely for exactly‑once semantics in production pipelines.
13 Nov 2025, 19:46 UTC

The Duplicate Event Problem
In a high‑throughput Kafka pipeline, a producer may retry sending a record when a transient network glitch or broker hiccup occurs. Without safeguards, each retry can create a duplicate in the target partition, breaking downstream consumers that expect a clean stream. Application‑level deduplication is error‑prone and adds latency.
Why Rely on Kafka’s Built‑in Guarantees
Kafka 0.11 introduced two complementary features that solve this at the broker level:
- Idempotent producer – guarantees that a single producer instance writes each record at most once to a partition, even after retries.
- Transactional producer – extends idempotence across multiple partitions and topics, allowing a batch of records to be committed atomically.
These guarantees are broker‑scoped: they ensure that a record is written only once to Kafka, but downstream systems still need to handle idempotency if they persist data elsewhere.
Configuring an Idempotent Producer
Enable idempotence with a minimal configuration. The following Java example shows the required properties and a simple publish loop.
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// --- Idempotence settings ---
props.put("enable.idempotence", "true"); // Enables sequence number tracking
props.put("max.in.flight.requests.per.connection", "5"); // Must be <=5 when idempotence is true
props.put("retries", Integer.MAX_VALUE); // Unlimited retries are safe
Producer<String, String> producer = new KafkaProducer<>(props);
for (int i = 0; i < 1000; i++) {
ProducerRecord<String, String> record = new ProducerRecord<>("my-topic", "key" + i, "value" + i);
producer.send(record, (metadata, exception) -> {
if (exception != null) {
// Handle send failure – the record may have been sent once already
}
});
}
producer.flush();
producer.close();
Key points to verify:
- Broker version ≥ 0.11 and client library ≥ 0.11.
- All producer instances use the same
client.idonly if you intend to share a producer ID (rare). - Set
max.in.flight.requests.per.connectionto 5 or less to avoid out‑of‑order duplicates.
Building a Transactional Pipeline
Transactional producers allow you to read from an input topic, transform the data, and write to one or more output topics atomically. Below is a concise example that demonstrates the full transaction life‑cycle.
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("group.id", "transactional-processor");
props.put("enable.idempotence", "true");
props.put("transactional.id", "processor-1"); // Must be unique per producer instance
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions(); // Contact transaction coordinator
try {
producer.beginTransaction();
// Read from input topic (simplified – real code would use a consumer loop)
ConsumerRecord<String, String> input = fetchRecord("input-topic");
// Transform
String transformed = transform(input.value());
// Write to output topics
producer.send(new ProducerRecord<>("output-topic", input.key(), transformed));
// Commit the whole transaction
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
throw e;
}
producer.close();
Consumer side:
- Set
isolation.level=read_committedto avoid seeing uncommitted records. - Use the same
group.idas the producer if you want to coordinate offsets.
Trade‑offs and Caveats
- Latency – Each transaction requires a round‑trip to the transaction coordinator. High throughput pipelines may see a 10–20 % increase.
- Transaction timeout – The broker’s
transaction.timeout.msdefaults to 15 min. If your application holds a transaction open longer, it will be aborted automatically. - Stale transactions – A crashed producer can leave a transaction open. Cleanup is automatic after timeout, but repeated failures can temporarily block new transactions.
- Exactly‑once is broker‑scoped. If you persist to a database, you still need idempotent writes there.
Practical Checklist Before Production
- Run
kafka-broker-api-versions.sh --bootstrap-server broker1:9092to confirmidempotent_producerandtransactional_producersupport. - Deploy the example to a staging cluster, trigger a network drop, and verify that no duplicate keys appear in the output topic.
- Monitor
producer.acksandtransaction.commit.lag.msmetrics via JMX to detect unexpected delays. - Set
transaction.timeout.msto a value that comfortably exceeds your maximum processing time per batch. - Ensure downstream consumers use
isolation.level=read_committedand handle idempotent writes if they write to external stores.
By enabling idempotent or transactional producers, you offload duplicate suppression from your application logic to Kafka, simplifying the code base and improving reliability in the face of transient failures.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.