Python websockets客户端如何为await异步接收函数创建独立线程
需求完全可行,你之前用threading尝试失败的核心原因是async异步函数必须绑定对应线程的asyncio事件循环,不能直接在普通线程中运行。另外你的原有代码还存在几处可修复的问题:函数名拼写错误、await语句写在异步函数外部、async with作用域结束后连接会自动断开、全局变量存连接容易出现悬空引用。
以下是两种适配不同场景的实现方案:
方案一:仅用asyncio原生并发(无额外线程,推荐)
asyncio本身就是为单线程下实现IO并发设计的,你的需求完全不需要开线程,用asyncio的后台任务就可以实现:接收消息作为后台任务跑,主逻辑同时执行其他任务。
修正后的代码:
from typing import Dict import websockets import asyncio import json URL = "ws://localhost:你的服务端口" # websocket地址前缀为ws/wss,不要填http地址 # 不要用全局变量存连接,避免作用域问题 def get_dict(): # 补全你自己的初始化参数逻辑 return {"type": "init"} async def receive_msg(ws): # 直接传入连接对象,不依赖全局变量 while True: msg = await ws.recv() print(f"Server: {msg}") async def listen(): input("按回车开始连接") async with websockets.connect(URL) as ws: # 发送初始化消息 msg_initial: Dict[str, str] = get_dict() await ws.send(json.dumps(msg_initial)) # 创建后台任务运行接收逻辑,不会阻塞当前协程 receive_task = asyncio.create_task(receive_msg(ws)) # 这里写你原本要执行的其他业务逻辑 print("主逻辑正常运行,接收消息已在后台启动") for i in range(10): print(f"主逻辑执行任务{i}") await asyncio.sleep(1) # 模拟你的其他异步操作 # 如需长期保持连接,等待接收任务结束即可 await receive_task if __name__ == "__main__": asyncio.run(listen())
说明:
asyncio.create_task会把接收消息的协程注册到当前事件循环作为后台任务,只要事件循环在运行就会自动调度,完全不需要额外开线程- 连接的生命周期和
async with绑定,出了作用域会自动断开,不会出现资源泄露
方案二:必须用独立线程实现(适配原有同步逻辑场景)
如果你原有主线程的逻辑都是同步代码,没法改成async协程,那可以在独立子线程里单独跑一个asyncio事件循环处理websocket接收:
代码示例:
from typing import Dict import websockets import asyncio import json import threading import queue URL = "ws://localhost:你的服务端口" # 线程安全的通信队列,主线程需要发消息时直接往队列里塞数据即可 send_queue = queue.Queue() def get_dict(): return {"type": "init"} # 子线程入口函数,内部单独运行事件循环 def websocket_thread(): async def ws_task(): async with websockets.connect(URL) as ws: msg_initial: Dict[str, str] = get_dict() await ws.send(json.dumps(msg_initial)) # 两个后台任务:分别处理收消息、发主线程传来的消息 async def recv_task(): while True: msg = await ws.recv() print(f"Server: {msg}") async def send_task(): while True: if not send_queue.empty(): msg = send_queue.get() await ws.send(json.dumps(msg)) await asyncio.sleep(0.1) # 两个任务并行运行 await asyncio.gather(recv_task(), send_task()) # 子线程单独创建专属的事件循环 loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) loop.run_until_complete(ws_task()) if __name__ == "__main__": # 启动websocket子线程,daemon=True表示主程序退出时自动销毁子线程 t = threading.Thread(target=websocket_thread, daemon=True) t.start() # 这里是你主线程的原有同步逻辑,完全不受影响 input("按回车停止程序\n") # 主线程发消息示例:send_queue.put({"content": "来自主线程的消息"})
内容的提问来源于stack exchange,提问作者Gui Reis
相关产品推荐
相关产品推荐

