【发布时间】:2019-05-16 21:14:11
【问题描述】:
因此,我的 kafka 消费者遇到了一些烦人的偏移提交案例。 我在我的项目中使用“kafka-node”。 我创建了一个主题。 在超过 2 个服务器的消费者组中创建了 2 个消费者。 自动提交设置为 false。 对于我的消费者收到的每条消息,他们都会启动一个异步进程,该进程可能需要 1~20 秒,当进程完成时,消费者会提交偏移量。 我的问题是: 有一个情景,其中, 消费者 1 收到一条消息并需要 20 秒来处理。 在这个过程的中间,他收到了另一条需要 1 秒来处理的消息。 他完成第二条消息处理,提交偏移量,然后立即崩溃。 导致之前的消息处理失败。 如果我重新运行消费者,他不会再次阅读第一条消息,因为第二条消息已经提交了大于第一条的offsst。 我怎样才能避免这种情况?
Kafkaconsumer.on('message', async(message)=>{
await SOMETHING_ASYNC_1~20SEC;
Kafkaconsumer.commit(()=>{});
});
【问题讨论】:
标签: apache-kafka message offset kafka-consumer-api