如何用asyncio实现带动态任务队列的并行URL下载?
实现固定并发的动态任务补全(asyncio + aiohttp)
要实现“一个任务完成后立即添加新任务”、保持固定并发数的需求,你可以用asyncio.Queue配合固定数量的worker协程来实现,这样能稳定维持指定的并发量,无需等待整批任务结束再启动新任务。
核心思路
- 用
asyncio.Queue存放所有待下载的URL,worker协程从队列中获取任务 - 启动固定数量(比如16个)的worker,每个worker循环处理队列中的URL
- 当一个worker完成当前下载任务后,会立即从队列取下一个URL继续处理,直到所有URL都处理完毕
完整代码示例
import asyncio import aiohttp async def download_one(session, url): """单个URL的下载逻辑,根据你的需求修改""" try: async with session.get(url, timeout=10) as resp: if resp.status == 200: # 示例:读取内容,你可以改成保存到文件等操作 content = await resp.read() print(f"✅ 完成下载: {url} (大小: {len(content)} bytes)") else: print(f"❌ 下载失败: {url},状态码: {resp.status}") except Exception as e: print(f"⚠️ 下载出错: {url},错误信息: {str(e)}") async def worker(session, queue): """worker协程:循环从队列取任务处理""" while True: # 从队列获取URL,队列为空时会阻塞等待 url = await queue.get() try: await download_one(session, url) finally: # 标记当前任务完成,让队列知道这个任务已处理完毕 queue.task_done() async def main(urls, max_concurrent=16): # 初始化队列,将所有URL放入队列 queue = asyncio.Queue() for url in urls: queue.put_nowait(url) # 创建aiohttp会话(复用连接,提升下载效率) async with aiohttp.ClientSession() as session: # 启动指定数量的worker任务 worker_tasks = [asyncio.create_task(worker(session, queue)) for _ in range(max_concurrent)] # 等待队列中所有任务都被处理完成 await queue.join() # 所有任务处理完毕后,取消所有worker(因为worker是无限循环) for task in worker_tasks: task.cancel() # 等待所有worker任务结束 await asyncio.gather(*worker_tasks, return_exceptions=True) if __name__ == "__main__": # 示例URL列表,替换成你的大量URL urls = [f"https://example.com/page/{i}" for i in range(1000)] asyncio.run(main(urls))
关键部分说明
asyncio.Queue:作为任务缓冲区,自动处理任务的分发和等待,无需手动管理任务的启动时机worker协程:通过无限循环持续从队列取任务,确保只要队列有任务就会被处理queue.join():阻塞直到队列中所有任务都被标记为task_done(),也就是所有URL都处理完成- 取消worker:因为worker是无限循环,所以在所有任务完成后必须手动取消这些任务,避免程序无法退出
内容的提问来源于stack exchange,提问作者masroore
相关产品推荐
相关产品推荐

