【发布时间】: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