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
相关产品推荐
相关产品推荐

