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

Python asyncio异步函数并发时单任务占CPU无法调度解决方案

问题根因

asyncio 是协作式单线程并发模型,不存在抢占式调度,协程只有执行到await关键字、主动让出事件循环控制权时,其他待调度协程才能获得执行机会。你的代码存在三个核心问题导致调度失效:

  • send_image中for im in self.camera:的逐帧读取是同步阻塞操作,cv2.resize()是纯CPU密集的同步计算,这两段代码执行过程中没有任何await点,会一直占住事件循环线程,完全不给其他协程留调度空间
  • 代码中await asyncio.sleep(1000)是笔误,参数单位为秒,1000秒约等于17分钟才会触发一次让出动作,完全不符合实时帧传输的需求
  • asyncio.gather()的用法错误,该方法返回值是所有协程返回值组成的列表,(done, pending)元组是asyncio.wait()的返回格式,同时缺少异常场景下的任务清理逻辑,容易出现任务泄漏
符合规范的实现方案

核心思路是把所有阻塞事件循环的同步操作(摄像头读帧、OpenCV图像处理、图像编码)全部扔到执行器中运行,避免阻塞事件循环主线程,同时修正协程调度和异常处理逻辑。

第一步:重构send_image协程

使用Python 3.7+标准支持的loop.run_in_executor将同步阻塞操作放到线程池执行,执行过程中事件循环可以正常调度收消息等其他任务:

import asyncio
import cv2
# 假设image_to_bytes是你已有的编码函数

async def send_image(self, ws):
    """Sends an image to the websocket."""
    loop = asyncio.get_running_loop()
    
    def sync_frame_process():
        # 所有同步阻塞操作全部放在这个内部函数里,扔到线程池执行
        im = next(self.camera)  # 同步读帧,阻塞逻辑不占事件循环
        h, w = im.shape[:2]
        resized = cv2.resize(im, (w // 4, h // 4))
        return image_to_bytes(resized)

    while True:
        # 等待同步处理完成,这个await点会让出控制权给其他协程
        frame_bytes = await loop.run_in_executor(None, sync_frame_process)
        await ws.send_bytes(frame_bytes)
        # 按实际帧率需求调整间隔,比如10帧/秒就睡0.1秒,不要写1000
        await asyncio.sleep(0.1)

如果后续加入重CPU计算逻辑(比如AI推理、复杂图像变换),可以把run_in_executor的第一个参数换成提前初始化好的进程池执行器,规避GIL对CPU密集任务的性能影响。

第二步:修正WebSocket端点的调度逻辑

显式创建任务,增加异常场景下的任务取消逻辑,避免任务泄漏:

from fastapi import WebSocket, WebSocketDisconnect

@app.websocket('/ws')
async def websocket_endpoint(websocket: WebSocket):
    """Handle a WebSocket connection."""
    backend = Backend()
    logger.info('Started backend.')
    await websocket.accept()
    send_task = asyncio.create_task(backend.send_image(websocket))
    recv_task = asyncio.create_task(backend.handle_message(websocket))
    try:
        await asyncio.gather(send_task, recv_task)
    except WebSocketDisconnect:
        pass
    finally:
        # 连接断开或出现异常时,强制取消未完成的任务,避免后台残留
        for task in (send_task, recv_task):
            if not task.done():
                task.cancel()
        # 等待任务取消完成,捕获取消异常避免抛出未处理错误
        await asyncio.gather(send_task, recv_task, return_exceptions=True)
        await websocket.close()
额外注意事项
  • 协程中禁止直接写任何无await的长耗时同步逻辑,不管是IO阻塞还是CPU计算,只要单次执行耗时超过1ms,就会拖慢整个事件循环的响应速度
  • handle_message中的json.loads如果处理的消息体超过1MB,也建议放到run_in_executor中执行,避免大JSON解析阻塞事件循环
  • 不要每次WebSocket连接都新建进程/线程池,全局初始化一次执行器复用即可,避免频繁创建销毁进程/线程的开销

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 21:33:19