提示

返回博客列表

Python 并发下载器设计与实战:从同步到 asyncio + 线程池的吞吐优化

Python 并发下载器设计与实战:从同步到 asyncio + 线程池的吞吐优化

做视频工具最怕的不是"下不下来",而是"下一个要等半天、下一堆直接把机器拖垮"。本文用一个可取消的并发下载器为例,讲清楚同步 → 线程池 → 协程的演进,以及限速、重试、断点续传和背压这些生产级细节。

一、同步下载的瓶颈

最朴素的写法:

import requests

def download(url, path):
    r = requests.get(url, timeout=30)
    r.raise_for_status()
    with open(path, 'wb') as f:
        f.write(r.content)

问题:requests.get 是阻塞的,下载一个 200MB 文件期间线程干等网络,CPU 几乎闲置;下 100 个文件就是 100 倍串行时间。本质是等 IO 的时间被白白浪费

二、方案选型:线程 vs 协程

维度 线程池 (ThreadPool) 协程 (asyncio+aiohttp)
适用场景 阻塞库(requests/文件IO) 原生异步网络
并发上限 受 GIL,约数十~数百 数千级
编程复杂度 低,改动小 中,需 async/await 全链路
切换开销 较高 极低

结论:下载是 IO 密集型且多用 requests 时,线程池最省事;纯异步管线用 asyncio

2.1 requests + ThreadPoolExecutor

from concurrent.futures import ThreadPoolExecutor, as_completed
import requests

def download(url, path):
    with requests.get(url, stream=True, timeout=30) as r:
        r.raise_for_status()
        with open(path, 'wb') as f:
            for chunk in r.iter_content(8192):
                f.write(chunk)

urls = [('http://x/a.mp4', 'a.mp4'), ...]
with ThreadPoolExecutor(max_workers=8) as ex:
    tasks = [ex.submit(download, u, p) for u, p in urls]
    for fut in as_completed(tasks):
        fut.result()  # 抛出来任何异常

stream=True + iter_content 避免把整文件读进内存;max_workers=8 控制并发。

2.2 aiohttp + asyncio

import asyncio, aiohttp

async def download(session, url, path):
    async with session.get(url) as r:
        r.raise_for_status()
        with open(path, 'wb') as f:
            async for chunk in r.content.iter_chunked(8192):
                f.write(chunk)

async def main(urls):
    async with aiohttp.ClientSession() as s:
        await asyncio.gather(*(download(s, u, p) for u, p in urls))

asyncio.gather 并发调度,单线程内靠事件循环切换,吞吐远高于线程池。

三、设计一个可取消的并发下载器

生产环境要能:随时取消、限速、失败重试、断点续传。核心抽象是"任务 + 信号量限速 + 事件取消"。

import asyncio, aiohttp

class Downloader:
    def __init__(self, max_concurrent=8, rate_limit=0):
        self.sem = asyncio.Semaphore(max_concurrent)
        self.cancel = asyncio.Event()

    async def one(self, session, url, path):
        if self.cancel.is_set():
            return
        async with self.sem:               # 并发上限 = 背压
            async with session.get(url) as r:
                r.raise_for_status()
                with open(path, 'wb') as f:
                    async for chunk in r.content.iter_chunked(8192):
                        if self.cancel.is_set():
                            return
                        f.write(chunk)

    def stop(self):
        self.cancel.set()

信号量 Semaphore 就是背压:并发数封顶,避免把网卡/磁盘打满。

3.1 限速

最简单的令牌桶:每写 N 字节 await asyncio.sleep(t) 平滑速率,或直接用 asyncio.throttle/第三方限流库。脚本级可用 pv -L 1m 包一层限速。

3.2 重试、断点续传与校验

  • 重试:捕获 ClientError/Timeout,指数退避 2**i 秒,最多 3 次;
  • 断点续传:headers={'Range': 'bytes=%d-' % size} 从已下载长度续传,配合 mode='ab'
  • 校验:下载完比对服务端 Content-LengthETag/SHA256,不一致则重下。
async def download_resume(session, url, path, max_retry=3):
    for i in range(max_retry):
        start = os.path.getsize(path) if os.path.exists(path) else 0
        headers = {'Range': f'bytes={start}-'} if start else {}
        try:
            async with session.get(url, headers=headers) as r:
                r.raise_for_status()
                mode = 'ab' if start else 'wb'
                with open(path, mode) as f:
                    f.seek(start)
                    async for chunk in r.content.iter_chunked(8192):
                        f.write(chunk)
            return
        except Exception:
            await asyncio.sleep(2 ** i)
    raise RuntimeError('download failed: ' + url)

四、与 FFmpeg 流水线的衔接

下载完常常要转码/切片。注意不要让下载和转码同时占满磁盘 IO:用流水线思想,下完一个就丢给转码子进程,双阶段并发:

async def pipeline(url, raw, out):
    await download_resume(session, url, raw)
    proc = await asyncio.create_subprocess_exec(
        'ffmpeg', '-i', raw, '-c:v', 'libx264', out)
    await proc.wait()

或者更稳:用消息队列(Redis/RabbitMQ)把"下载""转码""上传"解耦成独立 worker,互不阻塞。

五、监控指标与背压

上线后至少盯三个数:

  • 并发在途数:长期等于上限说明下游是瓶颈,要调 max_concurrent
  • 单任务耗时分布:长尾多=个别源慢,加超时与重试;
  • 磁盘/网卡水位:超阈值就降 max_workers,这就是背压的落地。

六、常见坑

  1. 忘记设超时,requests 默认无超时会永久挂起 → 永远传 timeout=
  2. 异常被 gather 吞掉 → 在 one()try/except 记录并 result() 暴露。
  3. 并发过高打满 FD → 系统 ulimit -n 调大或降并发。
  4. 写入不是原子 → 下载中文件被读,用临时名 .part 下完再 rename
  5. GIL 下线程池不等于真并行 CPU → 重计算交给进程池 ProcessPoolExecutor

七、小结

并发下载器的演进主线是:用线程/协程把"等 IO"的时间并行起来,用信号量做背压,用重试+断点续传保可靠,用流水线把下载和转码解耦。掌握这套骨架,无论是视频批量备份还是媒体处理平台,吞吐和稳定性都能上一个台阶。VidDown 的批量下载能力正是建立在这套设计之上。

想亲手试试?用 VidDown 一键解析下载

粘贴视频链接即可解析,多平台支持、网页端即用;下载桌面客户端解锁海外平台本地解析,开通会员更享不限次下载。

顶部