AsyncMock协程未移交控制权?多Worker任务分配测试异常排查
异步任务分配测试中单个Worker处理所有任务的问题分析与解决
问题场景
我实现了一套异步任务处理结构:
- Worker对象负责接收任务并异步执行,通常会await异步IO调用
- 每个Worker对应一个内部循环,从共享持久队列读取任务并下发给Worker
- 顶层通过
asyncio.gather等待所有内部循环完成
核心代码如下:
import asyncio async def loop_inner(worker, worker_num, tasks): while tasks: task = tasks.pop() print('Worker number {} handled task number {}'.format( worker_num, await worker.do_work(task))) async def loop(workers, tasks): tasks = [loop_inner(worker, worker_num, tasks) for worker_num, worker in enumerate(workers)] await asyncio.gather(*tasks)
实际运行时并行性表现良好,但用AsyncMock替换真实Worker测试任务分配逻辑时,所有任务都被单个Worker处理。测试代码:
from unittest import IsolatedAsyncioTestCase, main from unittest.mock import AsyncMock class TestCase(IsolatedAsyncioTestCase): async def test_yielding(self): tasks = list(range(10)) workers = [AsyncMock() for i in range(2)] for worker in workers: worker.do_work.side_effect = lambda task: task await loop(workers, tasks) main()
测试输出:
Worker number 0 handled task number 9 Worker number 0 handled task number 8 Worker number 0 handled task number 7 Worker number 0 handled task number 6 Worker number 0 handled task number 5 Worker number 0 handled task number 4 Worker number 0 handled task number 3 Worker number 0 handled task number 2 Worker number 0 handled task number 1 Worker number 0 handled task number 0
原因分析
- AsyncMock未真正让出事件循环控制权:你给
do_work设置的side_effect是普通lambda函数,返回的是同步值。虽然AsyncMock的方法默认是协程函数,但如果side_effect返回非可等待对象,await会直接获取结果,不会触发事件循环切换。第一个loop_inner协程会连续执行,不给其他协程运行机会。 - 共享列表的同步操作无挂起点:
tasks.pop()是同步操作,没有任何异步挂起逻辑。第一个协程启动后,会在同一个事件循环tick里把所有任务pop完,其他协程还没来得及执行。
解决方案
方案1:让Mock的协程真正挂起,让出控制权
修改测试代码,让do_work的side_effect是一个异步函数,或者返回一个会挂起的可等待对象(比如asyncio.sleep(0),仅用于触发事件循环切换,实际不等待时间):
from unittest import IsolatedAsyncioTestCase, main from unittest.mock import AsyncMock import asyncio class TestCase(IsolatedAsyncioTestCase): async def test_yielding(self): tasks = list(range(10)) workers = [AsyncMock() for i in range(2)] # 用异步函数作为side_effect,确保await时让出控制权 async def mock_do_work(task): await asyncio.sleep(0) return task for worker in workers: worker.do_work.side_effect = mock_do_work await loop(workers, tasks)
修改后,每个await worker.do_work(task)都会触发事件循环切换,其他loop_inner协程有机会执行,任务会被多个Worker分配处理。
方案2:使用异步队列替代共享列表(更贴近真实场景)
真实场景中你提到是“共享持久队列”,但示例里用了普通列表。换成asyncio.Queue(异步安全的队列),它的get()方法是异步的,当队列为空时会挂起协程,天然支持多个Worker公平获取任务:
先修改核心代码:
import asyncio async def loop_inner(worker, worker_num, task_queue): while True: try: task = task_queue.get_nowait() except asyncio.QueueEmpty: break print('Worker number {} handled task number {}'.format( worker_num, await worker.do_work(task))) task_queue.task_done() async def loop(workers, tasks): task_queue = asyncio.Queue() for task in tasks: task_queue.put_nowait(task) coros = [loop_inner(worker, worker_num, task_queue) for worker_num, worker in enumerate(workers)] await asyncio.gather(*coros) await task_queue.join()
测试代码可以保持Mock的逻辑,即使不用asyncio.sleep(0),队列的异步特性也会让多个Worker交替获取任务:
from unittest import IsolatedAsyncioTestCase, main from unittest.mock import AsyncMock class TestCase(IsolatedAsyncioTestCase): async def test_yielding(self): tasks = list(range(10)) workers = [AsyncMock() for i in range(2)] for worker in workers: worker.do_work.side_effect = lambda task: task await loop(workers, tasks)
这种方式更贴近真实生产环境的队列使用,测试结果也会更符合预期。
内容的提问来源于stack exchange,提问作者alexgolec
相关产品推荐
相关产品推荐

