不同线程下两个asyncio事件循环的通信方案问询
问题背景
开发一个包含后台任务与WebSocket服务器的程序,逻辑是客户端连接WebSocket触发事件,服务器完成任务后通知所有客户端并返回结果。运行时触发核心错误:
got Future <Future pending> attached to a different loop
当前实现的问题:
- WebSocket服务器基于websockets库,运行在asyncio事件循环的异步函数中;后台任务运行在同步线程中
- 在任务线程中新建asyncio事件循环,还设置了共享锁,但仍触发上述错误
- 尝试过用asyncio Queue,但该队列依赖异步函数调用,无法解决跨线程循环的问题
疑问:
- 是否存在不依赖事件循环的通信方式?
- 该场景的正确实现方案是什么?
- 不同线程的事件循环能否实现通信?
核心原因
报错本质是:WebSocket连接对象(websockets.WebSocketServerProtocol实例)绑定了WebSocket服务器所在的事件循环,后台线程新建的事件循环中调用client.send(),相当于让绑定到A循环的Future在B循环中运行,必然触发错误。
正确实现方案
不需要在后台线程新建事件循环,而是利用线程安全的API,让后台任务把消息提交到WebSocket服务器的事件循环中执行。asyncio提供的loop.run_coroutine_threadsafe()方法,专门用于跨线程向事件循环提交异步任务。
具体优化步骤
- 移除后台线程中新建事件循环的逻辑,直接复用WebSocket服务器的事件循环
- 使用
loop.run_coroutine_threadsafe()将发送消息的协程提交到WebSocket服务器的事件循环中执行(该方法线程安全) - 用线程锁替代asyncio锁管理客户端列表(后台线程是同步环境,asyncio锁仅支持异步场景)
修正后的Python代码
import asyncio import websockets import threading import time from random import randint class WebSocketServer: INSTANCE = None ADDR = "127.0.0.1" PORT = 7001 def __init__(self): self.clients = {} self.loop = asyncio.new_event_loop() self.clients_lock = threading.Lock() # 用线程锁管理客户端列表 async def add_client(self, client): with self.clients_lock: self.clients[client.id] = client return client async def remove_client(self, client): with self.clients_lock: if client.id in self.clients: del self.clients[client.id] async def handle_client(self, client): await self.add_client(client) try: while True: packet = await client.recv() print("Packet received", packet) await client.send(packet) print("Sending !") await asyncio.sleep(2) # 用asyncio.sleep替代time.sleep,避免阻塞事件循环 except Exception as e: print("Client disconnected", e) finally: await self.remove_client(client) def run(self): server = websockets.serve(self.handle_client, WebSocketServer.ADDR, WebSocketServer.PORT, loop=self.loop) self.loop.run_until_complete(server) self.loop.run_forever() def notify_clients(ws_server): while True: # 先复制客户端列表,避免遍历过程中列表被修改 with ws_server.clients_lock: clients = list(ws_server.clients.values()) for client in clients: # 用run_coroutine_threadsafe把协程提交到WebSocket的事件循环中执行 future = asyncio.run_coroutine_threadsafe(send_message(client), ws_server.loop) try: # 可选:等待协程执行完成,获取结果 future.result() print(f"Thread {threading.current_thread().ident} finished sending to client {client.id}") except Exception as e: print(f"Error sending message: {e}") time.sleep(randint(1, 3)) async def send_message(client): print("Trying to send message") try: await client.send(f"{threading.current_thread().ident} say hi") print("Message sent") except Exception as e: print(f"Error occurred: {e}") if __name__ == "__main__": ws_server = WebSocketServer() ws_thread = threading.Thread(target=ws_server.run) ws_thread.start() print("Running...") # 启动3个后台通知线程 for i in range(3): t = threading.Thread(target=notify_clients, args=(ws_server,)) t.start()
关键优化点
- 替换
time.sleep(2)为await asyncio.sleep(2):避免阻塞WebSocket服务器的事件循环,保证异步任务正常调度 - 用
threading.Lock管理客户端列表:适配后台同步线程的锁需求,asyncio锁无法在同步环境中使用 - 调用
asyncio.run_coroutine_threadsafe()提交跨线程任务:该方法返回concurrent.futures.Future对象,可同步等待执行结果,且线程安全 - 遍历前复制客户端列表:避免遍历过程中客户端断开连接导致列表结构变化引发异常
疑问解答
是否存在不依赖事件循环的通信方式?
不存在。WebSocket连接对象本身是asyncio异步资源,所有操作必须在它绑定的事件循环中执行,但可以通过线程安全的API将操作提交到对应循环,无需自行新建循环。该场景的正确实现方案是什么?
后台同步线程通过run_coroutine_threadsafe()将异步任务提交到WebSocket服务器的事件循环中执行,统一复用同一个事件循环处理所有WebSocket相关操作,这是最简单高效的方案。不同线程的事件循环能否实现通信?
可以,但当前场景完全没必要。如果确实需要跨循环通信,可结合asyncio.Queue与loop.call_soon_threadsafe()传递消息,但会增加复杂度。
内容的提问来源于stack exchange,提问作者Cory Carney

