如何在while True循环中使用websockets.broadcast()实现多客户端推送
问题根源分析
你的代码核心问题在于:
- 每个客户端连接后,
handler会调用echo函数的无限循环,而这个循环里没有任何异步让出控制权的操作(vidCap.read()是同步阻塞操作,且while True没有await耗时任务),导致asyncio事件循环被第一个客户端的echo完全占用,后续客户端的连接请求根本无法被处理,连clients.add(websocket)都执行不到。 - 尝试线程但无效,是因为同步线程没有和asyncio事件循环正确协同,导致广播消息被积压,直到服务停止才批量发送。
修改后的代码方案
Server.py
import websockets import cv2 import asyncio import time from concurrent.futures import ThreadPoolExecutor # 创建线程池,用于执行同步的CPU密集型操作(如DNN推理),避免阻塞asyncio事件循环 executor = ThreadPoolExecutor(max_workers=2) def predict(image): # 这里替换成你的实际DNN推理逻辑 # 模拟耗时推理 time.sleep(0.05) return "test" async def video_broadcast_task(clients): vidCap = cv2.VideoCapture('rtsp://xxx.xxx.xx') # 替换为你的RTSP地址或本地视频路径 try: while True: # 将同步的视频帧读取操作放到线程池执行,避免阻塞事件循环 ret, image = await asyncio.get_event_loop().run_in_executor(executor, vidCap.read) if not ret: print("视频读取失败/结束,尝试重启读取") vidCap.release() vidCap = cv2.VideoCapture('rtsp://xxx.xxx.xx') await asyncio.sleep(1) continue start = time.time() # 将同步的predict推理放到线程池执行 result = await asyncio.get_event_loop().run_in_executor(executor, predict, image) # 广播结果给所有在线客户端 if clients: websockets.broadcast(clients, result) end = time.time() print(f"exec time:{end - start:.6f} s") # 控制循环频率,避免无限制占用CPU,可根据视频帧率调整 await asyncio.sleep(0.01) finally: vidCap.release() async def handler(websocket, path): clients.add(websocket) print("新客户端连接") try: # 保持连接,等待客户端主动断开 await websocket.wait_closed() finally: clients.remove(websocket) print("客户端断开连接") async def serve(): global clients clients = set() # 启动独立的视频广播后台任务,与客户端连接逻辑解耦 asyncio.create_task(video_broadcast_task(clients)) start_server = await websockets.serve(handler, "localhost", 8765) await start_server.wait_closed() if __name__ == '__main__': asyncio.run(serve())
Client.py
import websockets import asyncio import time async def get_result(uri): try: async with websockets.connect(uri) as websocket: print("连接服务端成功") while True: start = time.time() recv_text = await websocket.recv() print(f"收到结果: {recv_text}") end = time.time() print(f"接收耗时:{end - start:.6f} s") except websockets.exceptions.ConnectionClosed: print("与服务端连接断开") except Exception as e: print(f"发生错误: {e}") if __name__ == '__main__': asyncio.run(get_result("ws://127.0.0.1:8765"))
关键修改说明
- 独立后台任务处理视频逻辑:将视频读取、推理、广播逻辑放到单独的
video_broadcast_task异步任务中,只运行一次,和客户端连接逻辑完全解耦,不会阻塞新客户端的连接请求。 - 同步操作线程池执行:把
vidCap.read()和predict这类同步阻塞操作放到线程池执行,通过run_in_executor让asyncio事件循环能同时处理客户端连接、消息收发等其他任务。 - 简化handler逻辑:
handler只负责维护客户端的加入/移除,不再执行视频处理循环,确保每个客户端连接都能被及时处理。 - 帧率与异常控制:添加视频读取失败的重试逻辑,以及循环间隔控制,避免CPU资源被无限制占用;客户端优化了异常处理,能明确区分连接断开和其他错误。
内容的提问来源于stack exchange,提问作者Ian Shih
相关产品推荐
相关产品推荐

