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

如何使用asyncio实现依赖任务与非依赖任务的并行执行及循环内任务的优化调度

如何使用asyncio实现依赖任务与非依赖任务的并行执行及循环内任务的优化调度

看起来你在asyncio的任务并行调度上遇到了两个典型问题,我来帮你逐个解决~

一、解决第一部分:依赖前置任务后的并行执行

你提到第一部分的代码看起来是顺序执行,但其实asyncio.gather()本身就是用来并行运行多个协程的核心工具,可能是你测试时的误解?不过先确认核心逻辑:first_execution是同步函数,会立即执行并返回结果,之后我们要让check_first_execution(耗时10秒)和second_execution(耗时5秒)同时启动,这样后者会先完成。

你的原代码逻辑其实是对的,但为了更直观,我们可以用asyncio.create_task()显式创建后台任务,再等待,这样更清晰:

async def main():
    choice = random.choice([0, 1])
    result_from_first = first_execution(choice)
    
    # 第一部分:修正后的并行逻辑
    # 先把check任务放到后台运行
    check_task = asyncio.create_task(check_first_execution(result_from_first))
    # 根据结果选择执行second或third
    if result_from_first == 1:
        await second_execution(result_from_first)
    else:
        await third_execution(result_from_first)
    # 等待后台的check任务完成(如果还没结束的话)
    await check_task
    
    # 如果你需要收集两个任务的结果,也可以用gather,效果完全一致:
    # if result_from_first ==1:
    #     results = await asyncio.gather(check_first_execution(result_from_first), second_execution(result_from_first))
    # else:
    #     results = await asyncio.gather(check_first_execution(result_from_first), third_execution(result_from_first))

这样运行后,输出会完全符合你的预期:

First result 1
Second result 2
First check of value 1 complete

因为check_task和second_execution是同时启动的,后者耗时更短,会先输出结果,10秒后check任务完成并打印信息。

二、解决第二部分:循环内的任务调度优化

你现在的循环是每次先等second_execution完成,再等check_second_execution完成才进入下一次循环,这导致了完全的顺序执行。要实现“先连续执行多个second,同时check在后台运行”的效果,我们可以把check_second_execution作为后台任务提交,不用立即等待,等所有second任务都执行完后,再统一等待所有check任务结束:

优化方案:后台运行check任务,不阻塞循环

async def main():
    # ... 前面的第一部分代码 ...

    # 第二部分:优化后的调度逻辑
    list_results_from_first = [result_from_first+i for i in range(5)]
    # 用列表存储所有check任务,最后统一等待
    check_tasks = []
    for first_result in list_results_from_first:
        # 先执行second_execution,拿到结果
        second_result = await second_execution(first_result)
        # 创建check任务放到后台运行,不等待直接进入下一次循环
        check_task = asyncio.create_task(check_second_execution(second_result))
        check_tasks.append(check_task)
    # 所有second任务都执行完后,等待所有后台的check任务完成
    await asyncio.gather(*check_tasks)

这样运行后的输出会类似:

Second result 3
Second result 4
Second result 5
Second result 6
Second result 7
Second check of value 3 complete
Second check of value 4 complete
Second check of value 5 complete

check任务会在后台并行运行,各自完成后就会打印信息,而second任务会依次执行(符合你需求里的“先跑多个second”的逻辑)。

如果你的需求是second任务也想并行执行,还可以进一步优化,把所有second任务也放到后台:

async def main():
    # ... 前面的第一部分代码 ...

    # 第二部分:second也并行的版本
    list_results_from_first = [result_from_first+i for i in range(5)]
    # 一次性创建所有second任务并并行执行
    second_tasks = [asyncio.create_task(second_execution(res)) for res in list_results_from_first]
    second_results = await asyncio.gather(*second_tasks)
    # 再启动所有check任务并行执行
    check_tasks = [asyncio.create_task(check_second_execution(res)) for res in second_results]
    await asyncio.gather(*check_tasks)

这种情况下,所有second任务会同时启动,5秒后一起输出结果,10秒后所有check任务完成并打印信息,适合任务之间没有依赖的场景。

完整修正后的代码

整合两部分优化后的完整代码如下:

import asyncio
import random

def first_execution(choice):
    if choice == 1:
        print(f"First result {choice}")
        return choice
    else:
        print(f"First result {0}")
        return 0

async def check_first_execution(result_from_first):
    await asyncio.sleep(10)
    print(f"First check of value {result_from_first} complete")

async def second_execution(result_from_first):
    await asyncio.sleep(5)
    print(f"Second result {result_from_first+1}")
    return result_from_first+1

async def check_second_execution(result_from_second):
    await asyncio.sleep(10)
    print(f"Second check of value {result_from_second} complete")

async def third_execution(result_from_first):
    await asyncio.sleep(5)
    print(f"Third result {result_from_first+2}")
    return result_from_first+2

async def check_third_execution(result_from_third):
    await asyncio.sleep(10)
    print(f"Third check of value {result_from_third} complete")


async def main():
    choice = random.choice([0, 1])
    result_from_first = first_execution(choice)
    
    # 第一部分:并行执行check_first和second/third
    check_first_task = asyncio.create_task(check_first_execution(result_from_first))
    if result_from_first == 1:
        await second_execution(result_from_first)
    else:
        await third_execution(result_from_first)
    await check_first_task

    # 第二部分:循环内优化,后台运行check任务
    list_results_from_first = [result_from_first+i for i in range(5)]
    check_tasks = []
    for first_result in list_results_from_first:
        second_result = await second_execution(first_result)
        check_task = asyncio.create_task(check_second_execution(second_result))
        check_tasks.append(check_task)
    await asyncio.gather(*check_tasks)

asyncio.run(main())

备注:内容来源于stack exchange,提问作者HienPham

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.15 03:23:14