Skip to main content
Create your own

Idempotent Kafka Producers for Exactly-Once Semantics

Hello! Welcome to the first lesson in our module on advanced Kafka topics.

Given your extensive experience in building high-load financial and payment systems, you're undoubtedly familiar with the critical need for data integrity. A duplicate payment transaction or a re-processed trading order can have significant consequences. This lesson tackles a fundamental aspect of ensuring that integrity within Kafka.

Our focus today is on implementing an idempotent Kafka producer. This is the first major step toward achieving exactly-once semantics (EOS). We will explore the mechanism that prevents message duplication from producer retries, how to configure it, and understand its precise guarantees and limitations.

Message Delivery Semantics: A Quick Recap

In distributed messaging, we generally talk about three delivery guarantees:

  • At-most-once: Messages may be lost but are never redelivered.
  • At-least-once: Messages are never lost but may be redelivered.
  • Exactly-once: Every message is delivered once and only once.

Kafka's default configuration provides at-least-once semantics. While this prevents data loss, it opens the door to data duplication, a problem we'll solve today.


1. The Problem: Duplicates from Producer Retries

In any distributed system, transient network failures are a fact of life. Let's consider a common scenario with the default at-least-once guarantee:

  1. A producer sends a message to a Kafka topic.
  2. The broker leader receives the message, writes it to its log, and replicates it.
  3. Before the broker can send an acknowledgment (ACK) back to the producer, the network connection fails, or the broker itself crashes.
  4. The producer, having not received an ACK within its timeout, assumes the write failed. It retries sending the exact same message.
  5. A new broker leader (or the same one after recovery) receives the retried message and, having no context of the previous attempt, appends it to the log again.

The result is a duplicate message. For many applications, like financial ledgers or billing systems, this is unacceptable.

The article "What is Kafka Exactly Once Semantics?" provides a good overview of this problem.

What is Kafka Exactly Once Semantics? How to Handle It?

To better visualize this issue, please read the following section from an article by Hevo Data.

Read the section titled 'What is the need for Kafka Exactly Once Semantics?'. Focus on the example of how a broker crash can lead to message duplication.

2. The Idempotent Producer: Mechanism and Guarantees

To solve the duplicate message problem, Kafka introduced the idempotent producer. An operation is idempotent if the result of performing it multiple times is the same as performing it once. In Kafka, this means retrying a send operation won't result in the message being written to the log more than once.

This is achieved through a simple but effective protocol using a Producer ID (PID) and sequence numbers.

To understand how this works, please read the following sections from our curated resources. They explain the core mechanics from slightly different perspectives, which will help solidify your understanding.

Kafka Exactly-Once Semantics Guide

This first reading, from the 'Kafka Exactly-Once Semantics Guide', introduces the idempotent producer and its core mechanics.

First, read the section '1. Idempotent Producer'. Then, review the definitions for 'PID (Producer ID)' and 'Sequence Number' in the 'Glossary of Terms' at the end of the article.

What is Kafka Exactly Once Semantics? How to Handle It?

This second reading, from the Hevo Data article, provides another excellent explanation of the same mechanism.

Read the section 'What is an Idempotent Producer?'. This will reinforce your understanding of how the PID and sequence number are used by the broker.

Key Takeaways on the Mechanism

Let's summarize the process:

  1. Initialization: When idempotence is enabled, the producer is assigned a unique Producer ID (PID) from the broker. This PID is retained across producer restarts, but a new PID is generated if the producer application fully stops and starts again.
  2. Message Batching: For each partition it sends to, the producer maintains a sequence number, starting from 0 and incrementing with each message. This (PID, partition, sequence number) triplet uniquely identifies each message.
  3. Broker-Side Check: The broker keeps track of the latest sequence number it has successfully accepted for each (PID, partition) pair.
    • If it receives a message with sequence number = last_accepted + 1, it's the expected next message. The broker accepts it.
    • If it receives a message with sequence number <= last_accepted, it's a duplicate (from a retry). The broker discards the message but sends a success ACK to the producer.
    • If it receives a message with sequence number > last_accepted + 1, it means a message was lost in transit. The broker rejects the message with an OutOfOrderSequenceException, and the producer must handle this fatal error.

This mechanism provides exactly-once, in-order semantics per partition for a single producer session.

3. Implementation and Configuration

Enabling the idempotent producer is surprisingly straightforward. You only need to set one property. However, doing so has important implications for other producer configurations.

