Kafka Broker无法接收Python Producer消息的排查求助
Kafka消费者无法接收消息的排查方案
问题背景
使用Docker Compose部署Kafka环境,生产者能发送消息但消费者收不到,相关配置和代码如下:
Docker Compose配置(compose_kafka.yml)
version: '3' services: zookeeper: image: wurstmeister/zookeeper container_name: zookeeper ports: - "2181:2181" environment: ZOO_MY_ID: 1 kafka: image: wurstmeister/kafka container_name: kafka ports: - "9092:9092" environment: KAFKA_ADVERTISED_HOST_NAME: 192.168.1.10 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 kafka_manager: image: kafkamanager/kafka-manager container_name: kafka-manager restart: always ports: - "9000:9000" environment: ZK_HOSTS: "192.168.1.10:2181" APPLICATION_SECRET: "random-secret"
生产者相关代码
Generate.py
from faker import Faker fake = Faker() class Registered_user: def get_registered_user(): return { "name": fake.name(), "address": fake.address(), "created_at": fake.year() }
Producer_registered_user.py
import time import json from kafka import KafkaProducer from fake_data import Generate def json_serializer(data): return json.dumps(data).encode("utf-8") producer = KafkaProducer(bootstrap_servers='192.168.1.10:9092', value_serializer=json_serializer) if __name__ == '__main__': while 1 == 1: user = Generate.Registered_user.get_registered_user() producer.send('registered_user', user) print(user) time.sleep(4)
消费者代码(Consumer_registered_user.py)
import json from kafka import KafkaConsumer if __name__ == '__main__': consumer = KafkaConsumer( bootstrap_servers='192.168.1.10:9092', auto_offset_reset="from-beginning", group_id="consumer-group-a" ) for message in consumer: print("User = {}".format(json.loads(message.value)))
排查步骤
1. 消费者未订阅目标Topic
当前消费者代码未指定要订阅的registered_user Topic,KafkaConsumer初始化后需显式订阅才能接收消息。修改消费者代码:
import json from kafka import KafkaConsumer if __name__ == '__main__': consumer = KafkaConsumer( 'registered_user', # 指定订阅的Topic bootstrap_servers='192.168.1.10:9092', auto_offset_reset="from-beginning", group_id="consumer-group-a" ) for message in consumer: print("User = {}".format(json.loads(message.value)))
也可在初始化后调用consumer.subscribe(['registered_user'])完成订阅。
2. 生产者消息未成功提交
send方法为异步操作,可能消息未完成发送就进入下一次循环。可添加回调并强制刷新确认发送状态:
import time import json from kafka import KafkaProducer from fake_data import Generate def json_serializer(data): return json.dumps(data).encode("utf-8") def on_send_success(record_metadata): print(f"消息已发送到Topic: {record_metadata.topic}, Partition: {record_metadata.partition}, Offset: {record_metadata.offset}") def on_send_error(excp): print(f"消息发送失败: {excp}") producer = KafkaProducer(bootstrap_servers='192.168.1.10:9092', value_serializer=json_serializer) if __name__ == '__main__': while 1 == 1: user = Generate.Registered_user.get_registered_user() producer.send('registered_user', user).add_callback(on_send_success).add_errback(on_send_error) producer.flush() # 确保消息提交到Kafka print(user) time.sleep(4)
3. Kafka Topic配置异常
检查Topic的副本同步状态,进入Kafka容器执行命令:
docker exec -it kafka /opt/kafka/bin/kafka-topics.sh --describe --topic registered_user --bootstrap-server localhost:9092
若ReplicationFactor为1但副本不在ISR(同步副本)列表中,消息无法持久化。可在Docker Compose的Kafka环境变量中添加:
environment: KAFKA_ADVERTISED_HOST_NAME: 192.168.1.10 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_DEFAULT_REPLICATION_FACTOR: 1 KAFKA_MIN_IN_SYNC_REPLICAS: 1
4. 消费者反序列化配置缺失
生产者用JSON序列化消息,但消费者未指定反序列化规则,手动解码可能出现格式问题。可显式配置反序列化器:
import json from kafka import KafkaConsumer if __name__ == '__main__': consumer = KafkaConsumer( 'registered_user', bootstrap_servers='192.168.1.10:9092', auto_offset_reset="from-beginning", group_id="consumer-group-a", value_deserializer=lambda m: json.loads(m.decode('utf-8')) ) for message in consumer: print("User = {}".format(message.value))
5. 网络连通性验证
确认宿主机与Kafka容器的网络连通:
telnet 192.168.1.10 9092
若不通,检查Docker端口映射是否正常,或容器内9092端口是否监听:
docker exec -it kafka netstat -tulpn | grep 9092
内容的提问来源于stack exchange,提问作者Erik hoeven
相关产品推荐
相关产品推荐

