You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

多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监听器

验证步骤

  1. 停止并删除原有容器:docker-compose down
  2. 启动新配置的集群:docker-compose up -d
  3. 等待集群稳定后,重新运行Node.js代码创建主题

内容的提问来源于stack exchange,提问作者yokus

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 23:49:52