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

asyncio异步任务优化:获取目标结果时取消pending任务的实现问题问询

问题解答

1. 为什么get_person任务没有异步运行?

你的代码里藏着一个致命的阻塞点:get_person_request方法使用了time.sleep(),这是同步阻塞调用,会直接卡死整个asyncio事件循环。

asyncio的事件循环是单线程调度的,所有异步任务都在这个线程里轮流执行。如果某个任务触发了同步阻塞操作(比如time.sleep、耗时CPU计算),事件循环就没法切换到其他任务,导致剩下的任务只能等这个阻塞操作结束才能启动——这就是为什么你看不到同时打印三个"Searching..."的原因,代码会先卡死在Bob的time.sleep(3)上,直到它完成才会轮到Alice和Jack的任务。

2. 如何处理asyncio.gather因取消任务抛出的CancelledError?

asyncio.gather()默认行为是"一损俱损":只要有一个任务被取消,它会立刻抛出CancelledError,同时取消所有其他任务,直接中断整个流程,自然拿不到已经找到的目标结果。

解决这个问题的核心是给asyncio.gather()加上return_exceptions=True参数。这个参数会让gather把任务的异常(包括CancelledError)作为返回值返回,而不是直接抛出,这样我们就能在后续结果里过滤出正常完成的任务数据。

另外你当前的取消逻辑也有小问题:没有判断任务是否已经完成,不过这个影响不大,但加上判断会更严谨。


修正后的完整代码

import asyncio

class Persons:
    def __init__(self):
        self.p = []
    
    def get_person_request(self, name):
        # 保留原有的同步逻辑,后续用异步方式调用
        if name == "Alice":
            print("Searching Alice")
            import time
            time.sleep(6)
            print("Returning Alice")
            return {'firstname': "Alice", 'surname': "Donnelly"}
        if name == "Bob":
            print("Searching Bob")
            time.sleep(3)
            print("Returning Bob")
            return {'firstname': "Bob", 'surname': "Murphy"}
        if name == "Jack":
            print("Searching Jack")
            time.sleep(8)
            print("Returning Jack")
            return {'firstname': "Jack", 'surname': "Connell"}
        return None

    async def get_person(self, n, _id):
        # 用asyncio.to_thread把同步阻塞逻辑放到线程池,不阻塞事件循环
        person = await asyncio.to_thread(self.get_person_request, n)
        if person["surname"] == "Murphy":
            # 只取消未完成的其他任务
            for i, task in self.p:
                if i != _id and not task.done():
                    task.cancel()
            return person

    async def get_persons(self, names):
        print("Setting tasks...")
        self.p = [(i, asyncio.create_task(self.get_person(a, i))) for i, a in enumerate(names)]
        print("Gathering async results...")
        # 添加return_exceptions=True,让异常作为返回值而非抛出
        persons = await asyncio.gather(*[task for _, task in self.p], return_exceptions=True)
        # 过滤出正常的字典结果
        return next(p for p in persons if isinstance(p, dict))

def test():
    val = asyncio.run(Persons().get_persons(['Bob', 'Alice', 'Jack']))
    print(val)

if __name__ == "__main__":
    test()

修正后的输出

Setting tasks...
Gathering async results...
Searching Bob
Searching Alice
Searching Jack
Returning Bob
{'firstname': 'Bob', 'surname': 'Murphy'}

关键修正点说明

  1. 用asyncio.to_thread包装同步调用:把get_person_request这个阻塞方法放到线程池执行,让事件循环能正常切换任务,三个请求真正并行启动,同时打印"Searching..."。
  2. 给asyncio.gather加return_exceptions=True:让被取消的任务返回异常对象而不是抛出,保证我们能拿到已经完成的Bob的结果。
  3. 优化取消逻辑:增加not task.done()判断,避免对已完成任务做无效取消操作。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 16:42:35