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

如何在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"))
关键修改说明
  1. 独立后台任务处理视频逻辑:将视频读取、推理、广播逻辑放到单独的video_broadcast_task异步任务中,只运行一次,和客户端连接逻辑完全解耦,不会阻塞新客户端的连接请求。
  2. 同步操作线程池执行:把vidCap.read()和predict这类同步阻塞操作放到线程池执行,通过run_in_executor让asyncio事件循环能同时处理客户端连接、消息收发等其他任务。
  3. 简化handler逻辑:handler只负责维护客户端的加入/移除,不再执行视频处理循环,确保每个客户端连接都能被及时处理。
  4. 帧率与异常控制:添加视频读取失败的重试逻辑,以及循环间隔控制,避免CPU资源被无限制占用;客户端优化了异常处理,能明确区分连接断开和其他错误。

内容的提问来源于stack exchange,提问作者Ian Shih

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 03:38:25