ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

asyncio实战:事件循环与协程如何让并发请求性能飙升?

asyncio实战:事件循环与协程如何让并发请求性能飙升? 说实话asyncio 这套东西我刚接触那会儿看着网上的例子总觉得像天书——明明每个单词都认识拼在一起就不知道程序是怎么跑起来的。后来是写了一个完整的案例把事件循环、Task、await 这些概念全部揉进去才真正把脑子里的线理顺了。所以这一篇我不想再讲零散的概念直接用两个实实在在的入门案例带你把 asyncio 从“知道”变成“会用”。如果你已经被同步代码的 I/O 等待折磨过或者继承老项目时看到 async 关键字有点发怵那这篇就是给你准备的。1. asyncio 到底解决了什么问题1.1 同步代码慢在哪儿一个实测对比先看一个最常见的场景写一个脚本要检查 100 个网站的可用状态。用同步的 requests就是一个个来请求发出去之后整个过程卡在那里等服务器回应哪怕对方慢得离谱你也只能干等着。我曾经真跑过一次100 个 URL 里混着几个响应要 20 秒的最终总耗时直接飙到两三分钟。而用 asyncio 重写一遍之后同样的 100 个请求总耗时基本等于其中响应最慢的那一个。注意不是 100 个请求的时间相加而是最慢那一个的时间。百来个请求处理完经常不到 5 秒。这个差距在批量调第三方 API、爬虫抓页面、大量小文件下载这类场景里体验是质变。但要注意这里说的“慢”指的是 I/O 密集型任务也就是程序在等待网络、磁盘、数据库响应。如果你要做的是大量 CPU 计算比如图像处理、复杂加密、大数据排序那 asyncio 帮不了忙后面我会专门聊这个。1.2 线程、进程、协程为什么协程更适合 I/O 密集很多人第一反应是“有线程啊用线程不也能并发吗”能但代价不一样。线程的切换是由操作系统决定的而且每个线程都有自己独立的内存栈开多了内存占用上去了切换时 CPU 还要浪费时间去保存和恢复现场又是一笔开销。你可以把协程理解为“用户态的轻量级线程”——它不归操作系统管而是由程序自己在合适的时候让出控制权。asyncio 的事件循环就像一个调度员轮询每个任务你去等这个 I/O 结果eslint去检查下一个任务那个 I/O 有结果了再回来继续执行。因为切换是代码里显式的代价极小所以哪怕你创建成千上万个协程也不会像开线程那样吃力。还有一点Python 有 GIL多线程在纯 CPU 计算上基本是负优化但在 I/O 等待的场景里GIL 其实不是主要瓶颈可线程切换的成本依旧在。协程因为压根不占内核资源切换成本低得多这也是它在 Python 异步编程里被推崇的主要原因。1.3 先泼一盆冷水什么场景别用 asyncio我不希望你看完这篇文章把所有代码都重写一遍。asyncio 是有适用边界的。纯 CPU 密集型任务比如视频转码、模型推理、复杂计算。这类任务没有“等待”await完全派不上用场。整个项目里只有 1-2 个请求同步就够了引入 asyncio 只会增加心智负担。依赖的第三方库是纯同步实现你又没法换成异步版比如某些内部工具包那硬用 asyncio 反而会把问题搞复杂。最适合 asyncio 的场景基本可以概括为一句话很多个“等别人”的操作彼此之间互不依赖可以同时发出请求。这种场景asyncio 能把你从等待的泥潭里拉出来。2. 五个核心概念串起 asyncio 的地基2.1 事件循环是调度中心asyncio.run()这个函数大家应该都不陌生它做的事其实有两件创建事件循环跑协程然后关闭。事件循环就是整个异步程序的运行机制——所有协程都会注册到里面由它决定谁先谁后谁该等待谁可以恢复。用生活里的事情来类比事件循环就像一个餐厅的前台顾客点完餐发起 I/O 请求前台让顾客去休息区等挂起协程又去安排下一桌客人点餐执行下一个协程饭菜做好了前台喊顾客来取I/O 完成恢复协程。前台不用等谁把菜吃完再做事它只需要不停地在“谁可以点餐”“谁的菜好了”之间周旋。2.2 协程函数调用后并不执行这是新手最容易懵的地方。你写了一个async def fetch_data(): ...然后在代码里调用它fetch_data()你以为它执行了实际上它只是创建了一个“协程对象”里面的代码一行都没跑。我在实际教人的时候经常看到新人卡在这里最后终端打出一行RuntimeWarning: coroutine was never awaited这才发现自己忘了加await。也正因为这个设计协程的运行时机完全由你掌控。你想让它先注册后面再 await或者直接丢给 Task 让它自动跑都很灵活。永远记住一句话协程函数是一条菜谱协程对象才是一份正在被准备的菜品。2.3 await 只能等待“可等待对象”在 asyncio 的世界里await后面能跟的东西有三类协程对象、Task、Future。泛称就是“可等待对象”。这里我不打算光讲理论直接给一个能跑的例子。假设我们要模拟三个请求最快的那个先返回import asyncio async def request(url, delay): await asyncio.sleep(delay) return fresponse from {url} async def main(): # 三个协程同时发起 task1 asyncio.create_task(request(a.com, 1)) task2 asyncio.create_task(request(b.com, 2)) task3 asyncio.create_task(request(c.com, 0.5)) # 谁先完成先处理谁 for coro in asyncio.as_completed([task1, task2, task3]): result await coro print(完成:, result) if __name__ __main__: asyncio.run(main())运行结果你会看到 c.com 最先打印然后是 a.com最后是 b.com。asyncio.as_completed就是一个“按完成顺序取结果”的迭代器特别适合处理耗时差异大的批量请求。2.4 Task 是对协程的包装直接await coroutine()等它执行完下一个协程才会继续这其实就是同步执行了。但很多时候我们想要的是先把这些协程全部丢进后台让它们同时跑谁先完成谁先返回。那就需要asyncio.create_task()。Task 本质上就是一个被事件循环调度执行的协程包装对象。创建 Task 之后协程立刻开始运行其实是进入调度队列不需要你再显式地await它才会启动。你可以把它想象成把菜谱递给了厨师厨师立刻开始切菜而不需要你站在旁边盯着每一道工序。2.5 gather、wait、as_completed 的选择这三个函数经常放在一起比较功能有重叠但侧重点不同。gather适合知道要等哪些任务并且想把结果按原顺序收集起来的情况。它的返回值是“所有结果按传入顺序排列的列表”。wait更底层它支持等待超时时间、控制是都完成还是任一完成还返回“已完成”“未完成”两组集合。as_completed用得最少但在需要“每完成一个就立刻处理对应结果”的场景里非常爽快。对照表函数返回内容适合场景asyncio.gather按传入顺序返回所有结果一旦有异常整个 gather 抛出需要等全部完成且关心所有结果asyncio.wait返回 (done, pending) 两个集合需要手动处理超时和未完成任务asyncio.as_completed一个迭代器按完成顺序产出结果每完成一个就处理一个及时处理中间状态3. 完整案例把 100 个请求从 2 分钟压到 3 秒3.1 需求与基线测试我拿一个很常见的需求来讲监控一批 API 接口的健康状态。假设有 100 个接口我们要向它们发请求统计响应状态码和耗时。同步版本我会用requests写个最朴素的 for 循环地基就这样打。import time import requests urls [https://httpbin.org/delay/1] * 100 # 每个请求至少等 1 秒 start time.perf_counter() for url in urls: resp requests.get(url, timeout5) print(resp.status_code, resp.elapsed.total_seconds()) print(f同步耗时: {time.perf_counter() - start:.2f}s)这段代码什么技巧都没有但它是我们后面所有优化的基线。100 个请求如果每个都等 1 秒以上总耗时至少 100 秒实际情况里还会更慢。下面我们一步步把它变成异步版本。3.2 第一版标准库 asyncio.to_thread 快速体验如果你暂时不想引入第三方异步 HTTP 库还有一条捷径asyncio.to_thread。它可以把一个普通同步函数放到线程池里执行再包装成协程来 await。严格来说这不是“真异步”但在和旧代码拼接时非常实用。import asyncio import time import requests async def fetch(url): resp await asyncio.to_thread(requests.get, url, timeout5) return resp.status_code async def main(): urls [https://httpbin.org/delay/1] * 100 tasks [asyncio.create_task(fetch(url)) for url in urls] results await asyncio.gather(*tasks) print(results) if __name__ __main__: asyncio.run(main())这段代码的耗时已经能降到和最后面 aiohttp 版本差不多的量级因为阻塞操作交给线程池去扛主循环没有被阻塞。但它也有隐患如果 1000 个请求全部丢进线程池线程切换成本就会重新冒头而且你本质上还是在吃操作系统线程的资源。所以它更适合作为“改造第一步”而不是最终方案。3.3 第二版aiohttp 真异步并发真正能发挥 asyncio 全部能力的是aiohttp这个库的 API 和 requests 很像但所有 I/O 操作都是非阻塞的。先pip install aiohttp然后看下面这段完整代码import asyncio import time import aiohttp async def fetch(session, url): async with session.get(url, timeout5) as resp: status resp.status text await resp.text() return status, len(text) async def main(): urls [https://httpbin.org/delay/1] * 100 async with aiohttp.ClientSession() as session: tasks [asyncio.create_task(fetch(session, url)) for url in urls] results await asyncio.gather(*tasks) for status, size in results: print(status, size) print(faiohttp 并发耗时: {time.perf_counter() - time.time_start():.2f}s) if __name__ __main__: asyncio.run(main())注意我用了async with来管理连接池的会话而不是为每个请求单独建立连接这是性能差异非常大的一个细节。如果每次session.get都新建一个 ClientSession开销会翻好几倍。测试下来100 个请求在这个版本里大概 2-3 秒就能跑完这就是真异步的效果。不过这个例子还有个隐患一旦某个请求异常超时gather会直接把整个任务组炸掉。后面我们加上容错和处理。3.4 第三版限流、超时、重试全加上真实项目里不能这么裸奔至少要考虑三件事不要让对方服务器以为你在攻击所以要加并发限制单个请求要么超时熔断要么重试异常不能把整个任务组带崩。我直接给一个项目里可以直接拿来改的版本import asyncio import aiohttp urls [https://httpbin.org/delay/1] * 100 async def fetch_with_retry(session, url, semaphore, retries3): async with semaphore: # 控制同时进行的请求数 for attempt in range(retries): try: async with session.get(url, timeout5) as resp: if resp.status 400: return resp.status, await resp.text() else: raise aiohttp.ClientError(fhttp {resp.status}) except (aiohttp.ClientError, asyncio.TimeoutError) as exc: if attempt retries - 1: return None, str(exc) await asyncio.sleep(0.5 * (attempt 1)) # 退避 return None, failed async def main(): semaphore asyncio.Semaphore(20) # 最多 20 个并发 async with aiohttp.ClientSession() as session: tasks [asyncio.create_task(fetch_with_retry(session, url, semaphore)) for url in urls] results await asyncio.gather(*tasks) success sum(1 for status, _ in results if status) print(f成功 {success}/{len(urls)}) if __name__ __main__: asyncio.run(main())这套代码的核心是Semaphore(20)它就像一个并发闸门保证同时只有 20 个请求在飞。重试时用了指数退避的简化版本第一次失败等 0.5 秒第二次等 1 秒第三次彻底放弃。整体逻辑已经和真实生产代码非常接近了。4. 进阶三板斧并发限制、超时与取消4.1 Semaphore 限流的底层逻辑信号量这个东西说穿了就是一个计数器。你创建asyncio.Semaphore(20)它初始值是 20。每次进入async with semaphore时如果计数大于 0 就把计数减 1放你过去如果已经是 0你就得在门口排队等前面的人出来时把计数加回 1再放你进去。为什么不能直接发 1000 个并发的请求不是因为 asyncio 扛不住而是对端服务器扛不住。很多服务端对单 IP 并发连接数有限制可能 50 个并发就把你拒之门外了。从你自己的角度连接池的连接数也有上限ClientSession 默认连接数在 100 左右超过之后照样排队。所以在学 asyncio 的同时把 Semaphore 用熟是实际项目和“玩具代码”的分水岭之一。4.2 wait_for 超时的坑与正确姿势asyncio 的wait_for用起来很简单result await asyncio.wait_for(coro, timeout10)。10 秒内没完成就抛出asyncio.TimeoutError。听起来很美好但有个坑超时后任务不会自动“消失”它会被取消cancel并等待其清理完成。如果你的协程里有收尾工作比如释放连接、写日志取消信号什么时候到是由事件循环决定的不是立刻终止。还有就是不要在wait_for里传一个简单的协程同时又塞进任务列表不然它会变成两个独立的对象执行。我有一次写了一批请求用create_task把协程丢进列表然后又对同一个协程对象wait_for(..., timeout3)结果任务执行了两遍接口被重复调了。根源就在于协程对象只能被 await 一次你传进去的引用还是原来那个。4.3 任务取消与 shield 的取舍任务取消的机制很优雅但也很容易踩坑。task.cancel()会在协程内抛一个asyncio.CancelledError协程可以选择在这个异常处正常退出也可以用try/finally做清理。关键点是如果你捕获了CancelledError又不重新抛出事件循环会认为任务已经被取消了但实际它在后台“苟活”着容易造成逻辑错乱。shield()则相反它保护一个协程不被取消。默认行为是取消信号直接打到 shield 上底层协程继续跑。但注意如果底层协程自己不去理会取消那 shield 也会跟着完蛋。我见过不少人在超时场景里加了 shield以为万事大吉结果底层用的是同步阻塞库取消信号根本不生效。所以要在“真异步”的前提下讨论取消混入同步阻塞逻辑后这些都是空谈。4.4 三种结果收集方式对比前面提到过 gather、wait、as_completed这里我再多说一层。我在实际使用中的习惯是这样的如果任务之间没有依赖但要等所有结果优先gather因为它返回结果列表的顺序稳定方便和 URL 列表一一对应。如果某个任务失败不能影响整体那就给 gather 加return_exceptionsTrue让它把异常封装进返回列表而不是往上抛。如果第一个成功的结果就能继续推进流程比如搜索接口并发多个服务商谁快用谁那就用asyncio.wait配合FIRST_COMPLETED拿到第一个就取消其余。这里给一个FIRST_COMPLETED的简要示例import asyncio async def service_a(): await asyncio.sleep(3) return a async def service_b(): await asyncio.sleep(1) return b async def main(): task_a asyncio.create_task(service_a()) task_b asyncio.create_task(service_b()) done, pending await asyncio.wait( {task_a, task_b}, return_whenasyncio.FIRST_COMPLETED, ) for task in pending: task.cancel() # 赢了就不用等慢的那个了 print(done.pop().result()) if __name__ __main__: asyncio.run(main())这种竞速模式在微服务调用、多路数据源选路时特别好用。5. 真实项目里最常见的 5 个坑5.1 RuntimeErrorEvent loop is closed你很可能在 Jupyter Notebook 里遇到过这个报错。原因很简单Jupyter 自己维护了一个事件循环而你在同一个进程里调用了两次asyncio.run()第二次调用时旧事件循环已经被关闭了新事件循环又想用同一个东西直接冲突。解决办法也简单在 Jupyter 里不要用asyncio.run()用await配合await asyncio.create_task(...)或者直接在 notebook 的顶层执行asyncio.run前面的代码块。另一个常见场景是 Flask/Django 这种同步框架里直接用 asyncio往已有的事件循环里塞任务时也会遇到类似问题。这时候用asyncio.set_event_loop(asyncio.new_event_loop())换个新循环或者改用nest_asyncio但这是补丁式方案心里要有数。5.2 协程里混用同步库导致“假异步”这是最大的一个坑。你写了一堆async def看着井井有条但里面突然有一行requests.get()这行同步请求一旦发出去整个事件循环就被阻塞住了其他协程全在坑里排队。你写的异步代码瞬间变成一张华丽的外皮。我排查过很多类似的问题表现是程序“异步改造”后耗时几乎没降。第一反应不是怀疑并发逻辑而是找找有没有同步阻塞调用混在里面。解决方案有三个换 aiohttp/httpx 这样的异步客户端或者用asyncio.to_thread把同步函数丢进线程池再或者如果你是在 FastAPI 等框架中可以直接用run_in_executor把阻塞任务隔离到独立线程。5.3 Task exception was never retrieved这个报错的含义是某个 Task 里抛了异常但没有任何地方接收这个异常。事件循环检测到后会在垃圾回收时打出一条警告。示例import asyncio async def bad(): raise ValueError(boom) async def main(): task asyncio.create_task(bad()) # 没有 await也没有加异常处理 await asyncio.sleep(0.1) asyncio.run(main())运行后你就会看到Task exception was never retrieved的警告。解决方法是所有 Task 都尽量有归属要么收集进列表后用gather或者wait统一处理要么给 Task 添加add_done_callback去检查异常要么在每个协程内部就做好 try/except把异常吞掉并记录日志。总之不要让任何异常悬空。5.4 回调地狱与 async/await 之间的关系老牌异步框架 Twisted 和 asyncio 的早期风格都鼓励用回调请求完成后调这个函数失败后调那个函数。回调一嵌套代码就变成了“箭头形”阅读和排查都很痛苦。asyncio 里虽然也有loop.add_reader这种偏底层的回调机制但我们现在写业务代码我的建议是能用 await 就不要用回调。这两个风格怎么选哪怕是协程之间的链式调用也优先用await把数据流写清楚而不是把一个函数的返回值作为参数传给另一个回调。唯一推荐回调的场景是在事件循环的底层 API、或者在自定义事件循环集成时普通业务代码几乎用不到。5.5 排查利器debug 模式与慢回调日志写 asyncio 程序出 bug 时靠 print 大法效率太低了。建议学会打开 asyncio 的调试模式在代码开头设置loop.set_debug(True)或者在环境变量里设置PYTHONASYNCIODEBUG1。开启之后事件循环会记录每个回调的执行时长如果发现某个协程执行时间超过默认阈值通常是 100ms日志里会标出“slow callback”。排查思路一般是先开 debug 模式看是哪个协程耗时异常再沿着日志定位是 I/O 等待还是在做 CPU 计算然后看是不是有同步阻塞库混入。我曾通过这种办法发现一个第三方 SDK 内部用的是同步requests藏得特别深不开 debug 根本发现不了。6. 走出 Python协程思想的跨语言迁移6.1 Kotlin 协程和 Flow 是什么虽然标题是 Python 的 asyncio但“协程”这套心智模型不止一个语言在用。现在 Android 开发里很火的 Kotlin 协程就是典型suspend fun相当于 Python 里的async defwithContext用来切换线程CoroutineScope相当于事件循环的作用域。而 Flow 则是 Kotlin 做异步数据流的一套方案类似 Python 里的异步生成器。很多 Python 程序员切过去没什么障碍因为两者的本质是一样的用顺序写的方式表达异步流程把线程切换交给框架而不是程序员。6.2 从 asyncio 到 Flow它们想解决同一件事如果你理解 asyncio 里“await 一个 I/O 操作时让出控制权”的逻辑那你理解 Kotlin 协程也不会太难。Flow 在它的基础上多了“上游发射数据、下游收集数据”的模型收集过程中一旦遇到网络请求这种 I/O照样会挂起并在合适的时机恢复。这两种语言里的异步观念完全同源我甚至觉得多学一门语言的异步实现能反过来加深对另一门语言的理解。当然这篇的篇幅有限我的建议是先把 Python 这套玩熟再跨过去看 Flow你会发现很多概念可以直接平移。6.3 用同理心去学任何语言的异步最后说点题外话。异步编程的难点不在语法而在于心智模型的切换你不能再按“第一行执行完再执行第二行”的老思路去读代码而要时刻去想“这里挂起了吗挂起多久谁在等它”。一旦接受了这个设定你学任何语言都有了一座桥。反过来如果只是背 API不理解事件循环和任务调度到了下一个语言又会回到从零开始的状态。结尾在我实际带新人的过程中协程这一块最常出现的问题是“看得懂例子不敢写逻辑”。我的建议一直很简单找一个自己手头正在做的同步小脚本比如批量检查接口、批量下文件、批量发通知照着这篇的思路逐步改造成 asyncio 版本跑通了再继续往里加限流、超时、重试和取消。当你把一个真实脚本改造成功的那一刻对事件循环、Task、Semaphore 这些概念的体感会和看完十篇教程都不同。如果真要说一个最容易被人忽略的小技巧那就是写 asyncio 代码时每个协程内部都要有独立的异常处理不要把风险都抛给最外面的gather。我见过太多线上事故都是从“某个请求异常导致整批任务都退出”开始的你花半小时给所有 create_task 加上兜底日志未来能省下几个通宵。
返回列表