skip to content

questions

6

In an event streaming platform like Kafka, what is a consumer offset, and how does it allow a consumer to replay events it has already processed?

level: juniorimportance: must knowfreq 65%

answer

  1. offset = position, not time
  2. log ≠ queue: data survives read
  3. commit before vs after = at-most vs at-least once
  4. seek() to replay
  5. __consumer_offsets topic

basics

~10 s

An offset is a bookmark marking how far a consumer has read in a partition. Since the log isn't deleted after reading, moving the bookmark backward lets the consumer re-read and reprocess old events.

solid answer

~40 s

Each partition in the log is an append-only, ordered sequence of records, and every record has a monotonically increasing offset - its position in that sequence. A consumer tracks its progress by committing the offset of the last record it successfully processed, usually to a special internal topic. Unlike a traditional queue where a message is removed once acknowledged, the log keeps records around for a retention period (time or size based) or indefinitely if compacted. Because the data is still there, a consumer can 'seek' to any offset - zero, a specific number, or a timestamp - and start consuming from that point again. This is how you replay a day's worth of events to backfill a new service, rebuild a corrupted cache, or reprocess after fixing a bug in the consumer logic.

go deeper

for a junior

Should explain offset as a per-partition position and know that committing lets a consumer resume where it left off; doesn't need to know the internal __consumer_offsets mechanics.

for a middle

Should articulate at-least-once vs at-most-once depending on commit timing, and know how to manually reset offsets to replay a range.

for a senior

Should reason about rebalance interactions, offset-commit strategies (manual/sync vs async, per-record vs batch), and design idempotent consumers so replay is safe.

for a principal

Should weigh replay-driven reprocessing against retention cost and downstream blast radius, and design offset-commit/transactional-write patterns that give exactly-once guarantees where the business genuinely requires them.

