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

Kafka连接报错connect ECONNREFUSED 127.0.0.1:9092,求解决方法

问题:Kafka生产者/消费者连接被拒绝(ECONNREFUSED)

问题描述

尝试通过Kafka生产者向指定topic发送消息,消费者监听该topic并记录数据,但出现连接拒绝错误,错误日志显示无法连接127.0.0.1:9092。

生产者代码

export class ProducerS {
  private kafkaProducer: KafkaProducer;
  constructor() {
    // Configuração do Kafka com os endereços dos brokers
    const kafka = new Kafka({clientId: 'kafkaP', brokers: ['localhost:9092'] });
    // Criação de um produtor Kafka
    this.kafkaProducer = kafka.producer()
  }

  // Método para conectar ao broker Kafka
  async connect(): Promise<void> {
    console.log('Conectou')
    await this.kafkaProducer.connect();
  }

  // Método para desconectar do broker Kafka
  async disconnect(): Promise<void> {
    console.log('Desconectou')
    await this.kafkaProducer.disconnect();
  }

  // Método para enviar uma mensagem para um tópico específico
  async sendM(topic: string, messages: any): Promise<void> {
    // Envia a mensagem para o tópico especificado
    await this.kafkaProducer.send({
      topic,
      messages: [{ value: JSON.stringify(messages) }],
    });
  }
}

消费者代码

export class PagamentoService {
    private kafka: Kafka;
    private consumer: Consumer;

    constructor(
        @InjectRepository(DetalheEntidade)
        private readonly detalheRepository: Repository<DetalheEntidade>,
    ){ 
        this.kafka = new Kafka({
            clientId: 'kafkaP',
            brokers: ['localhost:9092'],
        });
        this.consumer = this.kafka.consumer({groupId: 'pagamentos-group'});
        
    } 

    async pegarConsumer(){
        await this.consumer.connect()
        console.log('Opa amigao')
        await this.consumer.subscribe({topic: 'criar-detalhe'})
        await this.consumer.run({
            // eslint-disable-next-line @typescript-eslint/no-unused-vars
            eachMessage: async ({topic, partition, message}) => {
                const ids = await JSON.parse(message.value.toString())
                const topicS = await JSON.parse(topic.toString())
                console.log(ids)
                console.log(topicS)
            }
        })
    }

}

Docker Compose配置

version: '3.5'
services:
  kafka:
    build:
      context: ./
      dockerfile: ./apps/kafka/Dockerfile
    env_file:
      - .env
    depends_on:
      - postgres
    volumes:
      - .:/usr/src/app 
      - /usr/src/app/node_modules
    command: npm run start:dev kafka 
    
  kafkaprodute:
    build: 
      context: ./
      dockerfile: ./apps/kafkaprodute/Dockerfile
    ports:
      - '4000:5000'
    env_file:
      - .env
    depends_on:
      - kafka
    volumes:
      - .:/usr/src/app 
      - /usr/src/app/node_modules
    command: npm run start:dev kafkaprodute 
    
  postgres:
    image: postgres
    env_file:
      - .env
    ports:
      - '5432:5432'
    volumes:
      - ./db/data:/var/lib/postgresql/data

  postgres_admin:
    image: dpage/pgadmin4
    depends_on:
      - postgres
    env_file:
      - .env
    ports:
      - '15432:80'

  zookeeper:
    image: wurstmeister/zookeeper
    container_name: zookeeper
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
    ports:
      - '2181:2181'

  kafkaServer:
    image: wurstmeister/kafka
    container_name: kafkaServer
    ports:
      - '9092:9092'
    depends_on:
      - zookeeper
    environment:
      KAFKA_ADVERTISED_HOST_NAME: localhost
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
    volumes:
      - /var/run/docker.sock:/var/run/docker.sock

生产者调用代码

const usuarioSalvo = usuarioEntidade.id

const kafkaGerent = new ProducerS()
await kafkaGerent.connect()
await kafkaGerent.sendM(
  'criar-detalhe', {value: JSON.stringify({usuarioSalvo})}
)
await kafkaGerent.disconnect()

错误日志

ERROR [ServerKafka] ERROR [Connection] Connection error: connect ECONNREFUSED 127.0.0.1:9092 {"timestamp":"2024-04-24T15:08:35.624Z","logger":"kafkajs","broker":"localhost:9092","clientId":"nestjs-consumer-server","stack":"Error: connect ECONNREFUSED 127.0.0.1:9092\n    at TCPConnectWrap.afterConnect [as oncomplete] (node:net:1605:16)"}
kafka-1           | [Nest] 29  - 04/24/2024, 3:08:35 PM   ERROR [ServerKafka] ERROR [BrokerPool] Failed to connect to seed broker, trying another broker from the list: Connection error: connect ECONNREFUSED 127.0.0.1:9092 {"timestamp":"2024-04-24T15:08:35.625Z","logger":"kafkajs","retryCount":4,"retryTime":4152}

问题根源

  1. 容器网络寻址错误:生产者和消费者服务运行在Docker容器(kafka、kafkaprodute)中,容器内的localhost指向容器自身,而非宿主机,因此无法访问宿主机上的kafkaServer容器。
  2. Kafka广告地址配置错误:kafkaServer的KAFKA_ADVERTISED_HOST_NAME设为localhost,当容器内客户端连接后,Kafka会返回localhost:9092作为后续通信地址,但容器内无法将该地址解析到kafkaServer容器。

修复方案

1. 修正Kafka容器的网络配置

修改docker-compose.yml中kafkaServer的环境变量,区分内部容器通信和外部宿主机通信的监听地址:

kafkaServer:
  image: wurstmeister/kafka
  container_name: kafkaServer
  ports:
    - '9092:9092'  # 宿主机外部访问端口
    - '9093:9093'  # 容器内部通信端口
  depends_on:
    - zookeeper
  environment:
    KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
    KAFKA_LISTENERS: INSIDE://0.0.0.0:9093,OUTSIDE://0.0.0.0:9092
    KAFKA_ADVERTISED_LISTENERS: INSIDE://kafkaServer:9093,OUTSIDE://localhost:9092
    KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INSIDE:PLAINTEXT,OUTSIDE:PLAINTEXT
    KAFKA_INTER_BROKER_LISTENER_NAME: INSIDE
  volumes:
    - /var/run/docker.sock:/var/run/docker.sock

2. 修改生产者/消费者的Broker地址

容器内的服务需要使用容器服务名kafkaServer加内部端口9093连接:

生产者代码修改

const kafka = new Kafka({clientId: 'kafkaP', brokers: ['kafkaServer:9093'] });

消费者代码修改

this.kafka = new Kafka({
    clientId: 'kafkaP',
    brokers: ['kafkaServer:9093'],
});

3. 重启服务

执行以下命令重启所有容器,确保配置生效:

docker-compose down
docker-compose up -d

额外优化建议

  • 添加健康检查:给kafkaServer和zookeeper添加健康检查,避免依赖服务未就绪时启动生产者/消费者:
    zookeeper:
      # ... 原有配置
      healthcheck:
        test: ["CMD", "zkServer.sh", "status"]
        interval: 10s
        timeout: 5s
        retries: 5
    
    kafkaServer:
      # ... 原有配置
      healthcheck:
        test: ["CMD", "kafka-topics.sh", "--list", "--zookeeper", "zookeeper:2181"]
        interval: 10s
        timeout: 5s
        retries: 5
    
  • 自动创建Topic:在kafkaServer的环境变量中添加KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true',确保发送消息时自动创建不存在的topic。

内容的提问来源于stack exchange,提问作者J. Maciel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 04:37:02