skip to content

With a consumer using assign() instead of subscribe(), how do you control where consumption starts, since there's no rebalance to trigger position setup?

level: middleimportance: should knowfreq 35%

answer

  1. committed offset -> else auto.offset.reset
  2. no onPartitionsAssigned with assign()
  3. seek / seekToBeginning / seekToEnd right after assign()
  4. offsetsForTimes for time-based replay
  5. seekToBeginning/End are lazy (resolved on poll)

basics

~10 s

After assign(), the consumer starts from the last committed offset for that group.id if one exists; otherwise auto.offset.reset (earliest/latest) applies. You can override explicitly with seek(), seekToBeginning(), or seekToEnd() before polling.

solid answer

~40 s

assign() gives you the partitions but you decide the starting position. The default behavior is the same as subscribe(): if there's a committed offset for this group.id on that partition, resume there; otherwise apply auto.offset.reset (latest by default, or earliest). Because there's no onPartitionsAssigned callback firing (no rebalance), you set positions imperatively right after assign(): call seek(tp, offset) for an exact offset, seekToBeginning(partitions) / seekToEnd(partitions) for the ends, or offsetsForTimes() to resolve an offset by timestamp then seek to it. A common pattern for replay or external offset stores: assign(), then seek() to the offset you loaded from your own store, then poll(). Note seekToBeginning/End and position lookups are lazy/evaluated on the next poll. If you store offsets externally and never want Kafka's committed offsets, you simply always seek() after assign().

go deeper

for a junior

Know seek() lets you choose where to start reading.

for a middle

Explain the committed-offset-then-auto.offset.reset default and the seek family for assign().

for a senior

Use offsetsForTimes and external offset stores; know seek laziness and where to place it without a rebalance callback.

for a principal

Design replay/exactly-once flows with external offset stores and reason about leader-epoch-aware seeks.

## Controlling start position under assign() When you use `assign()`, you immediately own the partitions — but Kafka doesn't decide a starting **offset** for you automatically the way you might expect from a fresh consumer. Here's the full picture. ### Default resolution order On the **first poll()** after `assign()`, for each partition the consumer needs a position. It resolves it as: 1. **Last committed offset** for this `group.id` on that partition (read from `__consumer_offsets`), if one exists. The consumer resumes there. 2. Otherwise, **`auto.offset.reset`**: - `latest` (default) — start at the end (only new records). - `earliest` — start at the beginning of the partition. - `none` — throw `NoOffsetForPartitionException`. This is the **same** logic as a subscribe()-based consumer; assign() doesn't change offset semantics, only who picks the partitions. ### Imperative position control with seek() With subscribe(), you'd typically call `seek()` inside `onPartitionsAssigned`. With assign() **there is no rebalance and no callback**, so you call the seek methods **directly after assign()**, before the first poll that should honor them: - `seek(TopicPartition, long offset)` — start at an exact offset. - `seek(TopicPartition, OffsetAndMetadata)` — exact offset with metadata (leader epoch). - `seekToBeginning(Collection<TopicPartition>)` — earliest available offset (respects retention). - `seekToEnd(Collection<TopicPartition>)` — just past the last record (tail). - `offsetsForTimes(Map<TopicPartition, timestamp>)` — resolve the first offset at/after a timestamp, then `seek()` to it (time-based replay). **Laziness:** `seekToBeginning`/`seekToEnd` don't fetch immediately; the actual offset is resolved on the **next poll()** (or `position()` call). So call them before polling. ### Patterns - **Replay a partition from scratch:** `assign(tp)` → `seekToBeginning([tp])` → poll. - **External offset store (exactly-once-ish):** keep offsets in your DB; `assign(tp)` → `seek(tp, offsetFromDb)` → poll; never rely on Kafka's committed offsets. - **Time travel:** `offsetsForTimes` to convert a wall-clock time into an offset, then seek. ### Gotchas - If you call `poll()` *before* seeking, the consumer may already start consuming from the default-resolved position; seek after that point repositions from the next fetch. - `position(tp)` forces resolution and is handy to log where you'll start. - Committed offsets still depend on `group.id`; with assign() + no group.id you must manage offsets entirely yourself.

  • Where would you put seek() logic for a subscribe()-based consumer versus an assign()-based one?
    For subscribe(), put seek() inside the onPartitionsAssigned rebalance callback (it fires when you gain partitions). For assign(), there's no rebalance/callback, so call seek() imperatively right after assign() and before the first relevant poll().

saying these in an interview costs you the question

  • Saying assign() always starts at offset 0 / earliest by default — it follows committed offset then auto.offset.reset (default latest).
  • Expecting onPartitionsAssigned to fire under assign().
  • Believing seekToBeginning/End take effect immediately rather than on the next poll.
  • Thinking you cannot use committed offsets with assign() — you can if group.id is set.

context