【问题标题】:NodeJS : KafkaJSProtocolError: The group member's supported protocols are incompatible with those of existing membersNodeJS:KafkaJSProtocolError:组成员支持的协议与现有成员的协议不兼容
【发布时间】:2019-08-17 16:51:25
【问题描述】:

我正在尝试使用 MongoDB debezium 连接器从 Kafka 捕获数据,但是当我尝试使用 KafkaJS 读取数据时出现错误:

KafkaJSProtocolError: The group member's supported protocols are incompatible with those of existing members

我正在使用 docker 图像来捕获数据。

以下是步骤,我正在关注:

  1. 启动 Zookeeper

    docker run -it --rm --name zookeeper -p 2181:2181 -p 2888:2888 -p 3888:3888 debezium/zookeeper:latest
    
  2. 启动卡夫卡

    docker run -it --rm --name kafka -p 9092:9092 --link zookeeper:zookeeper debezium/kafka:latest
    
  3. 我已经在复制模式下运行了 MongoDB

  4. 启动 debezium Kafka 连接

    docker run -it --rm --name connect -p 8083:8083 -e GROUP_ID=1 -e CONFIG_STORAGE_TOPIC=my_connect_configs -e OFFSET_STORAGE_TOPIC=my_connect_offsets -e STATUS_STORAGE_TOPIC=my_connect_statuses --link zookeeper:zookeeper --link kafka:kafka  debezium/connect:latest
    
  5. 然后发布 MongoDB 连接器配置

    curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/ -d '{ "name": "mongodb-connector", "config": { "connector.class": "io.debezium.connector.mongodb.MongoDbConnector", "mongodb.hosts": "rs0/abc.com:27017", "mongodb.name": "fullfillment", "collection.whitelist": "mongodev.test", "mongodb.user": "kafka", "mongodb.password": "kafka01" } }'
    
  6. 如果我运行一个观察者 docker 容器,我可以在控制台中以 Json 格式数据

    docker run -it --name watchermongo --rm --link zookeeper:zookeeper --link kafka:kafka debezium/kafka:0.9 watch-topic -a -k fullfillment.mongodev.test
    

但我想在应用程序中捕获这些数据,以便我可以对其进行操作、处理并推送到 ElasticSearch。为此,我正在使用

https://github.com/tulios/kafkajs 

但是当我运行消费者代码时,我得到了错误。这是代码示例

//'use strict';






// clientId=connect-1, groupId=1

const { Kafka } = require('kafkajs')



const kafka = new Kafka({

  clientId: 'connect-1',

  brokers: ['localhost:9092', 'localhost:9093']

})


// Consuming

const consumer = kafka.consumer({ groupId: '1' })



var consumeMessage = async () => {



await consumer.connect()

await consumer.subscribe({ topic: 'fullfillment.mongodev.test' })



await consumer.run({

  eachMessage: async ({ topic, partition, message }) => {

    console.log({

      value: message.value.toString(),

    })

  },

})



}



consumeMessage();


KafkaJSProtocolError: The group member's supported protocols are incompatible with those of existing members

【问题讨论】:

  • 您正在运行什么消费者代码?您从 watch-topic 获取数据的事实表明 Debezium/Kafka 位工作正常。您遇到的错误来自 KafkaJS 以及您如何使用它。
  • 另外,用于使用带有kafka-connect-elasticsearch 的 Kafka Connect 写入 Elasticsearch。将其连接到您将处理过的数据写入其中的 Kafka 主题(或者直接连接到来自 Debezium 的主题,如果您只想将 MongoDB 镜像到 Elasticsearch)
  • 感谢@RobinMoffatt,我已经更新了使用 Nodejs 应用程序使用数据的代码。我也尝试了其他 kafka-connect-elasticsearch ,但我无法在我的虚拟机上安装它
  • 您可以尝试将groupId: '1' 更改为groupId: 'foobar' 吗?该错误消息表明组中有其他同名消费者。

标签: docker apache-kafka apache-kafka-connect debezium kafkajs


【解决方案1】:

您不应该在 Connect 和您的 KafkaJS 使用者中使用相同的 groupId。如果这样做,它们将属于同一个消费者组,这意味着消息只会被其中一个或另一个消费,即使它甚至可以工作。

如果您将 KafkaJS 使用者的 groupId 更改为独特的,它应该可以工作。

请注意,默认情况下,新的 KafkaJS 消费者组将从最新的偏移量开始消费,因此它不会消费已经产生的消息。您可以在 consumer.subscribe 调用中使用 fromBeginning 标志覆盖此行为。见https://kafka.js.org/docs/consuming#from-beginning

【讨论】:

  • 所以这意味着即使来自 Zookeeper、kafka docker 图像 .. 我没有传递任何客户端 id 或组 id .. 在使用它时.. 我可以使用任何随机客户端 id 和组 id 吗?
  • 感谢@tommy 的帮助。我刚刚将组 ID 从 1 更改为 2,我可以通过这个 NodeJS 应用程序在控制台中查看数据...非常感谢...任何建议如何在 ElasticSearch 中接收这些数据......所以自定义节点应用程序来完成这项肮脏的工作
猜你喜欢
  • 2018-12-31
  • 2022-12-12
  • 1970-01-01
  • 2023-04-01
  • 2016-07-13
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-03-04
相关资源
最近更新 更多