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
相关产品推荐
相关产品推荐

