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

Kafka-Python消费者端无法接收消息问题求助

问题:Docker中Kafka CLI正常收发消息,但Python消费者无法接收消息

我的Kafka运行在Docker环境中,通过CLI创建主题并启动生产者、消费者时,Kafka可以正常传输消息,但使用Python代码执行相同操作时,消费者端无法接收消息。

我的docker-compose文件如下:

version: '3'
services:
  zookeeper:
    image: zookeeper
    container_name: zookeeper
    ports:
      - "2181:2181"
    networks:
      - kafka-net
  kafka:
    image: wurstmeister/kafka
    container_name: kafka
    ports:
      - "9092:9092"
    environment:
      KAFKA_ADVERTISED_LISTENERS: INSIDE://kafka:9092,OUTSIDE://localhost:9093
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INSIDE:PLAINTEXT,OUTSIDE:PLAINTEXT
      KAFKA_LISTENERS: INSIDE://0.0.0.0:9092,OUTSIDE://0.0.0.0:9093
      KAFKA_INTER_BROKER_LISTENER_NAME: INSIDE
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_CREATE_TOPICS: "chatgpt:1:1"
    networks:
      - kafka-net
    expose:
      - 9092
networks:
  kafka-net:
    driver: bridge

生产者代码:

import json
import time
from kafka import KafkaProducer

producer = KafkaProducer(bootstrap_servers=['localhost:9092'],
                         value_serializer=lambda x:
                         json.dumps(x).encode('utf-8'))

producer.send("chatgpt", f"Hello 1", f"key1".encode("utf-8"))
producer.flush()

消费者代码:

import json
from kafka import KafkaConsumer

consumer = KafkaConsumer("chatgpt", bootstrap_servers=['localhost:9092'],
                         auto_offset_reset='earliest',
                         enable_auto_commit=True,
                         group_id='my-group',
                         value_deserializer=lambda x: json.loads(x.decode('utf-8')))

consumer.subscribe(['chatgpt'])
consumer.subscription()

for msz in consumer:
    print(msz)

问题排查与修复

1. 核心问题:Kafka监听端口与外部连接不匹配

你的Docker Compose配置中,Kafka声明了两个监听地址:

  • INSIDE://kafka:9092:供容器内部服务(比如ZooKeeper)连接
  • OUTSIDE://localhost:9093:供宿主机或外部客户端连接

但你只做了9092:9092的端口映射,并且Python代码连接的是localhost:9092(容器内部的INSIDE地址)。宿主机无法解析kafka:9092这个容器内部域名,导致Python生产者实际无法正确把消息发送到Kafka,自然消费者收不到。

而CLI能正常工作,是因为你大概率是进入Kafka容器内部执行的命令,直接使用了kafka:9092这个内部可用的地址,所以没有问题。

2. 次要问题:Python代码的参数与逻辑冗余

  • 生产者send方法参数顺序易出错:kafka-python的send方法签名是send(topic, value=None, key=None, ...),明确指定关键字参数更稳妥;另外添加get()方法可以等待发送确认,方便排查是否发送成功。
  • 消费者重复订阅主题:初始化KafkaConsumer时已经指定了"chatgpt"主题,后续的subscribe属于冗余操作。

修复后的配置与代码

修正Docker Compose的Kafka端口映射

在kafka服务的ports中添加9093:9093,让宿主机可以访问OUTSIDE监听端口:

kafka:
  image: wurstmeister/kafka
  container_name: kafka
  ports:
    - "9092:9092"
    - "9093:9093"  # 新增外部连接端口映射
  environment:
    KAFKA_ADVERTISED_LISTENERS: INSIDE://kafka:9092,OUTSIDE://localhost:9093
    KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INSIDE:PLAINTEXT,OUTSIDE:PLAINTEXT
    KAFKA_LISTENERS: INSIDE://0.0.0.0:9092,OUTSIDE://0.0.0.0:9093
    KAFKA_INTER_BROKER_LISTENER_NAME: INSIDE
    KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
    KAFKA_CREATE_TOPICS: "chatgpt:1:1"
  networks:
    - kafka-net

修复后的生产者代码

import json
from kafka import KafkaProducer

# 连接外部监听端口9093
producer = KafkaProducer(bootstrap_servers=['localhost:9093'],
                         value_serializer=lambda x: json.dumps(x).encode('utf-8'))

# 明确指定value和key参数,添加get()等待发送确认
future = producer.send("chatgpt", value="Hello 1", key="key1".encode("utf-8"))
# 等待发送结果,超时10秒
result = future.get(timeout=10)
producer.flush()

修复后的消费者代码

import json
from kafka import KafkaConsumer

# 连接外部监听端口9093
consumer = KafkaConsumer("chatgpt", bootstrap_servers=['localhost:9093'],
                         auto_offset_reset='earliest',
                         enable_auto_commit=True,
                         group_id='my-group',
                         value_deserializer=lambda x: json.loads(x.decode('utf-8')))

# 初始化时已指定主题,无需重复订阅
for msg in consumer:
    print(f"Key: {msg.key.decode('utf-8')}, Value: {msg.value}")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 00:17:03