Kafka Python消费者无数据消费及Docker对接Broker问题求助
本地API通过Kafka对接Docker环境问题
我正在寻求解决方案,以通过Kafka将本地API(localhost)对接至Docker环境。
现有代码情况
Producer代码(运行正常)
from kafka import KafkaProducer import json producer = KafkaProducer(bootstrap_servers='localhost:9092', value_serializer=lambda v: json.dumps(v).encode('utf-8')) with open('data/New York_0.json', 'r') as f: data = json.load(f) producer.send('my_topic', data) producer.flush()
Consumer代码(仅监听无数据消费)
from kafka import KafkaConsumer import json import traceback # 消费最新消息并手动提交偏移量 try: consumer = KafkaConsumer('my_topic', bootstrap_servers='localhost:9092', value_deserializer=lambda m: json.loads(m.decode('utf-8')), consumer_timeout_ms=6000, enable_auto_commit=False) print("Consumer created, start consuming") for message in consumer: print(message.value) except Exception as e: print(f"An error occurred: {e}") print(f"Exception Type: {type(e).__name__}") traceback.print_exc() finally: print("Closing the consumer") consumer.close()
排查情况
- 测试容器运行
test_producer.py时触发NoBrokersAvailable错误 - 关闭防火墙后问题未解决
解决方案
1. 修正Kafka容器的网络配置
Docker容器的localhost指向自身而非宿主机,需调整Kafka的监听配置,同时暴露宿主机可访问的端口:
docker run -d \ --name kafka \ --link zookeeper:zookeeper \ -p 9092:9092 \ -p 29092:29092 \ -e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092,PLAINTEXT_HOST://0.0.0.0:29092 \ -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092 \ -e KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 \ confluentinc/cp-kafka:latest
PLAINTEXT用于容器内部服务间通信PLAINTEXT_HOST用于宿主机和外部客户端访问
2. 调整代码中的bootstrap_servers
- 本地运行的Producer/Consumer:改为
bootstrap_servers='localhost:29092' - Docker容器内的客户端:改为
bootstrap_servers='kafka:9092'(通过容器服务名访问)
3. 验证Kafka连接状态
- 宿主机端查看主题列表:
docker exec kafka kafka-topics.sh --list --bootstrap-server localhost:29092 - 容器内发送测试消息:
输入消息后用本地Consumer验证是否能接收。docker exec -it kafka bash kafka-console-producer.sh --broker-list kafka:9092 --topic my_topic
4. 重置消费者偏移量
若Consumer能连接但无数据,可能是偏移量已到最新位置,重置到最早位置:
docker exec kafka kafka-consumer-groups.sh --bootstrap-server kafka:9092 --group your-consumer-group --reset-offsets --to-earliest --topic my_topic --execute
替换your-consumer-group为你的消费者组名,未指定时可通过kafka-consumer-groups.sh --list查看默认组名
内容的提问来源于stack exchange,提问作者Hawvb Ged
相关产品推荐
相关产品推荐