## What an offset is An event streaming platform such as Apache Kafka organizes data into topics, and each topic is split into one or more partitions. A **partition** is the fundamental unit of storage: an append-only, immutable sequence of records written to disk in the order they arrive. Every record in a partition is assigned a strictly increasing integer called its **offset**, starting at zero. The offset is nothing more than the record's position within that partition: - it has **no meaning across partitions**; - it has **no relationship to wall-clock time** (though records also carry a timestamp separately). ## How a consumer uses one 1. **Read.** A consumer reads a partition by asking the broker for records starting at a given offset, and the broker streams back records in order until the consumer catches up to the newest write (the 'log end offset' or high-water mark). 2. **Commit.** As the consumer processes records, it periodically commits its position - the offset of the next record it wants to read - to a durable store. In Kafka this is itself a compacted topic called `__consumer_offsets`, keyed by (consumer group, topic, partition), so the broker cluster survives consumer restarts without losing anyone's place. 3. **Resume.** On startup or after a **rebalance** (when partitions are reassigned among the members of a consumer group), each consumer looks up its last committed offset and resumes exactly there. ## Why the design exists This design exists because streaming platforms deliberately decouple producers from consumers in both time and multiplicity, which a traditional point-to-point queue does not. | Point-to-point queue | Log-based system | |---|---| | In a queue, once a message is acknowledged it is gone - if you need a second consumer, or need to recompute something after a bug fix, the data no longer exists. | A log-based system instead treats position tracking and data storage as two separate concerns: the log keeps data for as long as its retention policy says, while each consumer group independently tracks its own read position. | Because these are decoupled, a consumer - or a brand new consumer group that has never run before - can seek to any valid offset, including the very beginning, and replay history. This is what makes the log 'replayable': the ability to move a pointer backward and reprocess data you have already seen, without asking the producer to resend anything. ## The trade-off The trade-off is that you are trading storage cost and some operational complexity for this replay capability. Retaining weeks of high-volume data is not free - it costs disk, and if replication is on, it costs it several times over. You also inherit delivery-semantics decisions that a simple queue hides from you: - **At-most-once.** If you commit the offset before processing a record and the consumer crashes mid-processing, that record is silently skipped. - **At-least-once, the common default.** If you commit after processing, a crash between processing and committing causes the record to be reprocessed on restart. - **Exactly-once.** Getting exactly-once requires either idempotent downstream writes or a transactional mechanism that atomically ties the offset commit to the output write - genuinely harder to build and reason about. ## Failure modes Failure modes in production usually show up around three edges. 1. **First, 'offset out of range' errors.** These happen when a consumer's committed offset falls outside the currently retained range - typically because the consumer was down longer than the retention window, so the offset it wants no longer exists; the consumer then has to decide (via an auto.offset.reset-style policy) whether to jump to the earliest or latest available offset, silently skipping or duplicating a chunk of data. 2. **Second, uncommitted offset lag.** After a crash or slow consumer this causes a burst of reprocessing on restart or after a rebalance, which can overwhelm downstream systems if processing isn't idempotent. 3. **Third, committing offsets too eagerly** (e.g., auto-commit on a timer, decoupled from actual processing completion) silently drops records whenever the process dies between the commit and the completion of side effects like a database write. ## Reset and replay in practice A concrete example: an analytics pipeline discovers that its event-enrichment consumer had a bug for the last six hours, producing wrong output. Because the source topic retains seven days of raw events, the on-call engineer fixes the bug, resets the consumer group's offsets back to the timestamp six hours ago, and lets it reprocess - correcting the output without needing the producers to resend anything. That operational move - 'reset and replay' - is only possible because offsets and log retention are decoupled from what has already been consumed.

  • What happens if a consumer commits its offset right after reading a record but before finishing processing it, and then crashes?
    The record is lost from that consumer's perspective - it will never be processed because the committed offset already points past it, and on restart the consumer resumes after it. This is at-most-once delivery, and it's why most systems commit after processing completes rather than before.
  • Why can't a consumer just replay from offset zero every time it restarts?
    It could, but it would reprocess the entire retained history every restart, which is wasteful and, for stateful aggregations, would produce duplicate or inconsistent results unless every downstream effect is fully idempotent. Committing offsets lets a consumer resume incrementally instead of always starting over.
  • What causes an 'offset out of range' error and how do systems typically recover from it?
    It happens when a consumer's stored offset points to data that has already been deleted by the topic's retention policy, usually because the consumer was offline longer than the retention window. Recovery is policy-driven: reset to the earliest available offset (reprocess what's left) or the latest (skip the gap entirely), and both have data-loss or duplication implications the team must accept.

Think of the log as a very long shelf of numbered boxes that never get removed once placed. Your 'offset' is just a sticky note on the shelf saying 'I've unpacked up through box #482' - you can always move the sticky note back to an earlier box and unpack it again, because nobody threw the boxes away.

saying these in an interview costs you the question

  • Says offsets are timestamps
  • Thinks acknowledging a message deletes it from the log
  • Doesn't know the difference between committing before vs after processing
  • Believes replay requires the producer to resend data
  • Assumes all consumers in different groups share the same offset

context

open as a page

What does log compaction do to a Kafka-style topic, and how does it differ from simple time/size-based retention?

level: middleimportance: must knowfreq 60%

basics

~20 s

Compaction keeps only the newest record for each key and throws away older ones, instead of deleting everything older than some age. It's like keeping only the latest version of each row in a table, not a full history.

open as a page

In stream processing, what's the difference between tumbling, sliding, and session windows, and when would you reach for each?

level: middleimportance: must knowfreq 75%

basics

~20 s

A tumbling window is a fixed time slice that never overlaps the next one, like every 5 minutes. A sliding window overlaps and updates continuously, like a 'last 5 minutes' view refreshed often. A session window has no fixed size - it groups events until there's a gap of inactivity.

open as a page

How do stateful stream operators (like running aggregations or joins) maintain state across a potentially unbounded stream, and how do they recover that state after a crash or task reassignment?

level: seniorimportance: should knowfreq 55%

basics

~20 s

Stateful operators keep a local, on-disk 'memory' (like a running total) next to each processing task, and also write every change to a durable, replayable backup log. If the task dies or moves machines, it rebuilds its memory by replaying that backup log.

open as a page

In windowed stream processing, what is a watermark, and how does it let a system decide when a time window is 'done' despite events arriving out of order?

level: seniorimportance: should knowfreq 45%

basics

~20 s

A watermark is the stream processor's best guess of 'we won't see any more events older than this timestamp.' Once the watermark passes a window's end time, the system treats that window as closed and emits its result, accepting it might occasionally be wrong if a late straggler shows up after.

open as a page

At a system-design level, how do stream-processing frameworks achieve 'exactly-once' processing guarantees for stateful aggregations, tying together offset commits, state updates, and output writes - and what breaks that guarantee?

level: principalimportance: nice to knowfreq 30%

basics

~20 s

The system bundles 'I read this event,' 'I updated my running total,' and 'I wrote the result' into one all-or-nothing operation using transactions, so a crash can never leave things half-done - either all three happened or none did.

open as a page