
1. 任务背景两天 3000 万张怎么算都是不够用的时间有一天我接到一个很直接的需求大概有三千万张图片链接希望在两天内都拉回到本地做后续处理。3000w / 2 天听上去只是个大数但如果换成吞吐指标就很扎眼了•2 天 ≈ 172800 秒•3000 万张图 / 172800 秒 ≈ 173 张/秒这是“平均速度”。考虑到网络波动、失败重试、IO 抖动实际脚本得跑到三四百张/秒才有安全余量。当时我脑子里只有一个念头先写个能跑起来的再考虑优雅不优雅。2. 第一版最朴素的单线程 requests 脚本一开始我没想那么多就用最熟悉的组合requests for 循环。 大概长这样伪代码for url in urls: resp requests.get(url, timeout5) if resp.status_code 200: with open(filename, wb) as f: f.write(resp.content)这种写法的结果不出意料•CPU 几乎是闲着的大部分时间堵在网络 IO 上•并发为 1吞吐量非常有限•初步测了一下几百毫秒一张图离目标差了一个数量级这版脚本干了一件事 帮我证明了**“一条管道肯定不够”**。3. 第二版上线程池看起来“很并发”实际提升有限接着很自然地会想到上多线程。于是我把下载逻辑丢进ThreadPoolExecutor里每次开几十甚至上百个线程。吞吐是上去了但跑了半天我发现几个问题1.线程数量一上来切换开销越来越大2.requests 是同步 IO本质上还是在等网络3.某些站点对高并发不太友好开始出现各种 timeout、断连监控了一圈现象大概是•CPU 利用率时高时低没真正“吃满”•带宽利用率也不稳定•线程开太猛的时候目标站点明显“心情不好”这时候我开始有点意识到我不是简单缺“更多线程”而是没想清楚整个下载管线的瓶颈在哪里。4. 中途停下来先把瓶颈想清楚我做了几件事1.算账现在一秒大概能下多少张按这个速度全部跑完需要多久2.看资源利用率 CPU 是否始终在 70% 带宽是否接近上限3.看失败率和重试 是不是大量时间都花在 timeout 上 有没有被远端限流的迹象结果非常“典型”•CPU 远没吃满•带宽也没顶满•很多时间浪费在“排队等待网络返回 文件写入”上换句话说性能上不去并不是因为“机器不行”而是因为我没把机器喂饱。这一段反思之后我给自己定了几个设计目标1.尽量用 异步 IO 把网络等待时间利用起来2.结合 多进程突破单进程调度能力把 CPU 用满3.必须支持 断点续传中途挂了不能重新来过4.能通过几个参数来调优并发数、进程数、超时时间等最后的那份代码就是围绕这几个目标改出来的。5. 最终方案的整体结构多进程 进程内异步先说大思路再讲代码细节。这次我把流程拆成三层1.主进程层main 负责读 CSV、清洗 URL、过滤已下载文件、切任务 输出一个“干净的 URL 列表”再按 CPU 数量拆分成 N 份2.子进程层process_images 每个进程拿到自己的一段 URL 列表 在进程内再按批次拆分比如每批 1000 条 每批丢给异步事件循环去跑3.进程内异步层download_images / download_image 用 aiohttp 做 HTTP 请求 用 aiofiles 做文件写入 用 asyncio.Semaphore 控制并发 支持多代理轮询做简单的容灾结构图用一句话描述就是主进程做“分工”多进程做“吞吐”异步做“排队调度”。下面按这三层把代码里的关键点捋一下。6. 主进程数据清洗 任务拆分主流程入口是if __name__ __main__: main(urls_12.csv)main()做的事情可以概括为三块。6.1 清洗 URLdf pd.read_csv(csv_file) df.dropna(subset[url], inplaceTrue) df.drop_duplicates(url, inplaceTrue)•去掉空 URL•对 URL 去重这一步能避免无效请求浪费带宽。6.2 过滤已下载和空文件downloaded_images set(os.listdir(image_name)) image_urls df[url].tolist()再根据本地文件判断•如果文件不存在 → 需要下载•如果文件存在但大小为 0 → 当作上次失败需要重下这样脚本可以安全地反复执行天然支持断点续传。6.3 按 CPU 数量切块多进程处理num_processes cpu_count() num_per_process (num_urls num_processes - 1) // num_processes把 URL 列表均分成 N 份for i in range(num_processes): start_idx i * num_per_process end_idx (i 1) * num_per_process batch_urls image_urls[start_idx:end_idx] pool.apply_async(process_images, args(batch_urls,))到这里主进程的任务就结束了把清洗好的 URL 列表切成几块交给几个子进程各自去处理。7. 子进程负责“自己这一摊”的 URL 列表每个子进程进来的入口是process_images(image_urls)。这里我又做了一层批次拆分batch_size 1000 num_batches (len(image_urls) batch_size - 1) // batch_size原因很简单•一次性把几十万 URL 都塞进一个事件循环不太好控制•分批可以更好观察进度也便于在日志上排查问题每一批都通过loop.run_until_complete(download_images_batch(batch_urls))把当前这批的下载任务跑完。8. 进程内异步用 aiohttp 把网络空闲时间“抠干净”真正的下载发生在download_images()和download_image()里。8.1 并发控制Semaphore在每个进程里先定一个并发上限semaphore asyncio.Semaphore(20)这个值我实际跑下来大概是这样的经验•家庭网络、出口带宽一般10–20 比较安全•机房环境、带宽足可以一点点往上调观察 timeout 比例8.2 批量创建下载任务async with aiohttp.ClientSession(...) as session: tasks [] for url in batch_urls: filename ... task asyncio.ensure_future( download_image(session, url, filename, semaphore) ) tasks.append(task) for task in asyncio.as_completed(tasks): await task为什么用as_completed而不是gather•gather 某个任务异常时容易“一锅端”•as_completed 谁先完成谁返回更适合大量任务的下载场景8.3 单张图片下载逻辑download_image()里有几个关键点1.轮询代理列表2.进入请求前先 semaphore.acquire()结束后 release()3.用 aiofiles 写文件避免磁盘写阻塞事件循环4.捕获 TimeoutError、ServerDisconnectedError 等常见异常打印日志用于后续排查整体思路是把“等待网络返回”和“写磁盘”的时间让给事件循环来调度而不是锁死在一个线程里干等。9. 性能和调优从“理论上够”到“实测可以”最后说一下效果和调参过程。•在一台网络条件还可以的服务器上这套方案可以达到每分钟 1.7w 张左右只是量级参考•CPU 利用率可以保持在一个比较健康的区间•带宽利用率明显比多线程版本更稳定实践里我主要调了这几个参数1.进程数 CPU 核心数一般不需要太花哨2.每进程 semaphore 大小这个最敏感要结合网络和目标站点承受能力慢慢试3.timeout超时时间太短会导致大量无效重试太长又拖累整体吞吐4.batch_size批次大小主要影响内存占用和日志颗粒度调到一个比较满意的点之后整个任务能在预期时间内跑完也可以在中间停掉后重新启动继续跑。10. 小结回头看这次下载 3000w 图片的过程其实就是一步一步把**“直觉上的快”修正为“指标上的够用”**1.先写了一个单线程版本证明“肯定不够快”2.上了多线程发现问题不再是“线程多不多”而是资源利用率不均衡3.停下来算账、看监控从现象中找到真正的瓶颈4.重新设计多进程负责“并行度”异步负责“把 IO 缝隙填满”断点续传保障任务稳定性现在这份脚本不是“最优雅”但对“两天内把几千万张图片拉下来”这个具体目标来说它够用也容易维护。附完整源码import aiohttp import aiofiles import asyncio import os import pandas as pd from multiprocessing import cpu_count, Pool image_name images_12 proxies [ http://xxx:xx:xxxx, ] HEADERS { User-Agent: Mozilla/5.0 (Macintosh; Intel Mac OS X 10.15; rv:109.0) Gecko/20100101 Firefox/112.0, Accept: text/html,application/xhtmlxml,application/xml;q0.9,image/avif,image/webp,*/*;q0.8, Accept-Language: en-US, Accept-Encoding: gzip, deflate, br, Connection: keep-alive, Upgrade-Insecure-Requests: 1, Sec-Fetch-Dest: document, Sec-Fetch-Mode: navigate, Sec-Fetch-Site: none, Sec-Fetch-User: ?1 } async def download_image(session, url, filename, semaphore): for proxy_url in proxies: try: await semaphore.acquire() # 获取信号量 # async with session.get(url, timeout3) as response: # 不使用代理启用这行 async with session.get(url, timeout3, proxyproxy_url, headersHEADERS) as response: # 使用代理启用这行 if response.status 200: data await response.read() if data: async with aiofiles.open(filename, wb) as f: await f.write(data) break except aiohttp.ServerDisconnectedError: print(fServer disconnected when downloading {url}, retrying...) except asyncio.TimeoutError: print(f{url}: TimeoutError) except Exception as e: print(fAn error occurred when downloading {url}: {e}) finally: semaphore.release() # 释放信号量 async def download_images(batch_urls, semaphore): connector aiohttp.TCPConnector(limitNone) # 不限制连接池大小 async with aiohttp.ClientSession(timeout5, connectorconnector) as session: tasks [] for url in batch_urls: filename os.path.join(image_name, os.path.basename(url)).jpg task asyncio.ensure_future(download_image(session, url, filename, semaphore)) tasks.append(task) # 使用asyncio.as_completed迭代已完成的任务 for task in asyncio.as_completed(tasks): await task def process_images(image_urls): async def download_images_batch(batch_urls): semaphore asyncio.Semaphore(20) # 控制并发请求数量 await download_images(batch_urls, semaphore) loop asyncio.get_event_loop() # 分批下载 batch_size 1000 # 每批下载的数量 num_batches (len(image_urls) batch_size - 1) // batch_size for i in range(num_batches): start_idx i * batch_size end_idx (i 1) * batch_size batch_urls image_urls[start_idx:end_idx] # 异步下载图片 loop.run_until_complete(download_images_batch(batch_urls)) def main(csv_file): # 创建存储图片的目录 os.makedirs(image_name, exist_okTrue) # 读取 CSV 文件 df pd.read_csv(csv_file) # 一百万条记录 print(f读取csv文件总记录数{len(df)}) # 删除空的记录 df.dropna(subset[url], inplaceTrue) print(f去除空记录数后总数{len(df)}) # 对 df 中 url 列去重 df.drop_duplicates(url, inplaceTrue) print(f去重后总数{len(df)}) # 过滤掉已经下载的图片链接 downloaded_images set(os.listdir(image_name)) image_urls df[url].tolist() overwrite_urls [url for url in image_urls if url.split(/)[-1] .jpg in downloaded_images] print(已经下载图片数量: , len(overwrite_urls)) # 过滤掉已经下载的非零字节图片链接 image_urls [ url for url in image_urls if url.split(/)[-1] .jpg not in downloaded_images or os.path.getsize(os.path.join(image_name, url.split(/)[-1] .jpg)) 0 ] print(未下载图片数量: , len(image_urls)) # 使用多进程进行并行下载 num_processes cpu_count() with Pool(num_processes) as pool: num_urls len(image_urls) num_per_process (num_urls num_processes - 1) // num_processes # 每个进程处理一部分图片链接 results [] for i in range(num_processes): start_idx i * num_per_process end_idx (i 1) * num_per_process batch_urls image_urls[start_idx:end_idx] # 启动一个进程执行下载任务 result pool.apply_async(process_images, args(batch_urls,)) results.append(result) # 等待所有进程完成 for result in results: result.get() if __name__ __main__: semaphore: 默认20 信号量控制每个进城内的并发请求数量,根据网络情况调整 timeout: 默认5秒 超时时间根据网络情况调整 可以先打印控制台查看报TimeoutError的数量然后根据网络情况调整,可以去掉打印 使用另外起一个服务来查看图片数量 import asyncio import os image_name images_1 async def print_image_count(): while True: image_count len(os.listdir(image_name)) print(f当前图片数量{image_count}) await asyncio.sleep(60) # 每分钟打印一次 loop asyncio.get_event_loop() task loop.create_task(print_image_count()) loop.run_forever() 如果TimeoutError多的话降低semaphore的值如果TimeoutError基本没有的话可以适当增加semaphore的值 不需要代理的话可以注释掉使用代理 目前默认配置在网速良好的服务器测试速度是 一分钟下载1.7w张图片左右 main(urls_12.csv)