Learn Labs
4. Kafka Consumers: Reading Data from Kafka

4.9 Consuming from specific offsets

① Map every partition assigned to this consumer (consumer.assignment()) to the target timestamp.

consumer.seekToBeginning(Collection<TopicPartition> tp);  // replay all
consumer.seekToEnd(Collection<TopicPartition> tp);        // skip to now
consumer.seek(TopicPartition tp, long offset);            // exact position

Two motivating use cases from the book:

  • "a time-sensitive application could skip ahead a few records when falling behind"
  • "a consumer that writes data to a file could be reset back to a specific point in time in order to recover data if the file was lost"

Seek by timestamp — the recovery pattern

Long oneHourEarlier = Instant.now().atZone(ZoneId.systemDefault())
          .minusHours(1).toEpochSecond();

Map<TopicPartition, Long> partitionTimestampMap = consumer.assignment()
        .stream()
        .collect(Collectors.toMap(tp -> tp, tp -> oneHourEarlier));   // ①

Map<TopicPartition, OffsetAndTimestamp> offsetMap
        = consumer.offsetsForTimes(partitionTimestampMap);            // ②

for (Map.Entry<TopicPartition,OffsetAndTimestamp> entry: offsetMap.entrySet()) {
    consumer.seek(entry.getKey(), entry.getValue().offset());         // ③
}

① Map every partition assigned to this consumer (consumer.assignment()) to the target timestamp. ② offsetsForTimes() — "sends a request to the broker where a timestamp index is used to return the relevant offsets." ③ seek() each partition to the returned offset.

(Cross-reference Ch. 2: the resolution of timestamp→offset lookups depends on log segment size, because Kafka locates the segment that was being written at that time and returns the offset at its beginning. Big segments on a low-volume topic = imprecise seeks.)


On this page