【问题标题】:Kafka Streaming application reading only latest message after connection with KafkaKafka Streaming 应用程序在与 Kafka 连接后仅读取最新消息
【发布时间】:2018-12-12 01:20:43
【问题描述】:

我们正在使用 Kafka Streaming 库为 Kafka 主题上的传入消息构建实时通知系统,因此当流式应用程序运行时,它会实时处理主题中的所有传入消息并在遇到时发送通知某种预定义的传入消息。

如果流媒体应用程序关闭并重新启动,我们需要仅处理在流媒体应用程序初始化后到达的最近消息。这是为了避免处理流式应用程序未运行或关闭时未处理的旧记录。默认情况下,流式应用程序开始处理自上次提交偏移量以来的旧消息。 Kafka Streaming App 中是否有任何设置允许仅处理最新消息?

【问题讨论】:

  • 根据我的理解,默认情况下,kafka 消费者从消费者组偏移量的位置(最后提交的偏移量)选择。它不依赖于时间。所以在我的情况下,我可能需要拒绝时间戳早于消费者开始时间的记录。当然,这意味着我需要有带有相关时间戳的消息。根据 stackoverflow.com/questions/39514167/… 从 0.10.0 开始,可以将时间戳与消息相关联。
  • 解决方案在这个帖子中有详细说明:stackoverflow.com/questions/45075147/…
  • 如果您的群组 ID 是动态的,那么您每次启动流媒体应用时都会处理最新消息。

标签: apache-kafka streaming


【解决方案1】:

KafkaConsumer 的 'auto.offset.reset' 默认值为 'latest' 但您想使用 KafkaStreams,默认为“最早” 参考:https://github.com/apache/kafka/blob/trunk/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java#L634

因此, 如果设置 auto.offset.reset 是“最新的”,它将是你想要的。

【讨论】:

  • 实际上我尝试了“最新”,但是只有在没有消费者偏移量的情况下才有效。否则,如果消费者偏移量已经 > 0,则 auto.offset.reset 无效。客户端将从下一个偏移量中选择消息。
【解决方案2】:

你的假设是正确的。即使您将auto.offset.reset 设置为latest,您的应用程序也已经有了消费者偏移量。

因此,您必须使用带有这些选项 --reset-offsets --to-latest --executekafka-consumer-groups 命令将偏移量重置为最新。

检查不同的重置方案,您甚至可以从文件等重置为特定的日期时间,或按时间段。

【讨论】:

    猜你喜欢
    • 2017-12-01
    • 1970-01-01
    • 2018-01-09
    • 2015-02-04
    • 2017-04-03
    • 2021-05-08
    • 1970-01-01
    • 2019-04-27
    • 1970-01-01
    相关资源
    最近更新 更多