【发布时间】:2017-12-17 02:29:11
【问题描述】:
现在,我正在开发一个使用 Kafka Streams 处理消息的项目。我们的消息由两个标识符组成,一个是用户的标识符,第二个是用户所属的消息列表的标识符。
应用程序的故事是,客户希望通过我的应用程序向他们的用户列表发送推送消息。因此,客户端在向用户发送推送消息之前从外部源加载用户列表,之后开始发送。但客户端也希望随时停止正在进行的发送。
此时我的问题出现了。我想停止发送推送消息,而无需通过其流 api 消耗来自 kafka 的所有消息。我用谷歌搜索将流源的偏移量更改为最新,但找不到合适的方法来完成它。
我怎样才能完成它? 谢谢。
【问题讨论】:
-
不明白你想要达到什么。为什么需要更改偏移量?你能详细说明一下吗?
-
我想跳过未使用的消息,而不处理它们。我认为唯一的方法是改变流消费者的偏移量。但我找不到合适的方法来改变它。我正在为 kafka 流寻找类似 @987654321@ 的方式
标签: java apache-kafka apache-kafka-streams