【发布时间】:2021-08-17 17:48:29
【问题描述】:
当应用程序启动时,消费者已连接并且可以从日志中验证它正在使用消息。但过了一会儿,它停止消费消息,并且在消费者控制台中看到没有与应用程序地址连接的代理。
【问题讨论】:
标签: node.js apache-kafka librdkafka
当应用程序启动时,消费者已连接并且可以从日志中验证它正在使用消息。但过了一会儿,它停止消费消息,并且在消费者控制台中看到没有与应用程序地址连接的代理。
【问题讨论】:
标签: node.js apache-kafka librdkafka
example: if there are 3 partitions, use 3,6 or 9 consumers
实现断开事件,并在内部尝试连接代理并轮询消息
const consumer = new kafka.KafkaConsumer({
'group.id': 'test-consumer-group',
'metadata.broker.list': process.env.KAFKA_HOST,
'enable.auto.commit': false,
'enable.partition.eof': true
})
consumer.connect()
consumer.on('event.error', function (err) {
if (err == 'ETIMEDOUT') {
consumer.commit()
count = 0
}
consumer.disconnect()
logger('debug', {
message: "Error connecting to kafka consumer" + err
})
snmp_trap(1001)
})
consumer.on('disconnected', function (data) {
console.log("Disconnected. Reconnecting...");
consumer.connect();
});
consumer.on('ready', () => {
try {
consumer.subscribe([topicName]);
} catch (err) {
logger("debug", {
message: "error subscribing the topic: ",
topicName
})
}
try {
setInterval(() => {
consumer.consume(parseInt(process.env.POLL_SIZE || 100));
}, parseInt(process.env.POLL_INTERVAL || 1000))
} catch (err) {
logger("debug", {
message: "error starting consumer"
})
}
})
consumer.on('data', (data) => {
if (!data || !data.value) return
try {
let flag = offsetMap.get(`${data.topicName}${data.partition}${data.offset}`)
if (flag) {
logger("debug", `DUPLICATE OFFSET RECEIVED: ${flag}`)
consumer.commit({
offset: data.offset,
partition: data.partition,
topic: topicName
})
return
}
offsetMap.put(`${data.topicName}${data.partition}${data.offset}`, true)
logger("debug", {
message: "getting data" + data.value
})
msgQueue.push(data, function (err) {
// logger("debug", "error pushing data into queue")
})
logger("debug", {
message: `Current Queue size: ${msgQueue.length()}`
})
count++;
if (count > queueSize) {
consumer.commit({
offset: data.offset,
topic: topicName,
partition: data.partition
})
count = 0
}
} catch(err) {
console.log("ERR >>", err);
consumer.connect()
}
})
【讨论】: