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

如何用asyncio.as_completed()处理可变任务列表及失败任务优先级

解决方案

针对你提出的维持固定并行数、任务完成立即返回结果、失败任务优先重试的需求,我们可以通过优先队列+工作协程池的方式实现,以下是具体代码和说明:

核心思路

  1. 使用asyncio.PriorityQueue管理待执行任务:给失败任务设置更高优先级(优先级数值越小越优先),确保重试任务比新任务先被处理。
  2. 启动固定数量的工作协程,持续从队列中取任务执行。
  3. 任务执行完成后,通过结果队列立即返回结果;若任务失败,将其放回优先队列等待重试;初始任务一次性放入队列(若需动态补充可调整逻辑)。

完整代码实现

import asyncio
import random

async def task_work(task_id):
    """模拟任务执行:随机执行时间,10%失败概率"""
    await asyncio.sleep(random.uniform(0.1, 1.0))
    if random.random() < 0.1:
        print(f"任务 {task_id} 执行失败")
        return False, task_id
    print(f"任务 {task_id} 执行成功")
    return True, task_id

async def process_tasks(initial_tasks, max_parallel=8):
    # 优先队列:优先级0(失败重试任务)> 优先级1(新任务)
    task_queue = asyncio.PriorityQueue()
    
    # 初始化队列:放入所有初始任务,优先级设为1
    for task_id in initial_tasks:
        await task_queue.put((1, task_id))
    
    # 用于传递结果的队列,实现任务完成后立即返回
    result_queue = asyncio.Queue()

    async def worker():
        """工作协程:持续从队列取任务执行"""
        while True:
            priority, task_id = await task_queue.get()
            try:
                success, tid = await task_work(task_id)
                if success:
                    await result_queue.put(("success", tid))
                else:
                    # 失败任务放回队列,优先级设为0(优先执行)
                    await task_queue.put((0, tid))
                    await result_queue.put(("failed", tid))
            finally:
                task_queue.task_done()

    # 启动指定数量的工作协程
    workers = [asyncio.create_task(worker()) for _ in range(max_parallel)]

    # 生成结果:直到所有初始任务成功完成
    completed_count = 0
    total_initial = len(initial_tasks)
    while completed_count < total_initial:
        result = await result_queue.get()
        yield result
        if result[0] == "success":
            completed_count += 1

    # 清理工作协程
    for w in workers:
        w.cancel()
    await asyncio.gather(*workers, return_exceptions=True)

# 测试执行
async def main():
    initial_tasks = list(range(30))  # 30个初始任务
    async for res in process_tasks(initial_tasks, max_parallel=8):
        status, task_id = res
        print(f"结果:任务 {task_id} {status}")

if __name__ == "__main__":
    asyncio.run(main())

关键说明

  • 优先级队列:失败任务被放回队列时使用优先级0,新任务使用优先级1,确保重试任务优先被工作协程获取执行,解决了原方案中失败任务无法插队的问题。
  • 并行数维持:固定启动max_parallel个工作协程,只要队列中有任务就会被处理,始终保持指定的并行度。
  • 即时返回结果:通过result_queue传递执行结果,使用异步生成器yield立即返回,符合“任务完成后立即返回结果”的需求。
  • 终止条件:当所有初始任务都成功完成后,才会取消工作协程并结束,确保失败任务被反复重试直到成功。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 06:11:17