【发布时间】:2019-06-17 04:35:42
【问题描述】:
主题名称:testTopic
主题的消息总数:1
分区:8
消费者组名称:Consumer1
消费者语言:带有分区侦听器 impl 的 Java
基础设施:有 4 个 jvm 并行运行(这意味着,4 个使用者使用相同的组名运行)
问题:当我启动我的第一个消费者 Lister 回调方法并完成分区分配时,这个消费者开始处理我的消息。
举个例子,这个消费者持有一条消息 MSG-1,而我的处理器正在处理这条消息(我故意将 20 毫秒作为线程等待)。因此,没有将 MSG-1 提交到带有偏移量的主题。
消费者的属性 session.timeout.ms = 15 毫秒。
与此同时,消费者 2 开始了,
这个消费者开始,分配了分区(调用了正确的回调方法)并且没有消费消息,因为这两条消息由消费者 1 持有。
现在,从消费者心跳间隔超过和经纪人认为消费者 1 已死并重新分配消费者 2 的分区(全部) 现在回调在消费者 2 处调用的方法(分配和撤销)。同时,我的会话超时已过期,msg-1 和 msg-2 回到主题并拿起消费者 2。
现在,我已经处理了两次 msg-1 && msg-2.... 一次来自 consumer-1 和 consumer-2
我的问题是,
- Consumer-1 没有被分区撤销回调方法调用?
- 在我的线程睡眠完成后(从消费者 -1 开始),他正尝试使用分区提交偏移量......我们正在完成分区重新分配......您无法提交。这是正确的,但我怎样才能从消费者 1 中获取回调方法......
-纳雷什。
【问题讨论】:
标签: apache-kafka kafka-consumer-api