Apache Kafka for Data & AI Engineers

Offset Management & Commit Strategies — Surviving a Poison-Pill Message


The Requirement

Two problems have been sitting quietly under everything you have built so far.

Problem 1. Your consumer restarts. How does it know exactly where to continue? It must not redo everything, and it must not silently skip data it never handled.

Problem 2. Real event streams sometimes contain a bad message: broken JSON, a missing field, something a bug produced somewhere. What does your consumer do then? Does it crash? Does it get stuck on the same message forever?

Both problems come from the same mechanism: offset commits. Today we fix both.


What Committing an Offset Actually Means

Is a committed offset "the last message I read"?

No. A committed offset is the next offset you should start from.

If you have fully handled everything up to offset 40, you commit offset 41. After a restart, the consumer resumes from 41. (In confluent-kafka, when you store a message's offset, the library adds the +1 for you.)

Where does Kafka keep it?

In the internal topic __consumer_offsets, the same topic that picks your group coordinator. Each group has its own committed position per partition.

An offset actually goes through three stages. Keep these three words in mind, because the whole lecture depends on them:

StageWhereMeaning
ReceivedYour codepoll() gave you the message
StoredIn memory, inside the client"This offset is ready to be committed"
CommittedIn KafkaThe official saved position

Only a committed offset survives a restart.


auto.offset.reset — Clearing Up the Most Common Confusion

In the first lecture, you set auto.offset.reset to "earliest" with a short explanation. Now that you know what a committed offset is, the full explanation will finally make sense.

Does auto.offset.reset decide where a consumer starts every time it runs?

No. This is the biggest misconception about this setting. It matters in only one situation.

When does auto.offset.reset actually apply?When does auto.offset.reset actually apply?

Every time a consumer starts, it asks one question for each partition: "Is there already a valid committed offset for this group.id?"

  • If yes, the consumer resumes from it. auto.offset.reset is never even looked at.
  • If no, only then does auto.offset.reset decide the starting point.

When is there "no valid committed offset"?

Three common cases:

  1. A brand-new group. This group.id has never run before.
  2. The committed offset expired. Kafka deletes a group's committed offsets after the group has had no members for a while. The default is 7 days (the broker setting offsets.retention.minutes). So a consumer that was switched off for a week and then restarted behaves like a brand-new group.
  3. The offset is out of range. The data it points to no longer exists. For example, old messages were deleted by the topic's retention, or the topic was recreated with fewer messages.

That is why it is called a reset, not a default. It answers the question: "What do we do when we have no reliable position?"

The Three Options, and When to Choose Each

ValueBehaviorChoose this when...
earliestStart from the beginning of the topic's retained historyYou cannot afford to miss data: a new analytics or audit consumer that needs the full picture
latestStart from whatever is produced from now onYou only care about live activity: a real-time dashboard where old backlog is irrelevant
errorDo not guess. Raise an error for each partition and let your code decideHigh-stakes pipelines where silently picking a default could hide a real problem

Which one is the default if I set nothing?

latest. That is why we set earliest on purpose from the start.

What does error look like in practice?

The consumer reads nothing. Instead, poll() returns an error message for each partition, with the code _AUTO_OFFSET_RESET and the text "no previously committed offset available". These errors are non-fatal.

Output / Note

Careful: the value is spelled error, not none. If you type "none", the consumer refuses to start with Invalid value "none" for configuration property "auto.offset.reset".

Output / Note

Optional experiment: Take the first lecture's script. Change GROUP_ID to a brand-new name and set auto.offset.reset to "error". Run it. You will see one non-fatal error per partition, and no events. Then try "latest" (no events until you produce new ones) and "earliest" (all the history).

Let's think about it with our own pipeline. Suppose you start a new group to backfill six months of clickstream-events into an analytics table. You want earliest, because missing history defeats the purpose. But if you launch a live "what is trending right now" dashboard, replaying six months of old clicks at startup would be wrong. You want latest.


Now let's go back to committing. Even with a good starting point, the default commit behavior is more dangerous than it looks.


The Hidden Danger of the Defaults

confluent-kafka has two separate settings that work together here:

  • enable.auto.offset.store (default True). The moment poll() hands you a message, this marks that message's offset as "ready to commit" in memory. Not after you finished processing. The instant you receive it.
  • enable.auto.commit (default True). On its own timer, every auto.commit.interval.ms (default 5 seconds), whatever is marked "ready" is written to Kafka as the committed offset.

Default settings vs. the safe patternDefault settings vs. the safe pattern

Let's walk through the left side of the picture.

  1. poll() hands you a message. It is marked "ready" immediately, before you have looked at it.
  2. Your processing starts, and takes longer than expected. Maybe a slow downstream call. Maybe the message itself causes trouble.
  3. While you are still working, the auto-commit timer fires anyway. It commits an offset for a message that is not finished yet.
  4. Processing crashes, and you restart. The consumer resumes after that offset.

What happened to that message?

It is gone. It was not retried. It was not logged as missing. It was silently skipped, because Kafka was told it was done when it never was.

This is not just theory. In a test of exactly this scenario, a consumer with default settings crashed while processing offset 4. Kafka's committed offset for that partition was 5. The unfinished message had been marked done.


Let's Build It — The Correct Pattern

The fix is small, but it needs both settings to be handled on purpose.

Step 1: Turn Off Auto-Store, Keep Auto-Commit

python
consumer_config = { "bootstrap.servers": BOOTSTRAP_SERVERS, "group.id": GROUP_ID, "client.id": "safe-consumer", "auto.offset.reset": "earliest", # Keep auto-commit's periodic write-to-Kafka behavior. It's simpler # than calling commit() ourselves. "enable.auto.commit": True, # But turn OFF auto-storing, which otherwise marks a message # "ready to commit" the instant poll() hands it to us, before we've # done anything with it. Now nothing is marked ready until we say so. "enable.auto.offset.store": False, }

Why keep enable.auto.commit on?

Because it is simpler than calling commit() ourselves on a schedule. What we remove is only the part that marked messages as ready too early. Now nothing is marked ready until we say so.

Step 2: Store the Offset Only After Success

python
try: process_event(msg) # Only now, after successful processing, do we mark # this offset as safe to commit. consumer.store_offsets(msg)

What does store_offsets() do?

It is the manual version of what auto-store did automatically and too early. Called right after process_event() succeeds, it means the auto-commit timer can only ever commit offsets of messages that really finished. This one change closes the whole danger window from the picture.

The same test with the safe settings gives a different result. With the same crash at offset 4, Kafka's committed offset was 4. After a restart, the consumer starts again at message 4 and processes it properly.

Is this now perfect?

Almost, and it is good to be honest about the last gap. If the consumer crashes after processing a message but before store_offsets() runs, that message will be processed again after the restart. This pattern gives you at-least-once processing: nothing is lost, but a message can be seen twice. So write your processing to be safe when repeated, for example an "update" that gives the same result when done twice.

Is there another way?

Yes. You can turn off auto-commit and call consumer.commit() yourself, even after every message. It gives you exact control, but it is slower, because each commit is a round trip to Kafka. Auto-commit plus store_offsets() is a good balance for most consumers.

Step 3: Handle the Poison Pill

What is a poison pill?

A message that can never be processed successfully. Retrying does not help, because the message itself is the problem. For example, invalid JSON, or a missing field.

Look at the except block:

python
except Exception as e: reason = f"{type(e).__name__}: {e}" print(f"[safe-consumer] POISON PILL at partition {msg.partition()} offset {msg.offset()}: {reason}") send_to_dlq(msg, error_reason=reason) # We still store the offset here. We've "handled" this # message by dead-lettering it, so the consumer should # move on rather than retrying it forever. consumer.store_offsets(msg)

Why do we call store_offsets() here too?

This is deliberate. Look at the picture:

Poison pill: stuck forever vs. dead-lettered and moving onPoison pill: stuck forever vs. dead-lettered and moving on

Suppose we only logged the error and moved on, without storing the offset. We would avoid a crash right now. But on any future restart, the consumer would resume from before this message, and hit it again. And again. Forever. One message that can never succeed becomes a permanent blockage, triggered by every restart.

The fix is to send the bad message to a dead-letter topic. That is a separate topic that only holds messages which could not be processed, kept for someone to inspect later. And then we store the offset, because we have really handled the message, just not in the way we first planned.

python
def send_to_dlq(msg, error_reason: str): """ Sends the bad message to the dead-letter topic and WAITS until Kafka confirms it. If that fails, we raise, so the caller does not store the offset. A message must never be dropped without being saved. """ failures = [] def on_delivery(err, _msg): if err is not None: failures.append(err) dlq_producer.produce( topic=DLQ_TOPIC, key=msg.key(), value=msg.value(), headers={ "error_reason": error_reason.encode("utf-8"), "source_topic": msg.topic().encode("utf-8"), "source_partition": str(msg.partition()).encode("utf-8"), "source_offset": str(msg.offset()).encode("utf-8"), }, on_delivery=on_delivery, ) still_waiting = dlq_producer.flush(10) if failures or still_waiting: raise RuntimeError(f"Could not write to the dead-letter topic: {failures or 'timed out'}")

This reuses what you already know: a producer sending a message, just aimed at a topic for failures instead of successes. Notice three details:

  • Headers. We attach the error reason, and where the message came from (topic, partition, offset). Later, whoever inspects the dead-letter topic can see exactly what went wrong and where.
  • We wait for confirmation. flush() waits for the send to finish, and on_delivery tells us if it failed.
  • If the dead-letter write fails, we raise an error. The offset is not stored, and the consumer stops. That is intentional. If we cannot save the bad message, it is better to stop and try again after a restart than to lose the message.

Not Every Error Is a Poison Pill

One warning about except Exception. It catches everything: bad data, but also bugs in your own code, and temporary problems, like a database that is down for a minute.

  • A bad message will fail every time. Dead-lettering it is right.
  • A temporary problem would succeed if you tried again later. Sending those messages to the dead-letter topic would fill it with perfectly good events.

In a real system, you separate the two: retry (with a pause) for temporary failures, and dead-letter only the permanent ones. Our script keeps things simple and treats every processing error as permanent, so you can see the pattern clearly.


Running It

Step 1: Start Your Cluster

bash
docker compose start docker ps

Step 2: Create the Dead-Letter Topic

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

Step 3: Start the Safe Consumer

Download 5.4_safe_consumer_dlq.py into your chapter-5 folder and run it from your host terminal:

bash
python 5.4_safe_consumer_dlq.py

Step 4: Inject Two Poison Pills

In your broker terminal, start a console producer for clickstream-events:

bash
kafka-console-producer.sh --topic clickstream-events --bootstrap-server kafka-1:29092

Type these two lines, pressing Enter after each. The first is not valid JSON. The second is valid JSON, but the user_id field is missing:

this is not valid json at all
{"event_type": "click"}

Then press Ctrl+C to exit the producer.

Step 5: Watch It Survive

Go back to your consumer terminal. You should see two lines starting with POISON PILL at partition ... offset .... One says JSONDecodeError, and the other says KeyError: 'user_id'.

And most importantly, the consumer keeps running. It did not crash and it did not get stuck.

Now confirm that the bad messages reached the dead-letter topic. In the broker terminal:

bash
kafka-console-consumer.sh --topic clickstream-events-dlq --bootstrap-server kafka-1:29092 --from-beginning --timeout-ms 5000 --property print.headers=true

You should see both messages, with their headers (error_reason, source_topic, source_partition, source_offset). After 5 seconds without new messages, the command stops on its own.

Step 6: Prove the Consumer Moved On

Wait about 10 seconds, so the auto-commit timer has fired. Then run:

bash
kafka-consumer-groups.sh --describe --group recommendation-service --bootstrap-server kafka-1:29092

The LAG column should be 0 for every partition. The bad messages were counted as handled.

Now press Ctrl+C on the consumer, and start it again:

bash
python 5.4_safe_consumer_dlq.py

Do the poison pills come back?

No. This time there are no POISON PILL lines. Their offsets were stored, so the restart begins after them. Compare that with the "stuck forever" picture: this is the difference the store_offsets() call makes.

Step 7: Stop Everything

Press Ctrl+C on the consumer, then:

bash
docker compose stop

Common Consumer Pitfalls

Trusting default auto-commit blindly. You now know exactly why it is risky. It is not a vague warning. It is a specific gap between "received" and "actually processed".

Letting one exception crash the whole poll loop. One bad message should not take down a consumer. Always wrap processing in try/except, and decide on purpose what happens on failure.

Building a poison-pill trap without noticing. Catching an exception and just doing continue, without ever storing that offset, feels like it works. Until the next restart, when the same message fails again. Forever.

Dead-lettering without checking the send worked. If the dead-letter write fails silently and you still store the offset, the message is lost. Always confirm delivery.

Sending temporary failures to the dead-letter topic. A database outage would fill it with good messages. Retry temporary failures instead.

Never monitoring the dead-letter topic. A dead-letter topic that grows unnoticed is just a slower way to lose data. Someone has to watch it.


Checklist

  • Understood what committing an offset means: the next offset to read
  • Understood the three stages: received, stored, committed
  • Understood that auto.offset.reset applies only when there is no valid committed offset, not on every start
  • Know when there is no valid offset: a new group, expired offsets (7 days after the group empties), or an out-of-range offset
  • Know when to choose earliest, latest, or error (and that the default is latest)
  • Understood how enable.auto.offset.store and enable.auto.commit together can silently skip messages
  • Built a consumer that stores an offset only after successful processing, and know that this gives at-least-once processing
  • Understood why a poison pill needs its offset stored too, through a dead-letter path
  • Understood why the dead-letter write must be confirmed, and why temporary errors should not be dead-lettered
  • Ran the safe consumer, injected two poison pills, and watched it survive, and stay past them after a restart

You have now built, scaled, rebalanced, and hardened a real Kafka consumer. Next: a bug is found hours after the fact, and the fixed code has to process that past window again. You will learn the ways Kafka gives you to reprocess the past, and how to choose between them.