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

能否从Python websocket读取数据写入Kafka?收首条消息连接就关闭如何解决

问题根因
  • 你当前使用的同步Kafka生产者的Future.get()方法是阻塞调用,会直接卡住asyncio事件循环,导致WebSocket客户端无法执行心跳保活逻辑,服务端检测到心跳超时后主动断开了连接
  • 日志中没有抛出业务报错、连接关闭码为1000(正常关闭)也符合心跳超时的断开特征
  • 额外注意你提供的WebSocket客户端代码里有语法错误:wait ws_client.subscribe(topic)应该改为await ws_client.subscribe(topic)
修复方案

方案1:换用异步Kafka客户端(推荐)

使用适配asyncio生态的aiokafka库,完全异步的生产逻辑不会阻塞事件循环:

  1. 安装依赖
    pip install aiokafka
  2. 修改代码
from aiokafka import AIOKafkaProducer
import asyncio

# 全局初始化异步生产者,避免重复创建
producer = AIOKafkaProducer(bootstrap_servers="你的Kafka集群地址")

async def msg_handler(msg):
    if msg:
        try:
            # 异步等待发送结果,不会阻塞事件循环
            record_metadata = await producer.send_and_wait(kf_topic, value=msg)
            print(record_metadata)
        except Exception as er:
            print(er)

async def main():
    # 启动生产者
    await producer.start()
    ws_client = await WsClient.create(URL, client, msg_handler)
    topic = '/my-topic'
    await ws_client.subscribe(topic)
    while True:
        await asyncio.sleep(60)
    
    # 程序退出前关闭生产者
    await producer.stop()

if __name__ == "__main__":
    asyncio.run(main())

方案2:同步生产者逻辑异步化

如果要保留原有同步Kafka生产者,将阻塞的get操作丢到线程池执行,避免阻塞事件循环:

async def msg_handler(msg):
    if msg:
        try:
            loop = asyncio.get_running_loop()
            f = producer.send(kf_topic, value=msg)
            # 把同步阻塞操作放到线程池执行
            record_metadata = await loop.run_in_executor(None, f.get, 10)
            print(record_metadata)
        except KafkaError as er:
            print(er)

内容的提问来源于stack exchange,提问作者Sarath C Akurathi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 03:24:03