You need at-least-once delivery with no data loss in a consumer. How do you order processing and commits, and how do you bound duplicates?
answer
- process THEN commit = no loss
- duplicates inevitable -> idempotent handlers
- commit more often = smaller replay window
- commitAsync loop + commitSync finally
- EOS = transactions + sendOffsetsToTransaction
basics
~20 sDisable auto-commit, process the records first, then commit. Never commit before the work is durable. Because a crash between processing and commit replays the batch, make handlers idempotent. For exactly-once, use Kafka transactions or commit offsets atomically with your output to an external store.
solid answer
~40 sSet enable.auto.commit=false and always order it as poll -> process -> commit. Committing only after processing succeeds guarantees no loss: if you crash before the commit, the uncommitted records are re-delivered. The cost is duplicates on the replay window (everything since the last commit), so the consumer of side effects must be idempotent (dedup keys, upserts, conditional writes). To bound the duplicate window, commit more frequently (smaller batches) — trading throughput for less replay. Use commitAsync in the loop for speed and commitSync in finally/at rebalance for a durable boundary. For true exactly-once into Kafka, use a transactional producer and sendOffsetsToTransaction so the output and the offset commit are atomic; for an external sink, store the offset in the same transaction as the data and seek() to it on assignment.
go deeper
Know the rule: process first, then commit, to avoid losing records.
Explain that this yields at-least-once with possible duplicates and why idempotency is required.
Tune the commit cadence to bound the replay window and combine commitAsync/commitSync correctly, including at rebalance.
Design exactly-once via Kafka transactions (sendOffsetsToTransaction) or external atomic offset storage with seek() on assignment, and reason about throughput vs duplicate trade-offs.
## Goal: at-least-once with no loss 'No data loss' means every record is processed **at least once**; the acceptable cost is occasional **duplicates**. The single rule that delivers this: **commit the offset only after the record's effect is durable.** ## Correct ordering ``` enable.auto.commit=false loop: records = poll() process(records) // make side effects durable commit(lastOffset + 1) // only now ``` If the process step writes to a DB/topic and *then* you commit, a crash at any earlier point leaves the offset uncommitted, so on restart the records are re-delivered and reprocessed — **nothing is lost.** The inverse order (commit then process) gives at-most-once (loss on crash) and is almost never what you want here. ## Why duplicates are unavoidable and how to handle them Between the last successful commit and a crash, all processed records get replayed. You cannot eliminate this with plain commits — you must make processing **idempotent**: - Use a natural/business key and **upsert** instead of insert. - Track processed record ids (e.g., (topic,partition,offset) or a message id) and skip seen ones. - Make external calls idempotent (idempotency keys). ## Bounding the duplicate window The replay size = work since last commit. **Commit more often** (smaller batches, or per-record commit at the extreme) to shrink it — at the cost of commit overhead/throughput. Conversely, large batches with infrequent commits maximize throughput but enlarge the replay window. This is the core tuning knob. ## Commit mechanics for this pattern - **commitAsync()** in the hot loop: fast, no blocking; a dropped async commit just means a slightly larger replay window (still safe for at-least-once). - **commitSync()** in a finally block and inside onPartitionsRevoked: a guaranteed durable boundary at shutdown and before partition handoff. ## Moving to exactly-once Plain offset commits can't be atomic with arbitrary side effects, so 'exactly-once' needs one of: 1. **Kafka transactions**: a transactional producer brackets your output sends and `producer.sendOffsetsToTransaction(offsets, groupMetadata)` so the produced records **and** the consumer offset commit either all commit or all abort — the read-process-write loop is atomic (the basis of Kafka Streams EOS). 2. **External atomic store**: write the processed data and the consumer offset in the **same transaction** of your sink DB; on `onPartitionsAssigned`, `seek()` to the offset read back from that DB instead of relying on __consumer_offsets. This makes offset and output inseparable. ## Common mistakes - Committing inside the processing loop before the side effect is durable. - Relying on auto-commit for no-loss (it can commit ahead of processing). - Forgetting idempotency and being surprised by duplicates after a rebalance/crash.
- Why does committing after processing guarantee no loss but allow duplicates?If you crash before the commit, the offset still points before those records, so they are re-delivered — none are lost. But the records already processed before the crash get processed again on replay, hence duplicates.
- How would you upgrade this from at-least-once to exactly-once when writing back to Kafka?Use a transactional producer: begin a transaction, send the output records, call sendOffsetsToTransaction with the consumer offsets and group metadata, then commit. The output and offset commit are atomic, so no record is processed twice in effect.
saying these in an interview costs you the question
- Saying at-least-once means zero duplicates (it means no loss; duplicates are expected)
- Committing before processing for 'safety' (that's at-most-once, causes loss)
- Claiming plain commitSync gives exactly-once (needs transactions or an external atomic store)
- Ignoring idempotency and assuming the framework dedupes for you