【问题标题】:How to get last 5 days messages from Kafka using Java如何使用 Java 从 Kafka 获取最近 5 天的消息
【发布时间】:2017-09-05 14:08:53
【问题描述】:

我已将 Kafka 中某个主题的 TTL 设置为 7 天,我从 Kafka 获取数据并将其存储在数据库中,但从过去 5 天开始,我的数据库服务器已关闭,现在我必须获取最近 5 天的消息来自Kafka 并将它们存储在数据库中 注意:从过去 5 天开始,Kafka 没有问题。

【问题讨论】:

  • 您需要借助偏移值进行消费。例如,如果您上次读取的偏移量为 100,那么您需要从偏移量 101 开始使用它。
  • 如何在 Java 中使用这个偏移量概念以及我如何知道存储消息的最后一个偏移量值,因为我没有存储任何偏移量值

标签: java apache-kafka kafka-consumer-api


【解决方案1】:

首先调用 consumer.partitionsFor() 方法来获取主题的分区

https://kafka.apache.org/0110/javadoc/org/apache/kafka/clients/consumer/KafkaConsumer.html#partitionsFor(java.lang.String)

然后调用 consumer.offsetsForTimes() 获取每个分区的偏移量,以获取 5 天前成功处理最后一条消息时的时间戳。

https://kafka.apache.org/0110/javadoc/org/apache/kafka/clients/consumer/KafkaConsumer.html#offsetsForTimes(java.util.Map)

然后调用consumer.seek()来定位当前消费者在那个时间点的偏移量,继续调用poll(),像往常一样处理消息。

https://kafka.apache.org/0110/javadoc/org/apache/kafka/clients/consumer/KafkaConsumer.html#seek(org.apache.kafka.common.TopicPartition,%20long)

【讨论】:

    【解决方案2】:

    在上一个不错的答案中,我将添加该调用 partitionsFor 方法来获取您主题的分区,然后按照 @Hans 所说的那样做。

    【讨论】:

    • 谢谢。我更新了答案以包含正确的第一步。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-11-25
    • 2013-12-23
    • 1970-01-01
    • 2012-02-08
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多