Python asyncio并发编程实战:从事件循环到协程的I/O密集型任务优化

Python asyncio并发编程实战:从事件循环到协程的I/O密集型任务优化 1. 为什么并发会成为 Python 开发的绕不开的话题先说一个很多初学者的误区认为 Python 程序跑得慢就是语言本身不行。实际上绝大多数业务系统不是被 CPU 算力卡住的而是被 I/O 等待卡住的。比如爬虫等接口响应、Web 服务等数据库返回、文件读写、网络请求90% 的时间都花在“等”上。而 asyncio 就是 Python 官方提供的、专门解决“等”这个问题的高效方案。asyncio 是 Python 3.4 起引入标准库的异步 I/O 框架核心价值是让你用写同步代码的方式获得并发处理 I/O 密集型任务的能力。它的适用人群很广写爬虫的、写 Web 后端的、写自动化脚本的、做量化交易做数据采集的只要涉及网络请求或者外部服务调用都值得掌握它。我在实际项目中用 asyncio 的最大体会是它能帮你在单线程内同时管理成百上千个等待中的 I/O 任务而不需要为每个连接创建一个线程。要知道线程是有成本的创建、切换、销毁都有开销一个进程里开上千个线程光是上下文切换就能把 CPU 拖垮。而 asyncio 的事件循环event loop在单线程内用协作式调度来切换任务任务切换的开销比线程切换小一个量级。不过 asyncio 也不是万能的。它解决的是 I/O 密集型并发问题不是 CPU 密集型并行问题。如果你有一堆死循环或者大量浮点运算那该用多进程还是多进程asyncio 帮不了你。这个边界如果不搞清楚后面踩坑会踩得很痛苦。需要说明的是我下面讲的不是照抄文档的翻译版而是我自己在爬虫、API 网关、WebSocket 推送服务这些真实场景里用出来的经验。代码都是可以直接跑的但更重要的是理解它背后为什么这么设计。2. 并发模型的横向对比线程、进程与事件循环2.1 三种并发方案的底层差异Python 开发中常用的并发方案有三类多线程threading、多进程multiprocessing、异步 I/Oasyncio。很多人一上来就问“哪个快”这其实是伪命题关键要看你的任务是 CPU 密集还是 I/O 密集。先说多线程。Python 的线程受限于 GIL全局解释器锁同一时刻只有一个线程能执行 Python 字节码。所以多线程在线程切换之间没法真正并行执行 Python 代码它适合的是 I/O 密集型场景——因为线程在等待 I/O 时会释放 GIL让其他线程跑。但如果任务是纯计算多线程基本是负优化因为切换线程本身还有额外开销。再说多进程。每个进程有独立的解释器和内存空间可以绕过 GIL 实现真正的并行计算适合 CPU 密集型任务。但进程的开销很大创建慢、内存占用高进程间通信IPC也比线程间共享变量麻烦得多。asyncio 的实现思路完全不同它不依赖操作系统的线程或进程调度而是在进程内自己维护一个事件循环event loop。事件循环里挂着一堆待执行的任务Task每个任务执行到 I/O 等待时就主动让出控制权await事件循环立刻切换到下一个就绪的任务继续执行。这种协程切换是用户态完成的不需要陷入内核态所以开销极小。为了更直观地理解三者的区别我看过的一张对比表一直记忆犹新维度threadingmultiprocessingasyncio适用场景I/O 密集CPU 密集I/O 密集、高并发连接利用多核受 GIL 限制不行可以不行单线程任务切换开销中等内核态切换高进程切换极低用户态切换共享数据方式共享内存需加锁IPC 或共享内存同线程内天然共享无需锁最大并发量受线程数限制受进程数限制可支持数万任务学习成本低中较高2.2 为什么说 asyncio 是 I/O 密集型场景的最优解举个最直观的例子写一个爬虫需要抓取 1000 个网页。如果用同步方式每个请求平均耗时 200 毫秒串行跑完需要 200 秒。用多线程Python 线程切换加上 GIL 竞争实际可能要 20 秒到 30 秒。用 asyncio理论上如果网络带宽和对方服务器都扛得住可以在 1 秒左右全部发起然后在 200 毫秒左右全部收到响应整体耗时接近单个请求的耗时。这个数量级的差距来自哪里不是 asyncio 本身“更快”而是它避免了大量线程的管理开销并且把“等待”的时间压缩到极致。网络 IO 发出请求后asyncio 不是干等着而是立刻去处理下一个请求。等第一个请求的响应到了事件循环会通过回调通知对应的任务继续往下走。我个人的选型经验是当单机需要管理的并发连接数超过 500 时多线程方案基本不可用asyncio 是唯一现实的选项。比如 WebSocket 推送服务、长轮询网关、消息推送中间件这类的场景asyncio 几乎是标准答案。当然asyncio 也有它的短板。最明显的是异步代码里一旦混入了一个阻塞调用比如用 requests 而不是 aiohttp整个事件循环就卡住了所有并发任务全部停摆。所以使用 asyncio 是“用纪律换效率”——得时刻留意代码里有没有潜在的阻塞操作。这个我在后面的常见问题部分会展开讲。3. asyncio 核心机制拆解从事件循环到协程3.1 事件循环到底是怎么工作的事件循环event loop是 asyncio 的心脏它的工作机制可以类比成一个餐厅里只有一个服务员顾客点完单后服务员不会站在每桌旁边等菜做好而是记下这桌的需求然后去服务下一桌。哪桌的菜好了后厨喊一声服务员再去上菜。这样一个人就能服务很多桌。对应到 asyncio 里任务Task就是顾客I/O 操作就是做菜事件循环就是服务员。代码里每次await一个 I/O 操作就相当于告诉事件循环“我现在要等一个外部结果你先去干别的活等结果好了来叫我。” 在 Python 3.10 及以后版本你甚至不需要手动获取和设置事件循环asyncio.run()会自动帮你创建和管理。来看一个最小示例import asyncio async def say_hello(): print(Hello) await asyncio.sleep(1) # 模拟一个 I/O 等待 print(World) asyncio.run(say_hello())这里async def声明了一个协程函数调用它不会立刻执行而是返回一个协程对象。asyncio.run()会创建一个事件循环然后把协程对象包装成 Task 投进事件循环里运行。await asyncio.sleep(1)的作用是让出控制权给事件循环1 秒后继续执行。看到这里你可能觉得这跟普通函数没啥区别因为只有一个任务时确实是顺序执行它的威力要等有多个并发任务时才显现。3.2 理解 await、Task 和 Future 的关系很多人学 asyncio 时被await、Task、Future 这几个概念绕晕了我尽量讲得形象一点。先说 Future它是异步编程里最底层的概念可以理解成一个“快递单号”。你下单后拿到一个单号这个单号本身不包含快递内容但它代表一个“将来会到的结果”。你查物流时快递没到单号就处于“未完成”状态快递到了单号才显示“完成”这时候你才能拿着单号取件。在 asyncio 里Future 就是一个这样的占位对象它表示一个异步操作的最终结果。Task 是 Future 的子类但它更具体——它包装了一个协程对象负责调度这个协程的运行。你可以理解为“快递公司的调度系统”它知道你这个包裹从哪个仓库发出、走哪条路线、什么时候能到。Task 对象被创建后会立刻被调度到事件循环里运行不需要你手动启动。await关键字的作用是接受一个可等待对象协程、Task、Future然后暂停当前任务把控制权还给事件循环直到这个可等待对象完成。它有点像“挂起线程”——只不过挂起的是协程开销极小。还有一个关键 APIasyncio.create_task()。它可以把一个协程包装成 Task 并立刻调度这样多个协程就可以并发运行import asyncio async def fetch_page(page_num): print(f开始抓取第 {page_num} 页) await asyncio.sleep(1) print(f第 {page_num} 页抓取完成) return f内容-{page_num} async def main(): # 并发创建 3 个任务而不是逐个 await tasks [asyncio.create_task(fetch_page(i)) for i in range(1, 4)] results await asyncio.gather(*tasks) print(results) asyncio.run(main())运行结果是这样的三个开始抓取会立刻连续打出说明三个协程已经并发执行了而不是等第一个 sleep 完再创建下一个。asyncio.gather()会在所有任务完成后返回结果列表顺序和传入的任务顺序一致。3.3 可等待对象的三种类型与选择规则Python 官方把可等待对象分为三类协程coroutine、任务Task、Future。很多人写代码时纠结什么时候该用哪种我的建议很简单如果你在写一个小函数内部有一个或多个 I/O 等待点那直接定义async def协程就够了调用时用await。如果你要并发运行多个协程让它们同时在事件循环里跑那用asyncio.create_task()创建任务再配合gather()或asyncio.wait()来聚合。如果你需要等待某个底层事件或回调完成的信号或者自己在实现一个异步库那才需要直接操作 Future。附上我常用的一张决策流程先判断是不是 I/O 密集任务不是就不该用 asyncio是的话再看协程之间是并行关系还是串行关系串行就直接一步步 await并行就 create_task 后 gather。4. asyncio 实战从并发爬虫到超时控制4.1 实战一用 asyncio aiohttp 实现并发网页抓取光说不练假把式。我拿一个真实场景来说用户在社区里问“Python 并发爬虫怎么做”最典型的实现就是 asyncio aiohttp 组合。下面这段代码是我在写一个资讯聚合爬虫时用的简化版结构清晰适合作为参考模板。import asyncio import aiohttp # 测试用的目标地址 URLS [ https://httpbin.org/delay/2, https://httpbin.org/delay/2, https://httpbin.org/delay/2, https://httpbin.org/delay/2, https://httpbin.org/delay/2, ] async def fetch_one(session, url, index): try: async with session.get(url, timeout10) as response: # 模拟解析响应内容 text await response.text() return index, len(text), response.status except asyncio.TimeoutError: return index, -1, 408 except Exception as e: return index, -1, 500 async def main(): # 创建共享的会话复用连接池性能更高 async with aiohttp.ClientSession() as session: tasks [asyncio.create_task(fetch_one(session, url, i)) for i, url in enumerate(URLS)] results await asyncio.gather(*tasks) for idx, length, status in results: print(f任务 {idx}: 状态码{status}, 响应长度{length}) if __name__ __main__: asyncio.run(main())几个关键点在代码里其实已经有体现了但我还是想单独强调一下aiohttp.ClientSession必须复用。每次请求都新建 session 的话等于放弃了连接池复用性能和资源占用都很差。官方文档也明确建议一个应用生命周期内只创建一个 session。asyncio.gather(*tasks)内部会等待所有任务完成默认情况下如果一个任务抛异常它并不会立刻取消其他任务而是所有任务跑完后统一把它们的结果聚合起来。如果希望“只要有任意一个任务失败就立即返回”可以传return_exceptionsTrue这样异常会作为结果返回而不是抛出。并发数不是越大越好。本地测试时我试过一次性创建 5000 个任务结果目标服务器直接拒绝服务还把自己的文件描述符给耗尽了。后面我会讲如何用信号量限制并发。这段代码跑起来后五个请求的耗时大约等于一个请求的耗时因为每个请求 delay 2 秒但它们是并发发出去的。如果你把 URLS 换成真实的几百个待抓取页面效果会更明显。4.2 实战二用 Semaphore 信号量控制并发数上限在实际爬虫项目中盲目的高并发会面临几个问题目标网站的反爬策略、本地 socket 连接数上限、内存占用。我踩过的真实坑是用 asyncio 一次性对 3000 个 URL 全部创建任务结果本机的文件描述符不够用进程直接报OSError: [Errno 24] Too many open files。从那以后我给所有并发任务都加了信号量控制。asyncio.Semaphore就是一个并发闸门同一时间只允许 N 个协程进入临界区import asyncio import aiohttp URLS [fhttps://httpbin.org/delay/1 for _ in range(100)] CONCURRENCY 20 # 最大并发数 semaphore asyncio.Semaphore(CONCURRENCY) async def fetch_with_limit(session, url, index): async with semaphore: # 获取信号量超过阈值则等待 async with session.get(url, timeout10) as response: return index, response.status async def main(): async with aiohttp.ClientSession() as session: tasks [asyncio.create_task(fetch_with_limit(session, url, i)) for i, url in enumerate(URLS)] results await asyncio.gather(*tasks) print(f完成 {len(results)} 个请求最大并发限制为 {CONCURRENCY}) asyncio.run(main())注意信号量要定义在协程外面否则每个协程都拿自己的信号量等于没限制。async with semaphore是推荐写法它保证在进入临界区时获取信号量、离开临界区时自动释放即使中途抛异常也不会卡死信号量。并发数设多少合适这个没有标准答案取决于目标服务器的承载能力和你的网络环境。我一般从 10 起步逐级加大观察响应时间和错误率。如果发现大量 timeout 或者 429请求过多就把并发数降下来。4.3 实战三超时控制与任务取消的坑在 Python 3.8 之后给单个协程加超时最优雅的方式是asyncio.timeout()它是个异步上下文管理器用法非常直观import asyncio async def slow_operation(): await asyncio.sleep(10) return 完成 async def main(): try: async with asyncio.timeout(3): result await slow_operation() print(result) except TimeoutError: print(操作超时任务已被取消) asyncio.run(main())这里有一个很重要的机制asyncio.timeout(3)内部在超时之后会尝试取消当前任务而Task被取消时会在协程内抛出asyncio.CancelledError。如果你的协程内部有try/finally块要小心取消时的清理逻辑——有些资源清理代码可能因为取消异常而跳过或者抛出二次异常。我遇到过的最隐蔽的问题是在 finally 里做清理操作结果清理操作本身也是一个 await这时候任务已经被取消了再次 await 会直接抛CancelledError。解决方案是给清理操作加上asyncio.shield()保护或者把清理逻辑放到外面单独处理。4.4 常见并发编排模式gather、wait 与 TaskGroup除了gather还有几个任务编排 API 需要掌握它们在特定场景下比 gather 更顺手asyncio.wait()可以指定返回时机比如FIRST_COMPLETED表示只要有任意一个任务完成就返回适合做“多路竞争”场景比如同时请求多个数据源谁先返回就用谁的。asyncio.as_completed()返回一个迭代器每迭代一次就得到一个已完成的任务结果适合边下载边处理的情况。asyncio.TaskGroupPython 3.11 新增比 gather 更安全。TaskGroup 的特点是组内任何一个任务异常会自动取消组内所有仍在运行的任务避免任务悬挂。我用 TaskGroup 重写爬虫任务就是这个效果import asyncio async def worker(name, delay): await asyncio.sleep(delay) print(f{name} 完成) return name async def main(): # 使用 TaskGroup 管理一组任务 async with asyncio.TaskGroup() as tg: tg.create_task(worker(A, 2)) tg.create_task(worker(B, 1)) tg.create_task(worker(C, 3)) print(所有任务结束) asyncio.run(main())TaskGroup 的异常传播行为是如果其中一个任务抛出异常TaskGroup 会先取消其他所有任务然后统一抛出异常——这在很多场景下比 gather 更符合预期因为它能防止“你已经知道任务失败了一半但其他任务还在傻傻地跑”。4.5 实战四异步代码中的阻塞陷阱这是我今天最想叮嘱读者的一个点。asyncio 的并发能力建立在“所有任务都会主动让出控制权”的前提下。如果有一个任务从头到尾不碰到任何await它会一直霸占事件循环其他所有任务都会被饿死。最常见的例子是混用阻塞库。比如在协程里调requests.get()而不是aiohttp的异步请求。requests是纯同步库它的网络调用会阻塞整个线程事件循环被卡住其他所有并发任务全部暂停。这种 bug 隐蔽性极高——在并发量小、单次请求时间短时几乎无感一旦某个请求慢下来整个服务的响应时间会瞬间恶化。除了网络请求还有一些容易忽略的阻塞点文件读写普通open()是阻塞的、加密解密大文件、time.sleep()。在 async 代码里如果你想模拟等待应该用await asyncio.sleep()而不是time.sleep()因为time.sleep()是同步阻塞。如果确实有阻塞操作无法避免有两个思路一是用asyncio.to_thread()把它丢到线程池里跑不阻塞事件循环二是用loop.run_in_executor效果类似但更底层。Python 3.9 后更推荐asyncio.to_thread()因为它简洁。import asyncio import time async def main(): print(开始) # 不要在协程里直接调用 time.sleep() # time.sleep(1) # 这会阻塞整个事件循环 # 正确做法 1用 asyncio.sleep await asyncio.sleep(1) # 正确做法 2把阻塞调用丢到线程池 await asyncio.to_thread(time.sleep, 1) print(结束)这个区别我实际踩过坑。当时我在一个 WebSocket 推送服务里用time.sleep做了个节流操作结果服务处理消息的吞吐量从每秒几千条直接掉到几十条而且所有在线连接都卡顿。改成asyncio.sleep之后立刻恢复正常。5. 事件循环的生命周期与底层调度原理5.1 asyncio.run 到底做了什么很多初学者会忽略asyncio.run()内部的一些实现细节。这个函数并不是简简单单地“运行一个协程”它做了几件重要的事如果没有正在运行的事件循环就创建一个新的事件循环将传入的协程包装成 Task 并运行运行期间处理所有已注册的回调、I/O 事件、子进程等协程结束后关闭事件循环并取消所有遗留的 Task。我特别提醒一点asyncio.run()每次调用都会创建一个全新的事件循环所以如果你在一个事件循环里多次调用asyncio.run()会得到一个 RuntimeError除非在单独的线程中。正确做法是只调用一次asyncio.run(main())把整个异步生命周期都放进main()里管理。有时候你在 Jupyter Notebook 或者交互式环境里使用 asyncio可能会遇到事件循环已存在的问题。这时可以手动获取当前事件循环或者用nest_asyncio之类的补丁但这些都是非常规操作不建议在正式项目里用。5.2 回调、Handle 与 ready 队列事件循环内部维护了一个就绪队列ready queue和一组回调句柄Handle。当一个 Task 执行到await时它会把一个续体continuation注册到某个 I/O 事件或者定时器上然后把自己从就绪队列中摘除。I/O 事件到达时事件循环会把对应的回调重新投递到就绪队列等待下一轮调度。这个机制也解释了为什么 asyncio 是协作式调度而不是抢占式调度。线程是由操作系统抢占式调度的——时间片用完就被踢下来而协程是自己主动让出 CPU 的。所以协程代码里不能出现长时间不 await 的 CPU 密集循环否则会让其他协程饿死。从实测来看Python 3.11 之后asyncio 的事件循环效率比 3.8 时代有了明显提升尤其是在任务创建和调度的开销上。如果生产环境允许我建议尽量用新版本 Python不只是为了语法特性性能和稳定性都有实打实的改善。6. 常见问题与排查技巧实录6.1 RuntimeError: Event loop is already running这是新手遇到最多的问题之一往往出现在 Jupyter Notebook、GUI 程序、或者嵌套调用asyncio.run()的情况。原因很简单事件循环是单实例的一个线程内同一时刻只能有一个事件循环在运行。如果你在一个协程里又调用asyncio.run()就会触发这个错误。解决方法也很直接不要在协程里再调用asyncio.run()。把入口统一或者用await嵌套协程。如果确实需要在已经运行的循环里拿到新任务的返回值用asyncio.create_task()await。6.2 Task was destroyed but it is pending这个警告出现时说明某个 Task 在事件循环关闭时还没有完成。常见原因在main()里用create_task()创建了任务但没有等它完成就直接退出了。解决思路确保你创建的所有任务都在退出前完成或者明确地取消它们。推荐用asyncio.gather()或者 TaskGroup 来统一管理任务生命周期避免漏掉个别任务。6.3 异步代码比同步代码还慢这种情况我也遇到过。如果你的“并发”任务全是轻量级操作单次不到 1 毫秒那 asyncio 的调度开销反而会成为负担。加上 Python 协程本身比普通函数调用稍重一些所以极短任务用 asyncio 未必划算。但如果任务是网络请求这类动辄几十毫秒的操作asyncio 的优势会极为明显。关键在于判断任务里有多少时间花在等待上——等待的时间越长asyncio 的价值越大。6.4 排查阻塞的神器running 线程与事件循环 debug 模式如果怀疑代码里有隐藏的阻塞调用可以用两个工具第一个是asyncio.run(main(), debugTrue)开启 debug 模式后事件循环会记录每个回调的执行耗时如果某个回调超过 100 毫秒会打印警告能帮你快速定位阻塞点在哪里。第二个是faulthandler标准库配合PYTHONFAULTHANDLER1环境变量可以让程序在卡住时打印当前线程栈直接看到是哪一行代码在阻塞。6.5 文件描述符耗尽的排查爬虫高并发时最常见的错误是OSError: Too many open files。检查方法是ulimit -n看当前限制用lsof -p 进程号 | wc -l看实际打开的文件数。解决方法一个是调整系统 ulimit更根本的是用信号量限制并发从源头控制连接数量。7. 我的几条实战心得最后分享几个我这些年用下来觉得最有价值的经验。第一不要把 asyncio 当作“更快的 Python”。它不会让单个任务变快它只是让“等待”变得便宜让程序能同时等待很多事情。评估一个系统能不能受益于 asyncio最简单的方法是看它有多少时间花在 I/O 等待上。如果超过 30%asyncio 值得尝试如果 80% 以上asyncio 基本就是最佳选择。第二多线程、多进程和 asyncio 不是互斥的它们可以组合使用。比如一个图片处理服务用多进程利用多核 CPU 做压缩用 asyncio 做网络请求的并发调度用多线程跑一些无法改造的旧同步库。组合起来往往比只用任何一种方案都强。第三写 asyncio 代码时我通常会给自己立几条硬规矩协程函数名字带上async_前缀或者统一放在services目录里所有外部 I/O 调用必须走异步库所有时间等待必须用await asyncio.sleep()所有并发任务必须有信号量限制所有任务都要有超时兜底。这几条规矩帮我避免掉了大量线上事故。asyncio 的生态现在已经很成熟了像 httpx、aiohttp、FastAPI、Tortoise ORM 都有完善的异步支持从同步代码迁到异步代码并没有想象中那么难。如果你主要做数据采集、接口开发、消息推送这类工作花两周时间把 asyncio 用熟后面省下的时间会远远超过这个投入。