【发布时间】: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 是动态的,那么您每次启动流媒体应用时都会处理最新消息。