【问题标题】:Kafka message loss because of later message由于后面的消息而导致 Kafka 消息丢失
【发布时间】: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


    【解决方案1】:

    您实际上想通过使用async.queue 来限制消息和处理并发。

    1. 创建一个async.queue,其中包含消息处理器和一个并发(消息处理器本身被setImmediate 包裹,因此它不会冻结事件循环)
    2. 设置queue.drain 恢复消费者
    3. 消费者message 事件的处理程序暂停消费者并将消息推送到队列。

    kafka-node README 详细信息this here。

    可以找到与您的问题类似的示例实现here。

    【讨论】:

    • 首先谢谢!但是,我想同时处理多条消息。我需要在我的代码中进行并行计算,唯一的问题是崩溃时的偏移量。我为节点尝试了“kue”,并看到当我点击“done”时,它是 kafka 'commit' 的替代项,它运行良好,并且不会丢失下一条消息的偏移量。有没有办法用kafka做到这一点?
    猜你喜欢
    • 2022-01-10
    • 2014-08-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-04-13
    • 1970-01-01
    • 1970-01-01
    • 2017-10-30
    相关资源
    最近更新 更多