The following resource provides a clear Java code example and a table of the relevant configurations.

Kafka Exactly-Once Semantics Guide

Please study the following sections to see how to enable idempotence in code and understand the related configuration settings.

First, review the Java code snippet under the heading 'Enabling Idempotence (Producers)'. Then, carefully examine the configuration table under 'Key Configurations for EOS', paying close attention to the rows for 'enable.idempotence', 'acks', 'retries', and 'max.in.flight.requests.per.connection'.

The Configuration "Package Deal"

As you saw, setting enable.idempotence=true is more than just a single flag; it enforces a configuration profile required for idempotence to work correctly:

  • acks=all: This is mandatory. The broker must wait for all in-sync replicas (ISRs) to receive the message before acknowledging it to the producer. Why? If the leader accepted a message, updated its sequence number, and then crashed before replicating it, a new leader would have an old sequence number, breaking the guarantee. acks=all ensures that any acknowledged message will survive a leader failure.
  • retries > 0: The entire point of idempotence is to make retries safe. The producer will retry sending messages upon transient errors. This is usually set to a very large number (e.g., Integer.MAX_VALUE).
  • max.in.flight.requests.per.connection <= 5: Without idempotence, setting this value greater than 1 can cause message reordering on retries. However, with idempotence enabled, Kafka's protocol guarantees that messages will still be written to the log in the correct order, even with multiple in-flight requests. This is a key performance optimization that allows for higher throughput without sacrificing ordering per partition.

Here is a minimal Java configuration for an idempotent producer:

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.Properties;
import java.util.concurrent.ExecutionException;

public class IdempotentProducerExample {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

        // Enable Idempotence
        props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");

        // The following are set implicitly by the above, but shown for clarity:
        // props.put(ProducerConfig.ACKS_CONFIG, "all");
        // props.put(ProducerConfig.RETRIES_CONFIG, String.valueOf(Integer.MAX_VALUE));
        // props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, "5");

        try (Producer<String, String> producer = new KafkaProducer<>(props)) {
            for (int i = 0; i < 100; i++) {
                String key = "key-" + (i % 10); // Distribute across 10 keys
                String value = "message-" + i;
                ProducerRecord<String, String> record = new ProducerRecord<>("idempotent-topic", key, value);
                
                // .get() makes the send synchronous for this example
                producer.send(record).get(); 
                System.out.printf("Sent record with key %s and value %s\n", key, value);
            }
        } catch (InterruptedException | ExecutionException e) {
            // In a real application, handle exceptions gracefully.
            // An OutOfOrderSequenceException would manifest here and is typically fatal.
            e.printStackTrace();
        }
    }
}

4. Scope and Limitations

It is crucial to understand the precise scope of the guarantee provided by the idempotent producer:

  • Exactly-once delivery per partition: It prevents duplicates within a single topic-partition. It does not provide atomicity across multiple partitions. If a producer tries to write to partition A and then to partition B, the first write could succeed while the second fails, leaving the system in an inconsistent state.
  • Within a single producer session: The PID/sequence number mechanism works for the lifetime of a single producer instance. If your application process crashes and is restarted, it will be assigned a new PID. This means there's a possibility of duplication if the old "zombie" instance briefly wakes up and sends a message before the new one takes over.

These limitations are addressed by Kafka Transactions, which build upon the idempotent producer to provide atomic writes across multiple partitions and topics.


Conclusion

In this lesson, we've taken a significant step towards building truly reliable data pipelines with Kafka.

Key Takeaways:

  • The Problem: Default at-least-once delivery can lead to message duplication on producer retries, which is unacceptable for critical systems.
  • The Solution: The idempotent producer uses a Producer ID (PID) and per-partition sequence numbers to allow brokers to detect and discard duplicate messages from retries.
  • Configuration: You can enable this feature by simply setting enable.idempotence=true, which also enforces acks=all and enables safe retries.
  • The Guarantee: The idempotent producer provides exactly-once, in-order semantics for messages sent to a single partition within a single producer session.

Preview of the Next Lesson

While powerful, the idempotent producer alone cannot guarantee atomicity for operations that span multiple partitions or topics (a common pattern in "read-process-write" applications). In our next lesson, "Implement a transactional Kafka producer," we will build on today's concepts to explore Kafka Transactions, which provide these stronger atomic guarantees and are the final piece of the puzzle for end-to-end exactly-once semantics within Kafka.

Can't find a good explanation? Sign up and we'll make it for you

Sign up