【问题标题】:NodeJS Kafka Consumer is getting duplicate messages?NodeJS Kafka Consumer 收到重复消息?
【发布时间】:2019-11-20 08:16:15
【问题描述】:

我在我的 NodeJS 应用程序中使用Kafka-Node 来生成和使用消息。我启动了一个正在等待我的主题的消费者。然后我启动生产者并将消息发送到 Kafka。我的消费者正在将这些消息中的每一条都插入到 Postgres 数据库中。

对于单个消费者来说,这很好用。

当我停止消费者并继续生产时,我会在大约 30 秒后重新启动消费者。我注意到大约有十几条消息已经从原始消费者插入到数据库中。

我假设当我杀死消费者时,还有一些尚未提交的偏移量,这就是为什么第二个消费者正在拾取它们?

处理这种情况的最佳方法是什么?

var kafka = require('kafka-node');
var utilities = require('./utilities');
var topics = ['test-ingest', 'test-ingest2'];
var groupName = 'test';

var options = {
    groupId: groupName,
    autoCommit: false,
    sessionTimeout: 15000,
    fetchMaxBytes: 1024 * 1024,
    protocol: ['roundrobin'],
    fromOffset: 'latest',
    outOfRangeOffset: 'earliest'
};

var consumerGroup = new kafka.ConsumerGroup(options, topics);

// Print the message
consumerGroup.on('message', function (message) {

    // Submit our message into postgres - return a promise
    utilities.storeRecord(message).then((dbResult) => {

        // Commit the offset
        consumerGroup.commit((error, data) => {
            if (error) {
                console.error(error);
            } else {
                console.log('Commit success: ', data);
            }
        });

    });

});

【问题讨论】:

    标签: node.js apache-kafka kafka-consumer-api


    【解决方案1】:

    我不知道为什么fromOffset: 'latest' 不适合你。一个简单的解决方法是使用offset.fetchLatestOffsets 来获取最新的偏移量,然后从该点开始使用。

    【讨论】:

    • 嗯,由于某种原因,这似乎在开始时重复了消息,而不是在我停止并重新启动消费者后,
    • @SBB 你用的是同一个消费群吗?
    • 是的,我正在重新启动同一个消费者文件,其中定义了组名var groupName = 'test';
    • 遗憾的是,在您提出建议之前,重复发生在最后几条记录插入数据库之前。根据您的建议,重复发生在插入记录的开头。
    • @SBB 顺便问一下,你在哪里定义了 Kafka 主机?
    猜你喜欢
    • 1970-01-01
    • 2019-09-01
    • 2017-02-15
    • 1970-01-01
    • 2019-12-06
    • 2013-04-24
    • 2019-01-23
    • 1970-01-01
    • 2018-11-28
    相关资源
    最近更新 更多