Distinguish a consumer's committed offset from its current position. How do poll, position, and commit interact?
answer
- position = next poll cursor, in-memory
- committed = saved to __consumer_offsets
- committed = last processed + 1
- process-then-commit = at-least-once
- no commit -> auto.offset.reset
basics
~20 sThe current position is the offset of the next record the consumer will fetch — it lives in memory and moves forward as you poll. The committed offset is the position durably saved to Kafka, used to resume after a restart or rebalance. They can differ.
solid answer
~40 sA consumer's **position** for a partition is an in-memory cursor: the offset of the **next** record `poll()` will return. It advances as you consume. The **committed offset** is a position **persisted** to Kafka (the __consumer_offsets topic) via `commitSync`/`commitAsync` or auto-commit; it is the resume point used after a restart or partition reassignment. Crucially, the committed offset stores the offset of the *next* record to consume (last processed + 1). Position and committed offset are independent: you can have consumed up to position 500 in memory while only committing 400, meaning a crash would reprocess 400–499 (at-least-once). On rebalance/startup, the consumer seeks to the committed offset; if none exists, `auto.offset.reset` decides earliest vs latest. `position()` reads the live cursor; `committed()` queries the stored value.
go deeper
Know there's a live 'where I am now' (position) and a saved 'where I'll resume' (committed offset).
Explain committed = next offset, the __consumer_offsets store, and process-then-commit for at-least-once.
Discuss auto-commit pitfalls, seek-based reprocessing, and OFFSET_OUT_OF_RANGE when committed data ages out.
Design exactly-once/idempotent pipelines around the position/committed split and per-group offset storage.
## Two different cursors When a consumer reads a partition there are two distinct notions of 'where am I': ### 1. Current position (live, in-memory) The **position** is the offset of the **next record** the consumer will return from the next `poll()`. It is held in the consumer client's memory. Mechanics: - After `poll()` returns a batch ending at offset 99, the position becomes 100. - You can move it manually with `seek(partition, offset)`, `seekToBeginning`, or `seekToEnd`. - `consumer.position(tp)` returns this value. - It is **not durable** — it lives only in this client instance and is lost on crash. ### 2. Committed offset (durable, stored in Kafka) The **committed offset** is a position that has been **saved** so consumption can resume later. Kafka stores it in the internal compacted topic **`__consumer_offsets`**, keyed by `(group, topic, partition)`. Mechanics: - Written via `commitSync()`/`commitAsync()` (manual) or by **auto-commit** (`enable.auto.commit=true`, every `auto.commit.interval.ms`). - By convention the committed value is **last-processed-offset + 1** — i.e. the offset to *start from* next time, not the last record you handled. (When committing explicitly with an `OffsetAndMetadata`, you must pass the next offset.) - `consumer.committed(tp)` queries this stored value. ## How poll / position / commit interact 1. On assignment, the consumer fetches the **committed offset** for each partition and sets its **position** there. If there's no committed offset, it applies **`auto.offset.reset`** (`earliest` -> log-start-offset, `latest` -> LEO/HW, `none` -> exception). 2. `poll()` fetches records starting at the position and advances the position past them. 3. `commit*()` writes the current position (or an offset you specify) to `__consumer_offsets`. 4. On the next restart or rebalance, step 1 repeats using whatever was last committed. ## Why they diverge — delivery semantics - **At-least-once (default-ish):** process records, then commit. If you crash after processing but before committing, you reprocess from the committed offset -> duplicates possible. - **At-most-once:** commit before processing. A crash after commit but before processing -> records skipped. - **Auto-commit caveat:** auto-commit commits the position reached by `poll()`, which may be ahead of what you've actually finished processing, silently risking data loss on crash. Many teams disable it and commit explicitly after processing. ## Edge cases - Position can be **ahead** of the committed offset (in-flight, uncommitted work). - Position can be moved **behind** the committed offset with `seek()` to reprocess. - A committed offset that is now **below the log-start-offset** (data aged out) triggers `OFFSET_OUT_OF_RANGE` -> `auto.offset.reset` kicks in. - The committed offset is per **consumer group**: different groups track the same partition independently. - `position()` may need a fetch to the broker if the position isn't yet known; `committed()` always queries the coordinator.
- If you commit offset 400 but your in-memory position is 500, then crash, what happens on restart?The consumer resumes at the committed offset 400, so records 400–499 are re-delivered and reprocessed. This is at-least-once behavior; processing must be idempotent to avoid side-effect duplication.
- Why do many teams disable enable.auto.commit?Auto-commit periodically commits the position reached by poll(), which can be ahead of records you've actually finished processing. A crash then skips unprocessed records (data loss). Manual commit after processing gives precise at-least-once control.
saying these in an interview costs you the question
- Saying the committed offset stores the last processed offset rather than last + 1.
- Conflating position and committed offset as the same thing.
- Claiming auto-commit guarantees no data loss.
- Thinking committed offsets are per-consumer-instance rather than per-group.