You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何解决用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 01:37:43