能否从Python websocket读取数据写入Kafka?收首条消息连接就关闭如何解决
问题根因
- 你当前使用的同步Kafka生产者的
Future.get()方法是阻塞调用,会直接卡住asyncio事件循环,导致WebSocket客户端无法执行心跳保活逻辑,服务端检测到心跳超时后主动断开了连接 - 日志中没有抛出业务报错、连接关闭码为1000(正常关闭)也符合心跳超时的断开特征
- 额外注意你提供的WebSocket客户端代码里有语法错误:
wait ws_client.subscribe(topic)应该改为await ws_client.subscribe(topic)
修复方案
方案1:换用异步Kafka客户端(推荐)
使用适配asyncio生态的aiokafka库,完全异步的生产逻辑不会阻塞事件循环:
- 安装依赖
pip install aiokafka - 修改代码
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
相关产品推荐
相关产品推荐

