What is the difference between position(), committed(), and beginningOffsets()/endOffsets()?
answer
- position = next record to read (in-memory)
- committed = last persisted to __consumer_offsets
- endOffsets = high watermark = one past last record
- beginningOffsets = oldest retained
- lag = endOffsets - committed
basics
~10 sposition() is where the consumer will read next (in-memory). committed() is the last offset persisted to Kafka for the group. beginningOffsets()/endOffsets() are the oldest and newest (end) offsets currently in each partition.
solid answer
~40 sposition(TopicPartition) returns the consumer's current FETCH position — the offset of the next record poll() will return, held in memory and advanced as you consume; it can block to fetch the position if unknown. committed(partitions) returns the last COMMITTED offset stored in __consumer_offsets for this group, i.e. the durable progress that survives restart. These differ whenever you've consumed past your last commit. beginningOffsets(partitions) returns the earliest offset still retained; endOffsets(partitions) returns the END offset = high watermark = one past the last record (the offset the NEXT produced record will get). A frequent use is lag = endOffsets - position (or - committed). beginningOffsets/endOffsets and committed are blocking broker calls; position is mostly local but may fetch.
go deeper
Know position is where you'll read next and committed is the saved progress.
Distinguish all four offsets and compute lag from endOffsets minus committed/position.
Reason about position vs committed divergence for at-least-once duplicate windows and avoid blocking calls on the hot path.
Build lag-monitoring and bounded-replay tooling using these primitives, accounting for offset gaps from compaction/transactions.
## Four offsets, four meanings For any partition there are several distinct offsets people confuse: - **`position(tp)` — the fetch/current position.** The offset of the **next** record `poll()` will hand you. It lives in the consumer's memory and advances automatically as records are consumed (or jumps when you `seek`). If the position isn't yet known (e.g. just assigned), calling it can trigger a blocking fetch / apply `auto.offset.reset`. - **`committed(tp)` — the durable committed offset.** The last offset this **consumer group** persisted to the internal `__consumer_offsets` topic. This is what a new consumer (or this one after restart/rebalance) resumes from. It is a **blocking** call to the group coordinator. Convention: a committed offset of N means 'records < N are done; resume at N'. - **`beginningOffsets(tps)` — the log start.** The **earliest** offset still retained in each partition (records below it have been deleted by retention/compaction). Blocking broker call. - **`endOffsets(tps)` — the log end / high watermark.** **One past** the last record — the offset the *next* produced record will receive. It is NOT a record you can read yet. Blocking broker call. ## Why position != committed They diverge any time you've polled records but not yet committed. Example: you poll up to offset 950 (so `position` = 951) but last committed 900 (`committed` = 900). A crash now would resume at 900 and reprocess 900-950 (at-least-once). Tracking both is how you reason about duplicate processing and lag. ## Computing consumer lag Lag is how far behind the consumer is: - **Lag = endOffset - committedOffset** (durable lag, what monitoring usually shows), or - **Lag = endOffset - position** (live, in-flight lag). Since `endOffsets` is one-past-the-last-record and offsets can have gaps (compaction, transaction markers), lag is an approximation of 'records remaining', not an exact count. ## Blocking vs local - `position()` is mostly local but may block to resolve an unknown position. - `committed()`, `beginningOffsets()`, `endOffsets()` all make **blocking** broker/coordinator round-trips (with timeouts), so don't call them on a hot path per record. ## Practical pairing with seek `beginningOffsets`/`endOffsets` are often read to bound a `seek`: e.g. to start 100 records from the end, `seek(tp, endOffsets(tp) - 100)`, clamped to `beginningOffsets`.
- endOffsets returns 1000 for a partition. Can poll() return a record at offset 1000?Not yet. endOffsets is the high watermark — one past the last existing record. Offset 1000 is where the next produced record will go; currently the last readable record is 999.
- How would you compute consumer lag for a group?Lag per partition = endOffsets(tp) - committed offset (or - position for live lag), summed across partitions. It's the standard endOffset-minus-progress calculation; tools like kafka-consumer-groups --describe report it the same way.
saying these in an interview costs you the question
- Saying position() returns the last consumed offset (it's the NEXT one)
- Treating endOffsets as a readable record rather than one-past-the-last
- Conflating position with committed
- Assuming these are all cheap local calls (committed/beginning/endOffsets block)