skip to content

How do you make a consumer start reading from a specific point in time, and what does offsetsForTimes return at the edges?

level: seniorimportance: should knowfreq 48%

answer

  1. timestamp(ms) -> first offset with ts >= target
  2. returns OffsetAndTimestamp, then seek()
  3. null when no record at/after target -> handle it
  4. message.timestamp.type: CreateTime vs LogAppendTime
  5. uses broker time index; blocking call

basics

~20 s

Use offsetsForTimes(map of partition to timestamp) to look up the first offset whose record timestamp is >= the given time, then seek() to it. If no record has a timestamp at or after the target, it returns null for that partition.

solid answer

~40 s

offsetsForTimes(Map<TopicPartition, Long timestampMs>) asks brokers, per partition, for the earliest offset whose record timestamp is >= the requested timestamp, returning an OffsetAndTimestamp (offset + that record's timestamp). You then seek() each partition to the returned offset to begin time-based consumption. Edge cases: (1) if no message has a timestamp at or after the target (target is in the future or beyond the log), the value is NULL for that partition — you must handle it (commonly fall back to seekToEnd). (2) The lookup uses the message timestamp, governed by the topic's message.timestamp.type — CreateTime (producer-set) or LogAppendTime (broker-set). With CreateTime, out-of-order producer timestamps can make results imprecise. (3) It relies on the broker's time index. It's a blocking call subject to default.api.timeout.ms / the timeout overload.

go deeper

for a junior

Know that offsetsForTimes maps a timestamp to a starting offset and you then seek to it.

for a middle

Explain the >= semantics, OffsetAndTimestamp result, and that you must seek afterward.

for a senior

Handle the null edge case explicitly and reason about CreateTime vs LogAppendTime accuracy and blocking/timeout behavior.

for a principal

Design timestamp-driven replay/backfill flows accounting for timestamp type, time-index granularity, producer clock skew, and per-partition fallback policy.

## The goal: 'start from 9am yesterday' Offsets are integers, not times. To consume from a wall-clock point you must translate a timestamp into an offset. Kafka maintains, alongside the log, a **time index** mapping timestamps to offsets, which makes this lookup efficient. ## The API `offsetsForTimes(Map<TopicPartition, Long> timestampsToSearch)` takes, per partition, a Unix epoch **millisecond** timestamp. For each partition the broker returns an **`OffsetAndTimestamp`** containing: - `offset()` — the **earliest** offset whose record timestamp is **>= the requested timestamp**. - `timestamp()` — that record's actual timestamp (>= what you asked for). You then `seek(partition, result.offset())` for each partition to begin reading from that time. Typical flow: build the timestamp map for assigned partitions, call `offsetsForTimes`, then loop and `seek` (handling nulls). ## Which timestamp is used? The lookup compares against each record's **message timestamp**, whose meaning depends on the topic config **`message.timestamp.type`**: - **`CreateTime`** (default) — the timestamp the **producer** set (usually when the record was created). If producers send out-of-order or skewed timestamps, the time index and thus `offsetsForTimes` can be imprecise. - **`LogAppendTime`** — the **broker** stamps the record on append, giving monotonic, reliable times at the cost of losing the original event time. ## The null edge case (most-asked) If **no** record in a partition has a timestamp **>= the target** — e.g. you asked for a future time, or a time after the last produced record — `offsetsForTimes` returns **`null`** for that partition (the map value is null). You MUST handle this; calling `.offset()` on a null throws `NullPointerException`. A common policy is: on null, `seekToEnd` that partition (consume only new data), or skip it. Conversely, if the target is **before** the earliest retained record, you get the **first available** offset (effectively the beginning). ## Operational notes - It is a **blocking** network call to the brokers; there's an overload taking a `Duration` timeout, otherwise it uses `default.api.timeout.ms`. - The partitions don't need to be the consumer's assigned partitions for the lookup itself, but to `seek` them they must be assigned. - Precision is at the granularity of the time index and the actual record timestamps; you won't land 'exactly' at a timestamp, but at the first record at/after it.

  • offsetsForTimes returns null for one partition. What does that mean and how do you handle it?
    It means no record in that partition has a timestamp at or after your target (target is in the future or past the log end). You must not dereference it; a common handling is to seekToEnd that partition so you read only future records, or skip/log it.
  • Why might time-based lookups be imprecise on a topic using CreateTime?
    With CreateTime the producer sets the timestamp, so out-of-order or clock-skewed producers can place records whose timestamps don't increase monotonically with offset, making the time index — and thus the offset returned — only approximate. LogAppendTime avoids this by stamping at the broker.

saying these in an interview costs you the question

  • Saying it returns the offset with the closest timestamp (it's the first >= target)
  • Forgetting the null result and dereferencing it (NPE)
  • Assuming timestamps are seconds, not milliseconds
  • Ignoring message.timestamp.type's effect on accuracy
  • Thinking offsetsForTimes by itself moves the consumer (you must seek to the returned offset)

context