如何解决用kafka-python连接Red Panda时的KafkaTimeoutError
解决KafkaTimeoutError:Failed to update metadata after 60.0 secs
针对kafka-python的解决方案
1. 补全SSL证书配置
你已经指定了security_protocol='SSL',但缺少SSL证书相关参数,这会导致客户端无法完成完整的SSL握手,进而无法获取元数据。
Red Panda默认会生成自签名证书,你需要将Docker容器中的证书目录挂载到WSL可访问的路径,然后在生产者配置中添加以下参数:
producer = KafkaProducer( bootstrap_servers='192.168.2.28:19092', security_protocol='SSL', ssl_cafile='/path/to/ca.crt', # Red Panda的CA证书路径 ssl_certfile='/path/to/client.crt', # 客户端证书路径 ssl_keyfile='/path/to/client.key', # 客户端密钥路径 api_version=(2, 8, 0) # 改用Red Panda兼容的高版本API,替代旧的0.10.1 )
如果你的Red Panda容器是通过官方命令启动的,证书通常存储在容器内的/etc/redpanda/certs/目录。可以通过Docker卷映射让WSL访问:
docker run -d --name redpanda-1 \ -p 9092:9092 \ -p 19092:19092 \ -v redpanda-cert:/etc/redpanda/certs \ docker.redpanda.com/redpandadata/redpanda:latest \ redpanda start \ --overprovisioned \ --smp 1 \ --memory 1G \ --reserve-memory 0M \ --node-id 0 \ --kafka-addr PLAINTEXT://0.0.0.0:9092,SSL://0.0.0.0:19092 \ --advertise-kafka-addr PLAINTEXT://localhost:9092,SSL://192.168.2.28:19092
在WSL中,证书卷的路径通常是/var/lib/docker/volumes/redpanda-cert/_data/,可以直接引用这个路径下的证书文件。
2. 确认Topic存在
Red Panda默认开启auto.create.topics.enable,但如果该配置被禁用,发送消息到不存在的Topic会触发元数据超时。你可以通过rpk工具手动创建Topic:
# 进入Red Panda容器执行rpk docker exec -it redpanda-1 rpk topic create test
3. 验证网络连通性
确保WSL能正常访问Red Panda的19092端口:
nc -zv 192.168.2.28 19092
如果连接失败,检查Docker Desktop的端口映射是否正确,以及WSL的网络模式是否能访问宿主机IP。
使用confluent-kafka的替代方案
confluent-kafka对Kafka/Red Panda的兼容性更好,配置更简洁,推荐使用:
1. 安装依赖
pip install confluent-kafka
2. 生产者代码实现
from confluent_kafka import Producer # 配置参数 conf = { 'bootstrap.servers': '192.168.2.28:19092', 'security.protocol': 'SSL', 'ssl.ca.location': '/path/to/ca.crt', 'ssl.certificate.location': '/path/to/client.crt', 'ssl.key.location': '/path/to/client.key' } producer = Producer(conf) # 消息投递回调函数 def delivery_callback(err, msg): if err: print(f"消息投递失败: {err}") else: print(f"消息已投递到 {msg.topic()} 分区 {msg.partition()}") # 发送消息并触发回调 producer.produce('test', b'Hello, World!', callback=delivery_callback) # 强制刷新缓冲区,确保消息发送完成 producer.flush()
内容的提问来源于stack exchange,提问作者exquisitecharm
相关产品推荐
相关产品推荐

