Learn Labs
7. Reliable Data Delivery

7.5 Using consumers reliably

The reliability-first default is earliest.

The clean division of responsibility:

"data is only available to consumers after it has been committed to Kafka... This means that consumers get data that is guaranteed to be consistent. The only thing consumers are left to do is make sure they keep track of which messages they've read and which messages they haven't."

How consumption works mechanically: "a consumer is fetching a batch of messages, checking the last offset in the batch, and then requesting another batch starting from the last offset received. This guarantees that a Kafka consumer will always get new data in correct order without missing any messages."

⚠️ The single way consumers lose messages

"The main way consumers can lose messages is when committing offsets for events they've read but haven't COMPLETELY PROCESSED yet. This way, when another consumer picks up the work, it will SKIP those messages and they will NEVER GET PROCESSED. This is why paying careful attention to when and how offsets get committed is critical."

COMMITTED MESSAGES vs COMMITTED OFFSETS — don't conflate them

TermMeaning
Committed message"a message that was written to all in-sync replicas and is available to consumers"
Committed offset"offsets the consumer sent to Kafka to acknowledge that it received and processed all the messages in a partition up to this specific offset"

5.1 The four consumer configs that determine reliability

group.id

"if two consumers have the same group ID and subscribe to the same topic, each will be assigned a subset of the partitions... If we need a consumer to see, on its own, EVERY SINGLE MESSAGE in the topics it is subscribed to, it will need a UNIQUE group.id."

auto.offset.reset

"There are only two options here" (in the reliability framing):

earliestlatest
“start from the beginning of the partition whenever it doesn’t have a valid offset.”“start at the end of the partition.”
“can lead to the consumer processing a lot of messages twice, but it guarantees to minimize data loss.”“minimizes duplicate processing but almost certainly leads to some messages getting missed.”

The reliability-first default is earliest. (Ch. 4 adds a third option, none, which throws — often the best choice when neither loss nor mass reprocessing is acceptable and you want a human decision.)

enable.auto.commit — the real decision

"This is a big decision: are we going to let the consumer commit offsets for us based on schedule, or are we planning on committing offsets manually?"

  • ✓ Benefit: “one less thing to worry about.” And a real correctness guarantee, but Only under one condition: “When we do All the processing of consumed records Within the consumer poll loop, then the automatic offset commit Guarantees we will never accidentally commit an offset that we didn't process.”
  • ✗ Drawback: “we have No control over the number of duplicate records the application may process because it was stopped after processing some records but before the automated commit kicked in.”
  • ✗ Hard limit: “When the application has more complex processing, such as Passing records to another thread to process in the background, There is no choice but to use manual offset commit since the automatic commit may commit offsets for records the consumer has Read but perhaps has not processed yet.”

This is a genuinely useful clarification. Autocommit is not sloppy per se — it's correct for the strictly-in-poll-loop case and unsafe the instant processing escapes the poll loop. Handing records to a thread pool or an async client silently converts autocommit into a data-loss mechanism.

auto.commit.interval.ms (default 5 s)

"committing more frequently adds overhead but reduces the number of duplicates that can occur when a consumer stops."

And a fifth, indirect factor

"While not directly related to reliable data processing, it is difficult to consider a consumer reliable if it frequently stops consuming in order to rebalance." (→ Ch. 4: cooperative sticky assignor, static membership, max.poll.records.)

5.2 Explicit offset commits — five rules

Rule 1: Always commit offsets AFTER messages were processed

"If we do all the processing within the poll loop and don't maintain state between poll loops (e.g., for aggregation), this should be easy... If there are additional threads or stateful processing involved, this becomes more complex, especially since THE CONSUMER OBJECT IS NOT THREAD SAFE."

Rule 2: Commit frequency trades performance against duplicates

"Committing has significant performance overhead. It is similar to produce with acks=all, BUT ALL OFFSET COMMITS OF A SINGLE CONSUMER GROUP ARE PRODUCED TO THE SAME BROKER, WHICH CAN BECOME OVERLOADED."

"Committing after every message should only ever be done on very low-throughput topics."

Why commits are more expensive than they look:

  • every commit = a produce with acks=all
  • And every commit for a given group goes to The same broker — the group's coordinator, the __consumer_offsets partition leader
  • ⇒ commits do Not spread across the cluster. One broker absorbs them all.
  • ⇒ “commit after every message” on a high-throughput topic is a self-inflicted hot spot.
Rule 3: Commit the RIGHT offsets at the RIGHT time

"A common pitfall when committing in the middle of the poll loop is accidentally committing the last offset READ when polling and not the offset AFTER the last offset PROCESSED. Remember that it is critical to always commit offsets for messages AFTER they were processed — committing offsets for messages read but not processed can lead to the consumer MISSING messages."

Rule 4: Handle rebalances

"consumer rebalances will happen, and we need to handle them properly. This usually involves committing offsets before partitions are revoked and cleaning any state the application maintains when it is assigned new partitions." (→ Ch. 4's ConsumerRebalanceListener.)

Rule 5: Consumers may need to retry — and Kafka has no per-message ack

The fundamental constraint:

"unlike traditional pub/sub messaging systems, Kafka consumers commit offsets and do NOT 'ack' individual messages. This means that if we failed to process record #30 and succeeded in processing record #31, we should NOT commit offset #31 — this would result in marking as processed all the records up to #31 INCLUDING #30, which is usually not what we want."

… 282930 FAILED31 OK32 …commit 31?30 is implicitly processed — LOST FOREVERdon’t commit?31 will be reprocessed — but you’re stuck

There is no “ack 31, nack 30.” Offsets are a watermark, not a set.

Figure 7.5.2Rule 5: Consumers may need to retry — and Kafka has no per-message ack

The two patterns:

Pattern A — pause and retry in place

  1. Commit the last record we processed successfully.
  2. Store the records that still need processing in a buffer (“so the next poll won’t override them”).
  3. Use consumer.pause() “to ensure that additional polls won’t return data” — note that you must keep polling to stay alive; pause() lets you poll without receiving new records.
  4. Keep trying to process the records.

Pattern B — retry topic (a dead-letter queue)

  1. On a retriable error, write the record to a separate topic and continue.
  2. Handle retries either by a separate consumer group consuming the retry topic, or by one consumer subscribed to both the main topic and the retry topic, “but pause the retry topic between retries.”

“This pattern is similar to the dead-letter-queue system used in many messaging systems.”

Trade-off: Pattern A preserves ordering but head-of-line blocks (nothing progresses while you retry). Pattern B keeps the main stream flowing but abandons ordering for the failed records.

Rule 6: Consumers may need to maintain state

The moving-average example:

"if we want to calculate a moving average, we'll want to update the average every time we poll Kafka for new messages. If our process is restarted, we will need to not just start consuming from the last offset, but we'll also need to RECOVER THE MATCHING MOVING AVERAGE."

The DIY approach and its limitation:

"One way to do this is to write the latest accumulated value to a 'results' topic AT THE SAME TIME the application is committing the offset." → so a starting thread picks up the latest accumulated value and resumes.

"In Chapter 8, we discuss how an application can write results and commit offsets in a SINGLE TRANSACTION. In general, this is a rather complex problem to solve, and we recommend looking at a library like KAFKA STREAMS or FLINK, which provides high-level DSL-like APIs for aggregation, joins, windows, and other complex analytics."

(This is the same "offsets are only half your position" point Ch. 5 made with the shoe-counting example — and it's the motivation for Ch. 8's transactions and Ch. 14's Streams.)


On this page