Explain seek(), seekToBeginning(), and seekToEnd() — how do they work and when does the new position take effect?
answer
- must be assigned first
- lazy — resolves on next poll()
- seek overrides auto.offset.reset
- seek does NOT commit
- rebalance listener onPartitionsAssigned pattern
basics
~10 sThese methods manually set where the consumer reads next. seek() jumps to a specific offset; seekToBeginning() to the oldest; seekToEnd() to the newest. The change takes effect on the next poll(), and overrides auto.offset.reset.
solid answer
~40 sseek(TopicPartition, offset) sets the fetch position for an assigned partition to an exact offset; seekToBeginning(collection) and seekToEnd(collection) move to the earliest/latest offsets for the given partitions. Key facts: (1) the partition must be ASSIGNED first — you call these after a poll() or after partition assignment completes (otherwise IllegalStateException or no effect). (2) They are lazy: seekToBeginning/seekToEnd don't issue an RPC immediately; the actual beginning/end is resolved on the next poll(). (3) seek overrides auto.offset.reset — once you seek, the reset policy is moot for that partition. (4) seek does not commit; it only changes where you read, so a later commit will persist the new position. A common pattern is a ConsumerRebalanceListener.onPartitionsAssigned that calls seek to restore externally-stored offsets.
go deeper
Know the three methods set the read position and take effect on the next poll().
Explain the assignment precondition, lazy resolution, and that seek overrides auto.offset.reset but doesn't commit.
Use seek inside a rebalance listener to implement externally-stored offsets and poison-pill skipping.
Design exactly-once / external-offset architectures around seek + commit semantics and reason about rebalance timing.
## The fetch position Every assigned partition has a **fetch position** (also called the *current position*) — the offset of the next record `poll()` will return. Normally this advances automatically as you consume. The seek family lets you set it manually. ## The three methods - **`seek(TopicPartition partition, long offset)`** — set the position to an exact offset. The next `poll()` fetches starting at that offset. There is also an overload taking `OffsetAndMetadata`. - **`seekToBeginning(Collection<TopicPartition>)`** — set the position to the **first available** offset (oldest retained record). An empty collection means *all currently assigned* partitions. - **`seekToEnd(Collection<TopicPartition>)`** — set the position to the **end** offset (one past the last record / the high watermark for that partition), so you only read records produced after this point. ## When does it take effect? (the lazy evaluation gotcha) These calls do **not** synchronously talk to the broker. `seekToBeginning`/`seekToEnd` just mark the partition as 'needs reset to beginning/end'; the actual offsets are resolved by an RPC on the **next `poll()`**. Consequently: - Calling `position()` right after `seekToEnd()` may trigger that resolution (it can block), but you should not assume the seek 'happened' the instant you called it. - The new position is only used the next time you `poll()`. ## Preconditions The partition must be **assigned** to this consumer: - With `subscribe()` (group management), assignment happens during a rebalance. You can only seek partitions you actually own — the natural place is inside `ConsumerRebalanceListener.onPartitionsAssigned`, or after a `poll()` returns and you've confirmed assignment. - With `assign()` (manual assignment, no group coordination), you own the partitions immediately and can seek right away. Seeking a partition you don't own throws `IllegalStateException`. ## Relationship to auto.offset.reset and commits - **seek overrides auto.offset.reset.** Once you explicitly position a partition, the reset policy is irrelevant for it. - **seek does NOT commit.** It changes only the in-memory read position. If you want that position to survive a restart/rebalance under group management, you must commit (or store it externally and re-seek). This is the foundation of the 'store offsets in your own datastore' pattern: on `onPartitionsAssigned`, look up your stored offset and `seek()` to it, ignoring `__consumer_offsets`. ## Typical uses - Replaying from a known offset after a bug fix. - Skipping a poison record by `seek(tp, position + 1)`. - Restoring externally-managed offsets in a rebalance listener. - Reprocessing all history via `seekToBeginning`.
- You call seekToBeginning() but the consumer keeps reading from where it was. Why might that be?Most likely you seeked before the partition was assigned (no effect / IllegalStateException), or you set the position but auto-commit/another seek overrode it, or you never called poll() again so the lazy reset hasn't resolved. With subscribe(), seek inside onPartitionsAssigned or after assignment is established.
- Does seek() persist across a consumer restart?No. seek only changes the in-memory position. To persist you must commit the new offset, or store it externally and re-seek on assignment. Otherwise on restart the consumer resumes from the last committed offset (or auto.offset.reset if none).
saying these in an interview costs you the question
- Saying seek() commits the offset
- Claiming seekToBeginning() immediately contacts the broker and reads instantly
- Trying to seek a partition that isn't assigned yet
- Believing auto.offset.reset still applies after an explicit seek