Docker化Kafka服务中Python生产者send方法阻塞问题排查
Docker Kafka容器中Python生产者卡住的解决办法
核心问题分析
你遇到的卡住问题,主要是Kafka的网络配置错误加上服务启动顺序未确保Kafka就绪导致的:容器内的生产者拿到Kafka返回的连接地址是localhost,会尝试连接自己容器的localhost而非Kafka容器;同时depends_on仅保证启动顺序,不等待Kafka完全可用。
具体修复步骤
1. 修正Kafka的网络广告地址配置
修改docker-compose.yml中Kafka的环境变量,区分Docker内部和本地主机的连接地址:
kafka: image: wurstmeister/kafka container_name: kafka ports: - "9092:9092" environment: KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 # 配置双监听端口:内部容器用kafka服务名,本地主机用localhost KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,PLAINTEXT_HOST://0.0.0.0:9093 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:9092
这样容器内的生产者会通过kafka:9092正确连接Kafka容器,本地代码依然可以用localhost:9092连接。
2. 修复生产者代码的错误
你的代码存在两个明显问题:
- 未导入
time模块就调用sleep,会直接报错 topic变量未定义,发送消息时会抛出异常
修正后的代码片段:
import json import time from kafka import KafkaProducer time.sleep(15) # 延长等待时间,确保Kafka初始化完成 topic = "your_topic_name" # 替换成你的实际Topic名称 producer = KafkaProducer(bootstrap_servers=['kafka:9092'], api_version=(0, 10)) json_path = "/data" with open(json_path, 'r') as f: data = json.load(f) print('Ready to publish') for record in data: producer.send(topic, json.dumps(record).encode('utf-8')) print('Published message !!!') producer.flush()
3. 确保Kafka完全就绪后启动生产者
depends_on只保证启动顺序,不等待服务就绪,建议用wait-for-it工具替代sleep:
- 先修改
Dockerfile安装工具:
FROM python:3.10 WORKDIR /app COPY . /app RUN apt-get update && apt-get install -y wait-for-it \ && pip install --user pip==23.0.1 && pip install pipenv && pipenv install --system ENV ENVIRONMENT=production CMD ["wait-for-it", "kafka:9092", "--timeout=60", "--", "python3", "src/producer.py"]
这样生产者会等待Kafka的9092端口连通后再启动,比固定sleep更可靠。
4. 验证容器内网络连通性(排障用)
如果还是有问题,可进入publisher容器测试网络:
docker exec -it publisher bash # 测试是否能ping通Kafka容器 ping kafka # 测试9092端口是否开放 telnet kafka 9092
内容的提问来源于stack exchange,提问作者Alejandro
相关产品推荐
相关产品推荐

