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

如何用asyncio实现批量并发任务?代码无法并发问题排查

问题:asyncio批量任务无法并发执行的原因及解决方法

我需要用asyncio实现批量任务处理:批次之间顺序执行,但每个批次内的慢I/O任务要并发运行。现有代码可以运行,但批次内任务始终串行执行,无法并发。

原代码

import asyncio
import random
import time

async def process_single(item: str) -> str:
    """Slow I/O function to process a single element"""
    print(f"started: {item}")
    time.sleep(random.randint(1, 5))  # Simulating I/O-bound processing
    print(f"processed: {item}")
    return f"returned {item}"

async def process_batch(batch: list[str]) -> list[str]:
    """Executes a slow I/O bound function asynchronously on a batch of samples"""
    results = []
    for item in batch:
        result = await process_single(item)
        results.append(result)
    return results

def write_results(results: list[str]):
    print(results)

async def main():
    batches = [['a', 'b', 'c'], ['d', 'e', 'f']]

    for batch in batches:
        results = await process_batch(batch)
        write_results(results)

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

原输出(串行执行)

started: a
processed: a
started: b
processed: b
started: c
processed: c
['returned a', 'returned b', 'returned c']
started: d
processed: d
started: e
processed: e
started: f
processed: f
['returned d', 'returned e', 'returned f']

错误原因

  1. 误用同步阻塞函数:time.sleep()是同步阻塞函数,调用时会直接卡住整个事件循环,导致其他异步任务无法获得执行机会,完全失去了异步的意义。
  2. 串行await任务:process_batch里用for循环逐个await process_single(item),这会等待当前任务完全结束后才启动下一个,自然无法并发。

修正后的代码

import asyncio
import random

async def process_single(item: str) -> str:
    """Slow I/O function to process a single element"""
    print(f"started: {item}")
    # 替换为异步的sleep,让出事件循环控制权
    await asyncio.sleep(random.randint(1, 5))  
    print(f"processed: {item}")
    return f"returned {item}"

async def process_batch(batch: list[str]) -> list[str]:
    """并发执行批次内的所有任务"""
    # 使用asyncio.gather并发运行所有任务,一次性等待全部完成
    tasks = [process_single(item) for item in batch]
    results = await asyncio.gather(*tasks)
    return results

def write_results(results: list[str]):
    print(results)

async def main():
    batches = [['a', 'b', 'c'], ['d', 'e', 'f']]

    for batch in batches:
        results = await process_batch(batch)
        write_results(results)

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

预期并发输出示例

started: a
started: b
started: c
processed: b
processed: a
processed: c
['returned a', 'returned b', 'returned c']
started: d
started: e
started: f
processed: e
processed: d
processed: f
['returned d', 'returned e', 'returned f']

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 23:10:08