如何用asyncio和aiohttp异步遍历任务与文件
异步分页+多任务并发的性能优化方案
核心问题拆解
你遇到的问题本质是两个:
- 分页获取JobID时串行等待,没有和后续Job处理并行
- 对
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同时处理
关键优化点总结
- 分层并发控制:分页请求、Job处理、文件下载分别用信号量限制并发数,平衡性能和接口压力
- 边获取边处理:分页请求和Job处理并行,避免等待分页完成再处理
- 手动创建任务:用
asyncio.create_task和asyncio.gather实现真正的并发,不要依赖async for自动并发
内容的提问来源于stack exchange,提问作者Drphoton
相关产品推荐
相关产品推荐

