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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 16:07:32