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}
问题根源
- 容器网络寻址错误:生产者和消费者服务运行在Docker容器(
kafka、kafkaprodute)中,容器内的localhost指向容器自身,而非宿主机,因此无法访问宿主机上的kafkaServer容器。 - 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
相关产品推荐
相关产品推荐

