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

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

原因分析

  1. AsyncMock未真正让出事件循环控制权:你给do_work设置的side_effect是普通lambda函数,返回的是同步值。虽然AsyncMock的方法默认是协程函数,但如果side_effect返回非可等待对象,await会直接获取结果,不会触发事件循环切换。第一个loop_inner协程会连续执行,不给其他协程运行机会。
  2. 共享列表的同步操作无挂起点: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 17:45:45