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

如何从本地主机连接Docker Compose部署的KRaft模式Kafka?

本地主机无法连接KRaft模式Kafka问题排查与解决

问题描述

使用无ZooKeeper的KRaft模式Kafka时,Docker内部的api服务可正常连接,但Docker外部通过9094端口使用相同Node.js代码连接时,出现连接关闭错误。

环境配置

Docker Compose配置

services:
  kafka_auth:
    image: confluentinc/cp-kafka:7.2.1
    environment:
      KAFKA_NODE_ID: 2
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,EXTERNAL:PLAINTEXT
      KAFKA_LISTENERS: PLAINTEXT://kafka_auth:9092,CONTROLLER://kafka_auth:9093,EXTERNAL://0.0.0.0:9094
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka_auth:9092,EXTERNAL://127.0.0.1:9094
      KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
      KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka_mqtt:9093,2@kafka_auth:9093,3@kafka_server:9093'
      KAFKA_PROCESS_ROLES: 'broker,controller'
    volumes:
      - ./run_workaround.sh:/tmp/run_workaround.sh
    command: "bash -c '/tmp/run_workaround.sh && /etc/confluent/docker/run'"
    depends_on:
      - kafka_mqtt
    ports:
      - "9094:9094"

  api:
    image: node:18
    volumes:
      - ./:/app
    command: sh -c "yarn install && yarn start:debug"
    depends_on:
      - kafka_mqtt

run_workaround.sh脚本内容

sed -i '/KAFKA_ZOOKEEPER_CONNECT/d' /etc/confluent/docker/configure
sed -i 's/cub zk-ready/echo ignore zk-ready/' /etc/confluent/docker/ensure
echo "kafka-storage format --ignore-formatted -t CLUSTER_ID -c /etc/kafka/kafka.properties" >> /etc/confluent/docker/ensure

Node.js客户端代码

import {
  Kafka,
  Consumer,
} from 'kafkajs';
const kafka = new Kafka({
    clientId: 'clientId',
    brokers: process.env.KAFKA_BOOTSTRAP_SERVERS.split(','),
    requestTimeout: 3600000,
    retry: {
      maxRetryTime: 10000,
      initialRetryTime: 10000,
      retries: 999999999,
    },
});

const consumer = kafka.consumer({ groupId: 'groupId' });
consumer.connect().then(async () => {
  // 检查所有主题,不存在则创建
  await consumer.subscribe({ topics: ['topic1'] });
  await consumer.run({
    eachMessage: async ({ topic, message }) => {
      if (!message?.value) {
        return;
      }
      switch (topic) {
        case 'topic1':
          method(message.value.toString());
          break;
        default:
          break;
      }
    },
  });

错误信息

{"level":"ERROR","timestamp":"2023-08-06T01:16:06.815Z","logger":"kafkajs","message":"[BrokerPool] Closed connection","retryCount":27,"retryTime":10000}

排查与解决步骤

  • 统一集群ID:KRaft模式下所有节点必须使用相同集群ID。当前脚本用固定CLUSTER_ID,需检查另外两个节点(kafka_mqtt、kafka_server)的配置。生成统一UUID:kafka-storage random-uuid,替换脚本中CLUSTER_ID,并删除所有节点存储卷后重启。
  • 修正监听器配置:若使用WSL2或局域网环境,将KAFKA_ADVERTISED_LISTENERS中的EXTERNAL://127.0.0.1:9094替换为宿主机局域网IP(如192.168.x.x:9094);用netstat -an | grep 9094确认端口未被占用。
  • 检查集群健康状态:进入kafka_auth容器,执行kafka-metadata-quorum --bootstrap-server kafka_auth:9092 describe --status,确保三个节点均为ONLINE状态且集群有leader节点。
  • 调整客户端配置:在KafkaJS配置中增加socketTimeout: 30000避免超时;确认宿主机连接时KAFKA_BOOTSTRAP_SERVERS设为127.0.0.1:9094。
  • 验证网络连通性:宿主机执行nc -zv 127.0.0.1 9094,确认端口可连通,若不通检查Docker端口映射或防火墙规则。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 05:40:39