apache / apache/beam

Support for reading Kafka topics from any startReadTime in Java

Open
#21,610 4 comments 0 reactions 0 assignees View on GitHub
io java kafka new feature P2
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

[https://github.com/apache/beam/blob/fd8546355523f67eaddc22249606fdb982fe4938/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/ConsumerSpEL.java#L180-L198](https://github.com/apache/beam/blob/fd8546355523f67eaddc22249606fdb982fe4938/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/ConsumerSpEL.java#L180-L198)

 

Right now the 'startReadTime' config for KafkaIO.Read looks up an offset in every topic partition that is newer or equal to that timestamp. The problem is that if we use a timestamp that is so new, that we don't have any newer/equal message in the partition. In that case the code fails with an exception. Meanwhile in certain cases it makes no sense as we could actually make it work.

If we don't get an offset from calling `consumer.offsetsForTimes`, we should call `endOffsets`, and use the returned offset **** 1. That is actually the offset we will have to read next time.

Even if `endOffsets` can't return an offset we could use 0 as the offset to read from.

 

Am I missing something here? Is it okay to contribute this?

Imported from Jira [BEAM-14518](https://issues.apache.org/jira/browse/BEAM-14518). Original Jira may contain additional context.
Reported by: bnemeth.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.