多Kafka Broker环境下node-rdkafka创建Topic超时问题
Kafka多Broker集群创建带副本主题超时问题
问题场景
搭建Kafka环境学习副本机制,单Broker无副本配置时可正常运行,但创建带3副本的主题时出现超时错误。
应用错误日志
error creating topic: LibrdKafkaError: Local: Timed out at Function.createLibrdkafkaError [as create] (/Users/my-name/Developer/test/kafka-nodejs/kafka-local-docker/node_modules/node-rdkafka/lib/error.js:456:10) at /Users/my-name/Developer/test/kafka-nodejs/kafka-local-docker/node_modules/node-rdkafka/lib/admin.js:132:28 { code: -185, errno: -185, origin: 'kafka' }
Node.js代码(node-rdkafka)
const Kafka = require('node-rdkafka'); const admin = Kafka.AdminClient.create({ 'bootstrap.servers': 'localhost:9091,localhost:9092,localhost:9093', 'broker.version.fallback': '0.10.2.1', 'log.connection.close' : false }); const topicName = 'cool-topic-4'; const newTopic = { topic: topicName, num_partitions: 2, replication_factor: 3, config: { 'min.insync.replicas': '2', } } // 创建Kafka主题 admin.createTopic(newTopic, (err) => { if (err) console.log("error creating topic: ", err); else console.log("topic created: ", topicName); } );
Docker Compose配置
--- version: '3' services: zookeeper: image: confluentinc/cp-zookeeper:7.3.0 hostname: zookeeper container_name: zookeeper environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 kafka1: image: confluentinc/cp-kafka:7.3.0 container_name: kafka1 ports: - "9091:9091" depends_on: - zookeeper environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_NODE:PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9091,PLAINTEXT_NODE://kafka1:29091 kafka2: image: confluentinc/cp-kafka:7.3.0 container_name: kafka2 ports: - "9092:9092" depends_on: - zookeeper environment: KAFKA_BROKER_ID: 2 KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_NODE:PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092,PLAINTEXT_NODE://kafka2:29092 kafka3: image: confluentinc/cp-kafka:7.3.0 container_name: kafka3 ports: - "9093:9093" depends_on: - zookeeper environment: KAFKA_BROKER_ID: 3 KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_NODE:PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9093,PLAINTEXT_NODE://kafka3:29093
Kafka Broker日志片段
kafka2 | at kafka.utils.ShutdownableThread.run(ShutdownableThread.scala:96) kafka2 | [2023-07-28 15:56:38,308] INFO [Controller id=2, targetBrokerId=1] Client requested connection close from node 1 (org.apache.kafka.clients.NetworkClient) kafka2 | [2023-07-28 15:56:38,410] INFO [Controller id=2, targetBrokerId=1] Node 1 disconnected. (org.apache.kafka.clients.NetworkClient) kafka2 | [2023-07-28 15:56:38,410] WARN [Controller id=2, targetBrokerId=1] Connection to node 1 (localhost/127.0.0.1:9091) could not be established. Broker may not be available. (org.apache.kafka.clients.NetworkClient) kafka2 | [2023-07-28 15:56:38,410] INFO [Controller id=2, targetBrokerId=3] Node 3 disconnected. (org.apache.kafka.clients.NetworkClient) kafka2 | [2023-07-28 15:56:38,410] WARN [Controller id=2, targetBrokerId=3] Connection to node 3 (localhost/127.0.0.1:9093) could not be established. Broker may not be available. (org.apache.kafka.clients.NetworkClient) kafka2 | [2023-07-28 15:56:38,410] WARN [RequestSendThread controllerId=2] Controller 2's connection to broker localhost:9093 (id: 3 rack: null) was unsuccessful (kafka.controller.RequestSendThread) kafka2 | java.io.IOException: Connection to localhost:9093 (id: 3 rack: null) failed. kafka2 | at org.apache.kafka.clients.NetworkClientUtils.awaitReady(NetworkClientUtils.java:70) kafka2 | at kafka.controller.RequestSendThread.brokerReady(ControllerChannelManager.scala:292) kafka2 | at kafka.controller.RequestSendThread.doWork(ControllerChannelManager.scala:246) kafka2 | at kafka.utils.ShutdownableThread.run(ShutdownableThread.scala:96) kafka2 | [2023-07-28 15:56:38,410] INFO [Controller id=2, targetBrokerId=3] Client requested connection close from node 3 (org.apache.kafka.clients.NetworkClient) kafka2 | [2023-07-28 15:56:38,410] WARN [RequestSendThread controllerId=2] Controller 2's connection to broker localhost:9091 (id: 1 rack: null) was unsuccessful (kafka.controller.RequestSendThread) kafka2 | java.io.IOException: Connection to localhost:9091 (id: 1 rack: null) failed. kafka2 | at org.apache.kafka.clients.NetworkClientUtils.awaitReady(NetworkClientUtils.java:70)
问题原因
从Broker日志可见,Controller(kafka2)尝试连接其他Broker时使用了localhost:9091和localhost:9093,但Docker容器内部的localhost指向容器自身,无法通过宿主机端口映射访问其他Broker,导致Broker之间通信失败,最终主题创建超时。
解决方法
修改Docker Compose中每个Kafka节点的配置,明确区分外部客户端访问地址和Broker内部通信地址,并指定Broker间通信的监听器:
调整后的Docker Compose配置
--- version: '3' services: zookeeper: image: confluentinc/cp-zookeeper:7.3.0 hostname: zookeeper container_name: zookeeper environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 kafka1: image: confluentinc/cp-kafka:7.3.0 container_name: kafka1 ports: - "9091:9091" depends_on: - zookeeper environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_INTERNAL:PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9091,PLAINTEXT_INTERNAL://kafka1:29091 KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT_INTERNAL kafka2: image: confluentinc/cp-kafka:7.3.0 container_name: kafka2 ports: - "9092:9092" depends_on: - zookeeper environment: KAFKA_BROKER_ID: 2 KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_INTERNAL:PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092,PLAINTEXT_INTERNAL://kafka2:29092 KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT_INTERNAL kafka3: image: confluentinc/cp-kafka:7.3.0 container_name: kafka3 ports: - "9093:9093" depends_on: - zookeeper environment: KAFKA_BROKER_ID: 3 KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_INTERNAL:PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9093,PLAINTEXT_INTERNAL://kafka3:29093 KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT_INTERNAL
关键配置说明
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP:定义两种监听器协议,分别用于外部客户端和Broker内部通信KAFKA_ADVERTISED_LISTENERS:对外暴露两个地址,PLAINTEXT供宿主机客户端访问,PLAINTEXT_INTERNAL供Broker之间通过Docker服务名通信KAFKA_INTER_BROKER_LISTENER_NAME:指定Broker之间通信使用PLAINTEXT_INTERNAL监听器
验证步骤
- 停止并删除原有容器:
docker-compose down - 启动新配置的集群:
docker-compose up -d - 等待集群稳定后,重新运行Node.js代码创建主题
内容的提问来源于stack exchange,提问作者yokus
相关产品推荐
相关产品推荐

