Python协程实战:将IO密集型爬虫采集从49秒优化到2.5秒

Python协程实战:将IO密集型爬虫采集从49秒优化到2.5秒 同一个爬虫脚本同步版跑 20 个接口要 49 秒用 Python 协程改成异步并发后只用了 2.5 秒。这不是网络变快了也不是换了一台性能更强的机器只是把代码里那些傻等的时间重新用了起来。如果你平时用 Python 写爬虫、批量调第三方接口、处理大量小文件读写或者做各类数据采集工作应该体会过那种跑起来 CPU 占用率极低但就是一直卡着不出结果的难受。这类任务就是典型的 IO 密集型任务——代码的大部分时间不是在计算而是在等待输入输出完成等待期间 CPU 基本是空的。而 Python 协程恰恰就是用来解决这种空等浪费的。今天这篇不聊太深的理论我会从一次真实的采集卡顿经历切入讲清楚 IO 密集型任务为什么卡、协程为什么能解决这个问题、改造一个爬虫脚本需要动哪些代码再把选型对比和实战里最容易踩的四个坑一并梳理一遍。无论你是刚接触协程还是已经写了一段时间 asyncio这篇文章应该都能帮你少走弯路。1. 先看清卡顿的根源IO 密集型任务的时间都花在等待上1.1 一次真实的数据采集卡顿经历我长期写爬虫和服务端脚本有一次做市场数据采集需要逐个请求 15 个城市、每个城市 5 个维度共 75 个 HTTP 接口。第一版用同步代码实现跑一次要 6 分钟。让我印象最深的不是时间长而是资源占用特别低跑的时候 CPU 占用率不到 10%内存也很低程序就像挂在树上打盹一样明明什么活都没干却在慢慢熬时间。这个现象的原因就是同步 HTTP 请求的流程。代码发出请求后大部分时间都阻塞在 socket 接收数据上这个等待过程是纯粹的空转。我第一个想法是用多线程叠了 20 个线程后时间确实降到了 90 秒但当并发数继续往上加线程切换的开销和内存占用开始明显增长。后来把所有请求换成 Python 协程加 aiohttp采集耗时直接从 6 分钟降到 25 秒这个提升才是彻底解决级别的。这个案例非常典型。它说明当你的程序大量时间花在等待网络、磁盘、数据库等输入输出完成上无论机器性能多好同步代码的排队等待都会成为最大的瓶颈。而 Python 协程的设计目标恰好就是消除这种空等。1.2 怎么判断一个任务是不是 IO 密集型判断标准很简单我在实际项目里就两步任务跑起来之后打开任务管理器或资源监视器盯一下 CPU 占用率。如果任务执行期间 CPU 占用率长期低于 20%那大概率是 IO 密集型任务如果 CPU 占用率稳定在 80% 到 100%那它是 CPU 密集型任务。IO 密集型任务的特征非常明显程序大部分时间阻塞在系统调用上等待外部设备返回数据。典型例子包括爬虫等待 HTTP 响应、批量读写文件等待磁盘吞吐、数据库批量查询等待连接与结果集返回、调用大模型 API 等待网络传输、读取流式行情数据等待网络包到达。CPU 密集型任务则相反程序的大部分时间在计算需要大量算术、逻辑、内存操作比如图像转码、加密解密、数据聚合、向量计算。明白自己的任务是哪种类型是先于写代码的决策。类型判断错了后面所有并发选型都会跟着错。1.3 多线程的不足催生了协程很多朋友会疑惑Python 不是有多线程吗确实多线程在并发场景有一定效果尤其在等待 IO 的一段窗口里操作系统会切换到另一个线程运行。但它并不是 IO 密集任务的终极解原因有三个。第一线程是操作系统级别的调度对象每创建一个线程都要消耗内核资源生产环境里线程数超过几千以后上下文切换开始明显吃掉性能。第二多个线程修改共享变量时要加锁加锁的粒度、顺序稍有不慎就是死锁或数据不一致调试成本很高。第三线程之间通信用 queue 或者共享内存配合条件变量代码复杂度会显著上升。协程恰恰解决了这些痛点协程跑在用户态切换成本极小协程之间不需要操作系统参与调度在同一个线程内它们共享上下文也更自然。这就是在高并发 IO 场景下协程逐渐成为首选的核心原因。当然协程也有自己的学习门槛和需要注意的地方这个我们在后面的实战与避坑部分细说。2. 协程的运行机制事件循环、async/await 和 Task2.1 事件循环让一个线程处理成千上万个连接Python 协程并不是多线程的替代品这么简单。它背后是事件循环机制。事件循环是一个持续运行的调度器它维护一个就绪任务列表反复执行这样一套逻辑从就绪队列取出一个任务运行它直到遇到 await如果 await 的底层 IO 未完成把任务挂起交给系统监听IO 完成后把任务重新放入就绪队列循环往复。可以把事件循环比作餐厅里的服务员客人点完菜后服务员不会一直站在那一桌等后厨出菜而是继续去接待下一桌等出菜好了再回头上菜。正是这种不等菜的机制让一个线程可以同时服务上千个连接。你在很多高性能框架里看到的单进程支持数万连接底层基本都是这一套东西。2.2 async 和 await 的真实含义很多初学者看过 async/await 的文档但不太理解执行模型。我总结成四条规则记住了基本就能看懂协程代码。async def 声明一个协程函数。调用这个函数不会立即执行函数体而是返回一个协程对象。协程对象可以被 await。await 会进入函数体并执行。在协程函数内部await 关键字后面必须是一个可等待对象。典型可等待对象包括协程对象、Task、Future。执行到 await 的瞬间协程可能挂起让出执行权给事件循环当被等待的对象完成时协程从这里恢复执行。举一个直观的对比同步代码resp requests.get(url)在收到响应之前当前线程被阻塞什么事都干不了协程代码resp await session.get(url)在收到响应之前当前协程挂起但事件循环可以跑去执行其他协程任务。这一挂一跑就是并发的本质。2.3 Task 与 Future协程任务的执行单位Task 是把协程对象包装成可被事件循环调度的任务的产物。每当你调用asyncio.create_task(coro)事件循环就知道可以调度这个协程了。可以这样理解协程对象是一份待执行的菜谱Task 是后厨已经接单、正在排队等待制作的订单Future 是这张订单上写的预计可以上菜的时间点。实际写代码时最常遇到的是 Task。asyncio.gather(*tasks)用于同时等待多个任务asyncio.wait_for(task, timeout)用于给单个任务加超时asyncio.shield(task)用于防止任务被外部取消。这些函数在后面的实战代码里会反复出现先混个脸熟。2.4 协程之间真的不需要加锁吗协程运行在同一线程里切换只发生在 await 处。因此只要你保证在不含 await 的连续代码段里修改变量这个操作对你是原子的不会看到改了一半的数据。这就是很多协程代码看起来比多线程代码简洁的原因。但注意这不代表完全不需要同步机制。如果两个协程交替执行时有先后依赖或者你在 await 前后对共享变量做了非原子的读改写操作仍然可能得到错误结果。比如一个计数器在多个协程里做count 1如果的中间被 await 打断就会出现并发问题。所以在协程里设计共享状态时依然要有临界区的意识不能完全放飞自我。3. 实战把阻塞式采集脚本改造成协程并发脚本3.1 HTTP 客户端选型为什么从 requests 换成 aiohttprequests 是 Python 社区最常用的 HTTP 库但它是同步阻塞的。在协程函数里直接调用 requests一旦发出请求底层 socket 会阻塞整个操作系统线程事件循环被卡死其他协程一个也跑不动。所以协程改造的第一步就是换一个异步 HTTP 客户端。推荐 aiohttp它的用法和 requests 很接近区别在于session.get(url)需要 await响应对象需要用async with管理生命周期连接池建立在事件循环之上可以自动复用 TCP 连接。如果你的项目里大量代码已经基于 requests 编写想平滑过渡也可以试试 httpx它保留了 requests 式接口同时提供了异步客户端。不过在爬虫和高并发采集场景我更推荐 aiohttp因为它还能承担 WebSocket 客户端和服务端的角色生态更完整。3.2 同步版本与耗时基线先看一份最基础的数据采集同步代码。我这里用 httpbin.org 的延迟接口模拟真实网络中的响应耗时每个请求延迟 0 到 4 秒不等一共 20 个请求。import time import requests # 模拟 20 个响应耗时不同的接口 URLS [fhttps://httpbin.org/delay/{i % 5} for i in range(20)] def fetch(url): resp requests.get(url, timeout10) return resp.status_code def main(): t0 time.perf_counter() for url in URLS: fetch(url) elapsed time.perf_counter() - t0 print(f同步版本总耗时: {elapsed:.2f}s) main()运行结果很直白同步版本总耗时: 49.18s20 个请求串行执行平均每个请求等待 2.5 秒左右整个过程没有任何计算负载纯排队等待。3.3 协程版本与性能对比接下来是协程版本。核心思路是把 20 个请求一次性交给事件循环并发执行等所有请求都回来以后再统一处理结果。import asyncio import aiohttp import time URLS [fhttps://httpbin.org/delay/{i % 5} for i in range(20)] async def fetch(session, url): async with session.get(url, timeout10) as resp: return resp.status async def main(): async with aiohttp.ClientSession() as session: tasks [asyncio.create_task(fetch(session, url)) for url in URLS] results await asyncio.gather(*tasks) print(results) t0 time.perf_counter() asyncio.run(main()) elapsed time.perf_counter() - t0 print(f协程版本总耗时: {elapsed:.2f}s)运行结果[200, 200, 200, 200, 200, 200, 200, 200, 200, 200, 200, 200, 200, 200, 200, 200, 200, 200, 200, 200] 协程版本总耗时: 2.51s从 49 秒到 2.5 秒差不多 20 倍的提升。原因不复杂20 个请求几乎在同一时刻发出等待时间是并行的总耗时约等于最慢的那个请求耗时而不是 20 个请求耗时之和。3.4 信号量把并发限制在合理范围上一版示例依赖 httpbin 这个免费测试服务实际工程中我们不能毫无限制地盲目并发。这里有两方面的原因一是目标服务器可能扛不住出现大量连接拒绝甚至封 IP二是本机可用的文件描述符数量是有限的socket 连接不能无限创建。此时用asyncio.Semaphore来限制最大并发数import asyncio import aiohttp MAX_CONCURRENT 5 semaphore asyncio.Semaphore(MAX_CONCURRENT) async def fetch(session, url): async with semaphore: async with session.get(url, timeout10) as resp: return resp.status信号量本质上是一个令牌桶同一时刻只有 N 个协程能拿到令牌进入请求代码其余协程在信号量上排队等待。这个 N 一般根据目标服务器的承受能力和你的网络环境动态调整5 到 20 是常见的经验区间。不要一上来就设 1000先压测看延迟和错误率再慢慢往上调。3.5 超时、重试与失败隔离生产环境的必备补丁并发请求一旦多起来网络抖动、上游服务不稳定都会暴露出来。生产环境里的采集脚本质量往往取决于异常处理而不是并发模型本身。这里给出一个可复用的请求模式把超时和重试都包进去。import asyncio import aiohttp class FetchError(Exception): pass async def fetch_with_retry(session, url, retries3, timeout5): last_exc None for attempt in range(retries): try: async with asyncio.wait_for(session.get(url, timeouttimeout), timeouttimeout) as resp: if resp.status 200: return await resp.text() last_exc FetchError(fHTTP {resp.status}) except asyncio.TimeoutError: last_exc FetchError(ftimeout on attempt {attempt 1}) except aiohttp.ClientError as exc: last_exc exc await asyncio.sleep(0.5 * (attempt 1)) raise last_exc这里有几个细节值得注意asyncio.wait_for负责兜底整个请求流程防止某个连接长时间悬挂每次失败后的等待时间用 0.5 秒、1 秒、1.5 秒阶梯式递增避免重试风暴所有异常被收拢成统一的 FetchError上层拿到异常后可以做失败隔离不影响其他协程的结果收集。4. 方案选型协程、多线程、多进程到底怎么选4.1 一张表看清三者的核心区别经常有朋友问协程、多线程、多进程到底应该用哪个我把它们的核心区别整理成一张表供实际选型时参考。维度协程asyncio多线程多进程适用场景IO 密集型IO 密集型CPU 密集型、IOCPU 混合切换方式用户态协作式内核态抢占式内核态进程切换创建开销极小中等较大内存开销小中等大数据共享同线程内较安全需加锁、队列需 IPC、共享内存代码复杂度中需要异步思维中需处理锁中需处理序列化最大并发量数万级别几千级别受限于 CPU 核数选择的基本逻辑很简单纯 IO 密集、并发量要求高、追求资源占用低选协程IO 密集但代码已经基于线程模型写好、不想大改可以继续用多线程CPU 密集或两者混合优先考虑多进程或者协程加进程池的组合。4.2 GIL 到底拦住了什么GIL 是 Python 里的全局解释器锁它的存在保证了同一时刻只有一个线程在执行 Python 字节码。这意味着即使你有 4 个 CPU 核4 个线程做纯计算任务也不会同时跑满 4 个核它们轮流使用同一个核。这就是为什么多线程对 CPU 密集型任务几乎无效。但在 IO 密集场景GIL 的影响没有想象中那么大线程在等待网络、磁盘 IO 时会主动释放 GIL让其他线程进来执行。所以用多线程处理 IO 密集任务确实有效只是线程数膨胀后系统调用与上下文切换的开销会快速增大。协程的优势在于它的切换只需要保存和恢复少量寄存器状态不触发系统调用开销比线程切换小一个数量级。这也是在需要上万级并发连接时协程几乎成为不二之选的原因。4.3 混搭协程与线程池、进程池的组合拳工程上全用协程往往不是最优解。我遇到过不少任务大部分时间是网络 IO适合协程但拿到数据后要做一次 CPU 密集的清洗和特征计算这部分如果直接在协程里做会阻塞事件循环把整个异步性能拖垮。正确的做法是协程负责高并发 IOCPU 密集部分丢进线程池或进程池执行。用asyncio.to_thread或者loop.run_in_executor都能实现。import asyncio from concurrent.futures import ProcessPoolExecutor def cpu_intensive(data): # 这里做一些重计算比如文本清洗、特征提取 return heavy_compute(data) async def main(): loop asyncio.get_running_loop() with ProcessPoolExecutor(max_workers4) as pool: result await loop.run_in_executor(pool, cpu_intensive, data) return resultrun_in_executor可以把同步阻塞函数丢进线程池或进程池执行事件循环本身不受影响。这是很多生产级 Python 服务的标准做法异步协程处理连接和 IO进程池处理 CPU 密集计算两边各司其职。如果维护的是 Django 或 Flask 老项目不想做大的异步化改造也可以用线程池配合 requests 提升并发虽然效果不如纯 asyncio 加 aiohttp 极致但改动最小。5. 协程实战中绕不开的四个坑5.1 坑一协程里用了同步阻塞调用整个事件循环停摆这是每个协程新手都会踩的坑。你在一个 async 函数里写了time.sleep(1)看似是等待 1 秒实际上它会让当前线程真正地睡过去事件循环整个停摆所有其他协程都无法推进。同步阻塞调用包括time.sleep、requests.get、普通文件对象的read。它们不会等待事件循环的调度一旦调用就会占用整个线程。正确的姿势是用异步版本time.sleep换成await asyncio.sleeprequests.get换成await session.get普通文件读取换成 aiofiles 的异步接口。记住一条核心原则在协程里只允许异步阻塞不允许同步阻塞。异步阻塞如await asyncio.sleep会释放控制权给事件循环同步阻塞如time.sleep会卡死整个循环。排查这类问题有一个技巧如果一个异步脚本在某个地方长时间无响应而且 CPU 占用率极低优先检查是不是哪个协程函数里漏网了同步阻塞调用。5.2 坑二Jupyter 和交互式环境里 asyncio.run 重复报错asyncio.run是一个很便捷的入口但它有一个限制同一事件循环只能运行一次而且在已经运行着事件循环的环境里调用会报错。典型场景是 Jupyter Notebook第一次运行asyncio.run(main())没问题第二次运行就抛出类似 Event loop is closed 的异常。这是因为 Notebook 的内核在后台已经启动了一个事件循环再次 create 新循环会冲突。解决办法是手动创建并设置事件循环绕开asyncio.run的内部检查import asyncio loop asyncio.new_event_loop() asyncio.set_event_loop(loop) loop.run_until_complete(main())在交互式环境下这种写法比反复调用asyncio.run稳定得多。如果你使用的是老版本的 Python 环境也可以用nest_asyncio.apply()补丁允许嵌套循环但从可靠性来说new_event_loop的方式更干净。5.3 坑三超时控制不彻底任务悄悄悬挂很多人给协程加了asyncio.wait_for就觉得万事大吉但忽略了超时后的资源清理。wait_for在超时后会自动取消被包装的 Task但如果你在任务里手动捕获了CancelledError并且没有正确抛出任务可能无法真正停止连接和文件句柄就会泄漏。正确的异常处理姿势是这样import asyncio async def worker(): try: await asyncio.sleep(10) except asyncio.CancelledError: # 完成必要的清理工作后重新抛出确保任务真正取消 raise async def main(): try: result await asyncio.wait_for(worker(), timeout2) except asyncio.TimeoutError: print(任务超时已取消)清理工作时要注意不要在except asyncio.CancelledError里做耗时的操作也不要吞掉异常。asyncio 的取消机制依赖 CancelledError 在任务内部正确传播一旦被吞掉取消就失效了。这个坑在开发长连接服务时最致命因为表面上看任务取消了实际上连接还挂在事件循环上。5.4 坑四协程对象忘记 await静默变成警告新手很容易写出这种代码async def hello(): return hi def main(): print(hello()) # 这里只打印了一个协程对象地址调用hello()只是创建了一个协程对象它还没有被调度执行。必须 await 它或者用asyncio.create_task包起来。如果你忘了 await运行时通常会打印一条警告RuntimeWarning: coroutine hello was never awaited。生产环境建议把 warning 提升为 error让这类问题在测试阶段就暴露出来import warnings warnings.simplefilter(error, RuntimeWarning)在测试入口加上这一行任何忘写 await 的地方都会直接报错而不是留下一个静默的坑。这个习惯我保持了很久大大减少了协程代码的隐性 bug。结语我在协程工程化上的一点体会把协程从能跑通变成能上线我前后踩了不少坑最终形成了一套比较保守但稳定的工作流先写一个能跑的同步版本明确耗时基线和瓶颈再改造成异步版本一次只改一层避免协程和阻塞调用混在一起每一层都加上超时、限流和重试上线前先用小并发数压一遍确认目标服务跟得上再逐步提高并发。这套流程看起来很慢但其实是协程工程化落地最稳的路径。还有一个很实用的建议不要为了用协程而用协程。如果你的并发量只有十几个多线程已经能解决如果任务是 CPU 密集协程帮不上忙。只有在明确属于 IO 密集型、并发生数量级较高、资源占用敏感的场景里协程才能真正发挥出那种让人眼前一亮的性能优势。先判断类型再做选型最后写代码这个顺序一定不要颠倒。