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

如何用asyncio和aiohttp异步遍历任务与文件

异步分页+多任务并发的性能优化方案

核心问题拆解

你遇到的问题本质是两个:

  1. 分页获取JobID时串行等待,没有和后续Job处理并行
  2. 对async for的逻辑理解有误——它是串行遍历异步迭代器,不会自动触发并发,必须手动创建任务实现并行

一、异步分页:边获取边处理,最大化并发

不要等一页所有Job处理完再取下一页,而是在发起分页请求的同时,并行处理已获取的Job。用asyncio.Semaphore控制分页请求的并发数,避免压垮接口。

分页请求函数

async def fetch_page_job_ids(session, url, semaphore):
    async with semaphore:
        async with session.get(url) as resp:
            resp.raise_for_status()
            data = await resp.json()
            return data['job_ids'], data.get('next_url')

边分页边处理的主逻辑

在主流程中,每获取一页JobID就立刻为每个Job创建处理任务,同时继续发起下一页的请求:

async def main():
    initial_url = "https://your-api.com/api/jobs?page=1"
    page_sem = asyncio.Semaphore(5)  # 限制同时发起5个分页请求
    async with aiohttp.ClientSession() as session:
        current_url = initial_url
        job_tasks = []
        while current_url:
            # 异步获取当前页数据
            job_ids, current_url = await fetch_page_job_ids(session, current_url, page_sem)
            # 为当前页每个Job创建处理任务,加入任务列表
            job_tasks.extend([
                asyncio.create_task(process_single_job(session, job_id))
                for job_id in job_ids
            ])
            # 控制任务堆积:当待处理任务过多时,先等待部分完成
            if len(job_tasks) > 20:
                done, job_tasks = await asyncio.wait(job_tasks, return_when=asyncio.FIRST_COMPLETED, n=10)
        # 等待所有Job处理完成
        await asyncio.gather(*job_tasks)

二、Job处理:并发获取FileID+下载

每个Job的FileID查询和文件下载都要并发执行,同样用信号量控制并发数,避免接口限流。

单Job处理逻辑

# 全局信号量控制Job和下载的并发
job_sem = asyncio.Semaphore(10)  # 同时处理10个Job的FileID查询
download_sem = asyncio.Semaphore(15)  # 同时下载15个文件

def should_process_job(job_id):
    # 你的Job过滤逻辑,比如判断状态、类型等
    return True

def should_download_file(file_id):
    # 你的文件过滤逻辑
    return True

async def download_file(session, file_id):
    async with download_sem:
        async with session.get(f"https://your-api.com/api/files/{file_id}/download") as resp:
            resp.raise_for_status()
            content = await resp.read()
            with open(f"./downloads/{file_id}.bin", "wb") as f:
                f.write(content)

async def process_single_job(session, job_id):
    if not should_process_job(job_id):
        return
    async with job_sem:
        # 异步获取FileID列表
        async with session.get(f"https://your-api.com/api/jobs/{job_id}/files") as resp:
            resp.raise_for_status()
            file_ids = await resp.json()
        # 过滤需要下载的文件
        target_files = [fid for fid in file_ids if should_download_file(fid)]
        if not target_files:
            return
        # 并发下载文件
        await asyncio.gather(*[
            asyncio.create_task(download_file(session, fid))
            for fid in target_files
        ])

三、纠正async for的认知误区

async for的作用是遍历异步迭代器(比如从数据库异步拉取数据),但它的执行逻辑是串行的:每次迭代必须等待当前异步操作完成后,才会进入下一次迭代。

错误写法(串行执行):

async def wrong_process(job_iterable):
    async for job_id in job_iterable:
        await process_single_job(job_id)  # 等当前Job处理完才会处理下一个

正确写法(并发执行):

async def correct_process(job_ids):
    tasks = [asyncio.create_task(process_single_job(job_id)) for job_id in job_ids]
    await asyncio.gather(*tasks)  # 所有Job同时处理

关键优化点总结

  1. 分层并发控制:分页请求、Job处理、文件下载分别用信号量限制并发数,平衡性能和接口压力
  2. 边获取边处理:分页请求和Job处理并行,避免等待分页完成再处理
  3. 手动创建任务:用asyncio.create_task和asyncio.gather实现真正的并发,不要依赖async for自动并发

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 06:20:36