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-Length或ETag/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,这就是背压的落地。
六、常见坑
- 忘记设超时,
requests默认无超时会永久挂起 → 永远传timeout=。 - 异常被
gather吞掉 → 在one()里try/except记录并result()暴露。 - 并发过高打满 FD → 系统
ulimit -n调大或降并发。 - 写入不是原子 → 下载中文件被读,用临时名
.part下完再rename。 - GIL 下线程池不等于真并行 CPU → 重计算交给进程池
ProcessPoolExecutor。
七、小结
并发下载器的演进主线是:用线程/协程把"等 IO"的时间并行起来,用信号量做背压,用重试+断点续传保可靠,用流水线把下载和转码解耦。掌握这套骨架,无论是视频批量备份还是媒体处理平台,吞吐和稳定性都能上一个台阶。VidDown 的批量下载能力正是建立在这套设计之上。