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

Docker部署Kafka后API无法连接Broker创建主题故障排查

问题描述

通过docker-compose部署Kafka、Kafka UI和Zookeeper后,调用POST API期望在写入数据库的同时创建Kafka主题,但始终无法连接到Broker。

相关代码

创建Kafka主题函数

const addKafkaTopic = async (
    topicName: string, 
    topicIp: string, 
    topicPort: string
    ) => {
    console.log(`ip: ${topicIp} port: ${topicPort} topic ${topicName}`);

    const kafka = new Kafka ({
      clientId: 'myclient',
      brokers: [`${topicIp}:${topicPort}`]
    });

    const admin = kafka.admin();
    await admin.connect();

    admin.createTopics({
      topics: [{
        topic: topicName,
        numPartitions: 1,
        replicationFactor: 1
      }]
    });

    return await admin.disconnect();
  };

POST API代码

router.post('/postkafka', async(req, res) => {
    const name = req.body.topicName;
    const ip = req.body.topicIp;
    const port = req.body.topicPort;
    console.log('add called!');
    try {
      await addKafkaTopic(
        name,
        ip,
        port
        );
      console.log('func worked!');
      await db('kafka_table').insert({name, ip, port});
      res.status(201).json({ message: 'Data has been added to DB.' });
      console.log("insert has been made")
    } catch (err) {
      console.error('Error:', err);
      res.status(500).json({ error: 'Something went wrong while adding data to DB.' });
    }

    return res.end();
  });

docker-compose.yaml配置

kafka-ui:
    container_name: cnt_macro1_kafka-ui
    image: provectuslabs/kafka-ui:latest
    ports:
      - "8082:8080"
    environment:
      KAFKA_CLUSTERS_0_NAME: local
      KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka0:29092
      KAFKA_CLUSTERS_0_ZOOKEEPER: zookeeper0:2181
      KAFKA_CLUSTERS_0_JMXPORT: 9997
    depends_on:
      - zookeeper0
      - kafka0

zookeeper0:
    container_name: cnt_macro1_zookeeper0
    image: confluentinc/cp-zookeeper:latest
    ports:
      - "2181:2181"
    volumes:
      - ./zoo/data:/var/lib/zookeeper/data
      - ./zoo/log:/var/lib/zookeeper/log
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000

kafka0:
    container_name: cnt_macro1_kafka0
    image: confluentinc/cp-kafka:latest
    depends_on:
      - zookeeper0
    ports:
      - "9092:9092"
      - "9997:9997"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper0:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka0:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      JMX_PORT: 9997
      KAFKA_JMX_OPTS: -Dcom.sun.management.jmxremote -Dcom.sun.management.jmxremote.authenticate=false -Dcom.sun.management.jmxremote.ssl=false -Djava.rmi.server.hostname=kafka0
      KAFKA_ADVERTISED_HOST_NAME: 192.168.1.100

错误现象

通过Postman向http://localhost:7007/api/kafka/postkafka发送请求,尝试过0.0.0.0、127.0.0.1、本地IP、Docker IP搭配9092或29092端口,均出现连接错误:

  • [BrokerPool] Failed to connect to seed broker, trying another broker from the list: Connection timeout
  • Connection error: connect ECONNREFUSED 0.0.0.0:9092 while sending AdminClient request Kafka

此前曾用0.0.0.0:9092成功创建主题,但当前完全无法操作。


解决方案

核心问题分析

  1. Kafka的ADVERTISED_LISTENERS与ADVERTISED_HOST_NAME配置冲突,后者会覆盖前者的部分规则,导致Broker对外暴露的连接地址混乱
  2. 0.0.0.0是容器内部的监听地址,并非外部客户端可直接连接的有效地址
  3. 代码中createTopics异步操作未加await,会导致Admin连接在主题创建完成前提前断开

具体修复步骤

1. 修正Kafka环境变量配置

移除冲突的KAFKA_ADVERTISED_HOST_NAME,将PLAINTEXT_HOST的地址改为你的本地实际IP(192.168.1.100),确保外部客户端能正确识别Broker地址:

kafka0:
    # 其他配置保持不变
    environment:
      # 其他配置保持不变
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka0:29092,PLAINTEXT_HOST://192.168.1.100:9092
      # 删除 KAFKA_ADVERTISED_HOST_NAME: 192.168.1.100

2. 修复代码异步问题

在addKafkaTopic函数中,给createTopics添加await,确保主题创建完成后再断开Admin连接:

const addKafkaTopic = async (
    topicName: string, 
    topicIp: string, 
    topicPort: string
    ) => {
    console.log(`ip: ${topicIp} port: ${topicPort} topic ${topicName}`);

    const kafka = new Kafka ({
      clientId: 'myclient',
      brokers: [`${topicIp}:${topicPort}`]
    });

    const admin = kafka.admin();
    await admin.connect();

    // 添加await等待主题创建完成
    await admin.createTopics({
      topics: [{
        topic: topicName,
        numPartitions: 1,
        replicationFactor: 1
      }]
    });

    return await admin.disconnect();
  };

3. 调整客户端连接地址

发送POST请求时,使用192.168.1.100作为topicIp,9092作为topicPort。如果你的API服务也运行在Docker容器中,需将其加入Kafka所在的Docker网络,此时可直接用kafka0:29092作为连接地址。

4. 重启Docker服务

修改配置后执行以下命令重启所有服务:

docker-compose down
docker-compose up -d

验证步骤

  1. 访问Kafka UI(http://localhost:8082)确认集群正常运行
  2. 发送POST请求,检查是否能成功创建主题并写入数据库

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 22:18:10