Apache Kafka for Data & AI Engineers

Delivery Semantics & Idempotent Producers — Proving No Duplicates Under Retry


Today's Requirement

Your purchase_completed events do not just sit in Kafka. They feed two live systems.

  • A revenue dashboard adds them up in real time for the business.
  • A fraud-detection feature pipeline counts "purchases by this user in the last 10 minutes". It uses this to catch account takeovers and card testing.

Both systems need the stream to be exactly right. What does that mean?

Output / Note

No missing events, and no duplicate events.

Think about what a mistake does:

  • A missing event makes the revenue dashboard show a number that is too low. It can also hide a pattern that a fraud model should have caught.
  • A duplicate event makes the revenue number too high. It also corrupts the fraud feature, and can flag an innocent customer as a fraud suspect.

In a pipeline that feeds live dashboards and live features, a bad event is not cosmetic. It is wrong data flowing into decisions. Today we learn exactly why duplicates and gaps happen, and we test how Kafka prevents them.


What Do We Mean by "Duplicate"?

We have to be precise, because it is easy to picture the wrong scenario.

We are not talking about your own code calling produce() twice for the same purchase. Suppose your checkout service has a bug and calls charge_customer() twice, and each call produces an event. Kafka sees two deliberate instructions from your code, and it writes both. That is correct behavior. No messaging system can guess that two calls were "meant" to be one. Nothing in this lecture fixes that kind of bug. Good application logic does, for example a unique order ID that stops a double charge.

What we are talking about is narrower and stranger:

Output / Note

The producer's own retry mechanism accidentally creating a second copy of a message that you sent only once.

Your code calls produce() one time. Somewhere between your code and the topic, a duplicate appears anyway. It is not your mistake.


How Kafka Protects You From Missing Events

What if a message does not reach the broker?

Kafka's answer is two things working together: acknowledgment and retry.

  • The producer sends a batch to the broker.
  • The broker saves it and sends back an acknowledgment.
  • If the acknowledgment arrives, the producer knows the message is safe.
  • If it does not arrive, the producer automatically sends the message again (a retry).

You did not write any retry code. The producer does it for you. It is a deliberate reliability feature. Without it, a small network problem could silently drop real data.

How a Retry Actually Works

Let's look at the retry step by step, because it explains both the protection and the risk.

How does the producer decide that something failed?

There are three kinds of failure:

  1. The connection to the broker breaks, and the client gets an error.
  2. The connection stays open, but no answer comes back for too long.
  3. The broker answers with a temporary error, for example because it could not get the confirmations it needed in time.

These settings control the timing:

SettingDefaultWhat it does
socket.timeout.ms60 secondsHow long the client waits for the answer to a request. If nothing comes, the request has failed
request.timeout.ms30 secondsHow long the broker waits for the other replicas to confirm the write (it matters when acks=all). If they do not confirm in time, the broker answers with an error
retry.backoff.ms100 msThe pause before a retry. It grows step by step, up to retry.backoff.max.ms (1 second by default)
retriesabout 2 billionHow many retries are allowed. Effectively unlimited
message.timeout.ms5 minutesThe total time one message may take, retries included. After this, the producer gives up

So the flow is:

  1. Send the batch, and wait for the answer.
  2. A failure is detected.
  3. Wait for the backoff time (about 100 ms, a little longer each time).
  4. Send the same batch again.
  5. Repeat until success, or until message.timeout.ms (5 minutes) has passed.

If the five minutes pass, the producer gives up, and your delivery callback receives an error. That is why we always check err in the callback. So "at-least-once" does not mean "always delivered, whatever happens". It means: the producer tries as hard as it can, and tells you if it failed. For most systems, these defaults are well chosen, and you rarely change them. But it is important to know how they work.


How a Good Feature Causes a Real Problem

Now picture exactly what can go wrong.

Your background thread sends a batch. The broker receives it and writes it successfully. It sends the acknowledgment back. But that acknowledgment gets lost on the way. It could be a network blip, or the broker could crash right after writing.

What does the producer see?

Silence. It cannot tell the difference between "the write failed" and "the write worked, but the answer got lost". Both look the same.

So what should it do? Giving up is not safer, because it could lose a real purchase. So it retries.

But look at the side effect. The original write did succeed. The broker has no way to know that this new request is a resend. So it does the only thing it can do: it writes the message again.

