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 positionTwo 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.)