【问题标题】:Proper way of reading messages from Kafka topic and then closing从Kafka主题读取消息然后关闭的正确方法
【发布时间】:2019-05-31 01:54:19
【问题描述】:

我正在使用 express 和 kafka-node 构建一个简单的 node.js API,当收到 HTTP 请求然后关闭连接时,它会从请求的 Kafka 主题和消费者组返回未读消息。我不需要也不希望消费者继续等待新消息。

在 kafka-node 中,检查是否已到达主题末尾的正确方法是什么,如果是,则关闭与代理的连接并退出应用程序以防止读取新消息?

这是我的 consumer.js。它与 kafka-node 文档中给出的示例几乎相同。

"use strict";

const kafka = require("kafka-node");

let topicName = "testTopic-01",
  groupName = "testGroup-01",
  consumerOptions = {
    kafkaHost: "localhost: 9092",
    groupId: groupName,
    sessionTimeout: 15000,
    protocol: ["roundrobin"],
    fromOffset: "earliest",
    encoding: "utf8"
  };

const consumerGroup = new kafka.ConsumerGroup(consumerOptions, topicName);

consumerGroup.on("message", message => {
  console.log(`Message: ${message.value}`);
});
consumerGroup.on("error", error => {
  console.error(error);
});

console.log(`Consumer started on topic ${topicName} on group ${groupName}`);

【问题讨论】:

    标签: node.js apache-kafka


    【解决方案1】:

    您可以使用#Offset 获取主题分区的当前偏移量。通过比较您分配的主题分区的如此获取的偏移量,您就可以知道相应主题分区中的最后一条消息是什么。

    请记住,如果您有多个并行消费者,您应该跟踪消费者组中的消费者被分配到的主题分区 (#fetchCommits)。

    【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-01-12
    • 1970-01-01
    • 2021-03-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-03-14
    • 2018-07-27
    相关资源
    最近更新 更多