Output / Note

A feature that exists to protect your data just created two copies of it.


How the Broker Recognizes a Retry

This is the fix, and it is elegant. It is called an idempotent producer. (Idempotent means: doing it twice has the same result as doing it once.)

When idempotence is on, two things happen that you never see:

  1. When the producer starts, the broker gives it a unique Producer ID (PID).
  2. Every message the producer sends to a partition gets a sequence number: 0, then 1, then 2, growing by one, separately for each partition.

How the broker recognizes a retryHow the broker recognizes a retry

The broker remembers, for each producer and each partition, the latest sequence numbers it has accepted. For example: "from PID 77, I accepted up to sequence 42".

Now the lost-acknowledgment story again:

  1. The producer sends sequence 42 from PID 77. The broker writes it.
  2. The acknowledgment is lost.
  3. The producer retries. It sends sequence 42 from PID 77 again.
  4. The broker checks: "I already have that." It does not write it again. But it still answers "success", because as far as the producer needs to know, it was delivered.

The retry is safe. No duplicate. The sequence numbers are added by the producer's background thread, and the checking is done by the broker. You write no extra code.

The Important Fine Print

Why did we say "your own duplicate produce() calls" are not covered?

Because this guarantee works only when the same producer resends the exact same attempt. Two produce() calls from your code get two different sequence numbers, on purpose. The broker cannot know they were "supposed" to be the same. Idempotence does not look at the message content to find look-alikes. It tracks attempts, by sequence number, per PID.

Three more points about the scope:

  • It works within one producer session. If the producer process restarts, it gets a new PID, and the broker's memory of the old one does not help.
  • It protects against duplicates from retries. It does not promise delivery. If the message cannot be delivered within message.timeout.ms, it fails, and your callback tells you.
  • It is exactly-once at the level of one producer writing to one partition. Full end-to-end exactly-once (a consumer reads, transforms, and produces again) is a bigger subject with transactions. That comes later in the course.

Three Ways a Producer Can Handle Retries

Now you know the mechanism. Here are the three named choices, and the settings behind each.

Three delivery semanticsThree delivery semantics

At-Most-Once: Send and Forget

Do not wait for any answer, so there is nothing to retry. The setting is acks=0. The producer does not expect an acknowledgment, and never retries.

  • It can never create a duplicate.
  • But a message can vanish, and nobody knows. For our fraud feature, a real purchase would silently not be counted.
  • It is the fastest option.

Good for data you can afford to lose a little of. For example, sensor readings from a million devices, where you only look for trends, and losing 0.1% changes nothing.

At-Least-Once: Retry Whenever in Doubt

This is what most people mean by "reliable". The producer sends, waits for the acknowledgment, and retries if it does not come. It is the default behavior of confluent-kafka. You do not have to configure anything special.

  • If the producer reports success, the broker has the message.
  • But, exactly as we just saw, a retry can create a duplicate.

Good when duplicates are harmless. For example, an "order confirmed" status event. If you send it three times, the latest status is still "order confirmed". Nothing breaks.

Output / Note

About the acks setting. acks decides how many confirmations the broker collects before it answers. acks=1 means only the partition leader. acks=all means all the in-sync replicas. In the video, the default was described as acks=1. In confluent-kafka, the default is actually acks=all. We use acks=1 explicitly in the unsafe script below, on purpose. We say more about acks in the tuning lecture.

There is one more thing to know about acks=1. The leader confirms the write before the followers have copied it. If the leader crashes at that moment, the message can be lost even though the producer was told it succeeded. So acks=1 at-least-once can lose data in a leader failure. acks=all protects against this. min.insync.replicas (next section) decides how strong that protection really is.

Idempotent Delivery: The Safe Retry

Set enable.idempotence=True, and the client does the rest. It automatically:

  • sets acks=all,
  • keeps the retries unlimited,
  • limits max.in.flight.requests.per.connection to 5 (at most 5 batches waiting for an answer at the same time), so the sequence numbers stay meaningful.

Retries still happen exactly as before. They are simply safe.

Output / Note

Careful: these settings are adjusted for you only if you did not set them yourself in a conflicting way. If you set enable.idempotence=True and also acks=1, the producer refuses to start. The error is: acks must be set to all when enable.idempotence is true. The same happens if you set max.in.flight above 5. This is why the idempotent script below has no acks line.

One Broker-Side Setting: min.insync.replicas

With acks=all, the leader waits for every replica that is currently in sync (the ISR, from the chapter on brokers). But what if only the leader is left in the ISR? Then "all in-sync replicas" is just the leader, and acks=all is no stronger than acks=1.

min.insync.replicas (a topic or broker setting) sets the minimum size of the ISR for a write to be accepted. For example, 2 means: if fewer than two replicas are in sync, the write is refused, so you do not get a false sense of safety. It does not change anything about duplicates. But it decides how meaningful acks=all is.


Let's Build It

We build two nearly identical producers, one without idempotence and one with it, so we can watch the difference and not just take it on trust.

The Shared Setup

Both scripts start like every producer you have built:

python
BOOTSTRAP_SERVERS = "localhost:9092,localhost:9094,localhost:9095"

The Unsafe Version

python
TOPIC = "purchase-events-unsafe" producer_config = { "bootstrap.servers": BOOTSTRAP_SERVERS, "client.id": "purchase-producer-unsafe", "enable.idempotence": False, # librdkafka's actual default — being explicit here "acks": 1, "message.timeout.ms": 300000, }

Three settings matter here:

  • enable.idempotence is False. That is the real default. We write it to be explicit.
  • acks is 1. We choose it on purpose, to get the weaker, faster kind of at-least-once.
  • message.timeout.ms is 300000 milliseconds: five minutes, which is also the default. It is how long the producer may keep retrying a message. After that, it gives up, and the callback reports an error.

The Idempotent Version

python
TOPIC = "purchase-events-idempotent" producer_config = { "bootstrap.servers": BOOTSTRAP_SERVERS, "client.id": "purchase-producer-idempotent", "enable.idempotence": True, "message.timeout.ms": 300000, }

Two differences. enable.idempotence is True. And the acks line is gone, because idempotence sets acks=all by itself, as we saw. Everything else is the same.

The Shared Loop: 100,000 Uniquely Numbered Events

python
def delivery_report(err, msg): if err is not None: print(f"Delivery failed for purchase_id={msg.key().decode('utf-8')}: {err}") def main(): for purchase_id in range(100000): event = { "purchase_id": purchase_id, "user_id": "user-501", "amount": 49.99, "purchased_at": time.time(), } producer.produce( topic=TOPIC, key=str(purchase_id).encode("utf-8"), value=json.dumps(event).encode("utf-8"), callback=delivery_report, ) producer.poll(0) time.sleep(0.01) producer.flush() print("All purchase events sent (or failed) — check the count in the topic.")

A few things to notice:

  • The callback prints only failures. With 100,000 messages, we do not want the screen full of success lines. A quiet screen means no failures.
  • Every event has a unique purchase_id. So counting the messages in the topic tells us at once whether something was duplicated or lost. We do not have to inspect each one.
  • time.sleep(0.01) waits 10 milliseconds after each message. So 100,000 messages take about 16 to 17 minutes. That is on purpose: it gives us a long time window to break things.

Let's Prove It

Now let's watch this against your own cluster. We will run the producer for about 16 minutes, and kill brokers while it runs, to force retries.

Step 1: Start Your Cluster

bash
docker compose start docker ps

Step 2: Create Two Topics

We use two separate topics, so the "before" and "after" runs never mix. We give them their own topics because purchase events need stricter durability than a page view.

bash
docker exec -it kafka-1 bash export PATH=$PATH:/opt/kafka/bin
bash
kafka-topics.sh --create --topic purchase-events-unsafe \ --bootstrap-server kafka-1:29092 --partitions 3 --replication-factor 3 kafka-topics.sh --create --topic purchase-events-idempotent \ --bootstrap-server kafka-1:29092 --partitions 3 --replication-factor 3

Step 3: Run the Experiment Without Idempotence

Open a terminal on your host machine (not docker exec), in the 01_producers folder:

bash
python 3-purchase_producer_no_idempotence.py

While it runs, open a second host terminal and create some failures. Kill two brokers, wait a little, then start them again:

bash
docker kill kafka-1 docker kill kafka-3

Wait 20 to 30 seconds. Then:

bash
docker start kafka-3 docker start kafka-1

Wait a few minutes, then do it again with a different broker, for example:

bash
docker kill kafka-2

and later

bash
docker start kafka-2

Check with docker ps that all three brokers are up again. You can repeat this as many times as you like during the run. The more failures, the better the test.

What will you see?

While brokers are down, the producer terminal prints connection errors. That is normal. It is the client noticing failures and retrying. When the brokers come back, it recovers and continues.

A few rules to follow:

  • Keep each outage short, well under 5 minutes. If it is longer than message.timeout.ms, the producer gives up on some messages, and you will see Delivery failed for purchase_id=... lines.
  • With two of three brokers down, the cluster may stop accepting writes until you start at least one again. That is fine. The producer keeps retrying.

Wait for the script to finish. It prints All purchase events sent (or failed). If no Delivery failed line appeared, every message was reported as delivered.

Step 4: Count What Landed

In your broker terminal:

bash
kafka-console-consumer.sh \ --topic purchase-events-unsafe \ --bootstrap-server kafka-1:29092 \ --from-beginning \ --timeout-ms 5000 | wc -l

We sent exactly 100,000 uniquely numbered events. What can the count tell us?

  • Exactly 100,000: nothing was lost or duplicated in this run.
  • More than 100,000: a retry created duplicates. Your revenue dashboard would be inflated.
  • Fewer than 100,000, with no Delivery failed line: a message was acknowledged and then lost, for example when a leader crashed before the followers had copied it. That is the acks=1 risk.

Be honest about what you see. In the recorded run, the count was exactly 100,000, and no duplicates appeared. A duplicate needs a very exact timing: the write must succeed, and the acknowledgment must be lost, and then the retry must happen. Failures like this are hard to reproduce on purpose. If you also see exactly 100,000, that is fine. It does not mean the risk does not exist. In production, with far more messages and constant network noise, it does happen. You can repeat the experiment, with more kills, to try to catch one.

Step 5: Run the Same Experiment With Idempotence

bash
python 4-purchase_producer_idempotent.py

Create the same kind of failures while it runs (kill and restart kafka-1, kafka-3, and kafka-2, as before). Wait for it to finish.

Step 6: Count Again

bash
kafka-console-consumer.sh \ --topic purchase-events-idempotent \ --bootstrap-server kafka-1:29092 \ --from-beginning \ --timeout-ms 5000 | wc -l

This count can never be higher than 100,000, however many times you repeat the experiment, because duplicates from retries are stopped. It will be exactly 100,000, unless a Delivery failed line printed. That would mean a message could not be delivered within message.timeout.ms, and the callback told you so.

Repeating an Experiment

Each run adds its 100,000 events to the topic. If you run a script again on the same topic, the count doubles. To start clean, delete and recreate the topic in the broker shell, for example:

bash
kafka-topics.sh --delete --topic purchase-events-unsafe --bootstrap-server kafka-1:29092 kafka-topics.sh --create --topic purchase-events-unsafe \ --bootstrap-server kafka-1:29092 --partitions 3 --replication-factor 3

Step 7: Stop Your Cluster

bash
docker compose stop

Common Mistakes at This Stage

Believing idempotence catches your own duplicate produce() calls. It only makes the producer's own retries safe.

Believing at-least-once can never lose data. With acks=1, a leader crash can lose an acknowledged message.

Setting acks=1 and enable.idempotence=True together. The producer will not start. Remove the acks line, or set it to all.

Never checking the callback. A failed delivery is reported only there. A quiet screen is good only if your callback really prints failures.

Expecting idempotence to survive a restart. A new producer process has a new Producer ID.


Checklist

  • Understood the difference between a duplicate from your own code and a duplicate from a retry, and that only the second is solved here
  • Understood how acknowledgment and retry work, including the timing settings and the five-minute limit
  • Understood how a lost acknowledgment plus a retry creates a duplicate write
  • Understood the Producer ID plus sequence number mechanism
  • Know the scope: one producer session, one partition, retries only
  • Understood the three delivery semantics and the settings behind each (including that the default acks is all)
  • Know why acks=1 can lose data, and that min.insync.replicas decides how strong acks=all is
  • Understood the two differences between the two experiment scripts
  • Ran the unsafe experiment with broker failures and inspected the count
  • Ran the idempotent experiment and confirmed the count did not exceed 100,000

You did not just read about idempotent producers. You tried to break one, on purpose, against real failing brokers. Next: we tune throughput and durability directly, with batching, compression and acks, and we look at the most common ways producers go wrong in production.