尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Python asyncio实战:事件循环与协程并发技巧拆解

Python asyncio实战:事件循环与协程并发技巧拆解 1. 内容整体设计与思路拆解为什么你总觉得 asyncio 难学?说真的我见过太多人一上来就抱着asyncio官方文档啃啃了三天还是只会async def加await一遇到真实业务就卡壳。为什么因为异步编程不是语法问题是思维问题。你得从“一个线程执行到底”的同步思维切换到“事情挂起、事情恢复”的事件循环思维。这篇文章我不想重复文档里那些你已经看腻了的例子我打算用实际项目里真正会用到的场景拆解 10 个关键技巧把asyncio从“看得懂”变成“用得上”。在开始之前我默认你已经了解 Python 基础语法知道asyncio是 Python 官方的异步 IO 库核心是事件循环event loop、协程coroutine、任务Task和 Future。如果你还不太清楚这些概念也没关系读完第 2 章你会有直观理解。为什么要用asyncio最直接的原因是 IO 密集型场景。比如你要抓 1000 个网页、批量查询数据库、并发调用第三方 API同步写法要等每一个请求完全结束才能发下一个1000 个请求就是 1000 个网络往返的等待时间。用asyncio可以在等待网络响应的间隙去处理其他请求多个请求交错执行总耗时大幅缩短。它比多线程更轻量因为线程切换由操作系统调度开销大而协程切换由事件循环在用户态完成开销极小。这套设计思路的核心就是要把“阻塞”变成“挂起”。普通函数一旦调用了time.sleep()整个线程就卡住了什么也干不了而在协程里调用await asyncio.sleep()事件循环会把这个协程挂起转去运行其他就绪的协程等时间到了再回来继续。这就是异步编程的精髓不是快而是不等待。明白了这一点后面所有技巧都能顺着这条主线理解。2. 核心细节解析与实操要点事件循环、协程与任务到底怎么配合2.1 事件循环究竟是什么它和普通函数调用有什么本质区别很多人把事件循环想得很玄其实它就是一条消息队列轮询机制。你往事件循环里丢了一批待执行的东西协程、回调、IO 事件它不停地在队列里检查哪些已经有结果了哪些还在等待中。等待中的就去挂起有结果的就去恢复执行。整个过程是单线程的所以没有多线程的锁、竞态、线程切换这些头疼问题。理解这个机制你就能明白为什么asyncio.run()几乎是每个程序的入口。asyncio.run()会创建一个新的事件循环运行传入的协程等协程结束后关闭事件循环。在 3.10 版本之后它还会执行一次loop.shutdown_asyncgens()清理工作。你只需要记住任何异步程序入口一定是一个会被 await 的协程然后用asyncio.run(coroutine)启动。import asyncio async def main(): print(开始执行) await asyncio.sleep(1) print(结束执行) # 3.7 推荐的启动方式不要再用 loop asyncio.get_event_loop() asyncio.run(main())这里有一个几乎所有人都踩过的坑不要在事件循环里再创建另一个事件循环。比如你把asyncio.run()写进了一个已经被事件循环驱动的协程内部Python 直接抛RuntimeError。很多开发者第一次碰到这个错就懵了其实原因很简单——事件循环是单例的每个线程同一时刻只能有一个正在运行的事件循环。2.2 async def 和 await 的正确理解不是“并行”是“让出控制权”我先纠正一个非常大的误区asyncio不是把代码变成并行执行而是把等待的时间利用起来。比如你调用一个第三方 API正常情况下要等 2 秒才有响应同步代码里这 2 秒就白白浪费了异步代码里await会把 CPU 交还给事件循环事件循环立刻调度另一个协程去处理它的逻辑。所以它的本质是“并发”不是“并行”。await只能在async def声明的协程函数里使用这一点新手经常搞混。比如你写了一个普通函数想在里边await something()Python 直接语法报错。解决办法是把这个函数改成async def。但要注意一旦改成异步函数调用方式就变了——原来func()直接执行现在必须await func()或者在事件循环里执行。async def fetch_data(): return {status: 200, data: ok} async def main(): result await fetch_data() # 必须用 await 调用 print(result) asyncio.run(main())这里我有个经验不要为了“显得专业”把所有函数都改成 async这会带来巨大的心智负担。异步函数有传染性调它的地方也得是异步的。如果某个函数只是做纯计算没有 IO 等待维持普通函数就好用def声明然后在协程里直接调用它。异步技术是用来解决等待问题的不是用来炫技的。2.3 Task 和 Future协程不会被主动执行Task 才是调度单位这是我见过最多人忽略的一个细节。你写了一个async def然后调用它实际上协程对象被创建了但并没有开始运行。协程必须被包装成 Task或者直接被await事件循环才会去驱动它。import asyncio async def foo(): print(foo 开始) await asyncio.sleep(1) print(foo 结束) async def main(): print(main 开始) task asyncio.create_task(foo()) # 把协程包装成 Task立刻调度 print(main 继续执行不会被 foo 阻塞) await task # 等待 task 完成 print(main 结束) asyncio.run(main())注意asyncio.create_task()只能在事件循环运行期间调用因为它需要获取当前正在运行的事件循环。在 3.11 之前的版本也有人用asyncio.ensure_future()这个函数兼容性更好既接受协程也接受 Future但日常用create_task足够了。Task 的核心价值在于把协程变成可以独立调度、独立取消、独立查询状态的单元。你可以在一个协程里创建多个 Task让它们交错执行然后统一等待。掌握了这一点你才算真正迈进了异步编程的大门。3. 实操过程与核心环节实现从并发爬虫到超时控制一步步写出来3.1 第一个实战用 asyncio 并发抓取多个网页理论讲完了直接上真实场景。假设你要批量抓取 10 个 URL 的内容同步写法的耗时是所有请求耗时的总和而用 asyncio总耗时约等于最慢的那个请求。import asyncio import aiohttp URLS [ https://httpbin.org/delay/1, https://httpbin.org/delay/2, https://httpbin.org/delay/3, ] async def fetch(session, url): async with session.get(url) as resp: print(f请求 {url} 状态码: {resp.status}) return await resp.text() 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, return_exceptionsTrue) print(f抓取完成共 {len(results)} 个结果) asyncio.run(main())这里有几个关键点。第一使用aiohttp.ClientSession时建议用async with管理生命周期否则连接池可能没有正确释放。第二asyncio.gather()如果其中一个任务抛出异常默认会立刻向外部抛出导致其他未完成的任务被取消。如果你希望某个请求失败不影响其他请求一定要加return_exceptionsTrue这样任务失败的结果会被放在返回列表里而不是直接抛异常。3.2 你真的需要把并发数限制住信号量 Semaphore 的妙用并发请求不是越多越好。你一口气创建 1000 个 Task可能直接把目标网站连死也可能触发对方的反爬机制更可能耗尽你自己机器的文件描述符。这时候就要用asyncio.Semaphore来做限流。我自己的经验是抓取普通网站时并发限制在 10~20 比较稳妥访问 API 接口时先看对方文档的 Rate Limit再据此设置信号量。import asyncio import aiohttp LIMIT 5 semaphore asyncio.Semaphore(LIMIT) async def fetch_with_limit(session, url): async with semaphore: async with session.get(url) as resp: print(f并发受限当前请求: {url}) return await resp.text() async def main(): urls [fhttps://httpbin.org/delay/1?i{i} for i in range(20)] async with aiohttp.ClientSession() as session: tasks [fetch_with_limit(session, url) for url in urls] await asyncio.gather(*tasks, return_exceptionsTrue) asyncio.run(main())信号量的工作原理类似于商场门口的限流栏杆一次只能放 5 个人进去。async with semaphore会等当前正在执行的协程数小于 LIMIT 时才继续否则挂起当前协程等别人释放信号量。这是一个极其实用、但被大多数教程忽略的技巧。3.3 给所有协程加上超时控制防止卡死异步程序最怕遇到永远不返回的请求。如果你有一个协程因为网络问题一直挂着整个程序就卡在那里哪怕其他任务早就完成了。asyncio提供了多种超时控制方式。从 Python 3.11 开始官方强烈推荐使用asyncio.timeout()上下文管理器语法更加直观。import asyncio async def slow_operation(): await asyncio.sleep(10) return done async def main(): try: async with asyncio.timeout(3): result await slow_operation() print(result) except TimeoutError: print(操作超时已取消 slow_operation) asyncio.run(main())注意asyncio.timeout()抛出的是内置的TimeoutError它会在退出上下文时自动取消被包裹的协程。如果你用的是 Python 3.10 及以下版本可以用asyncio.wait_for(coro, timeout)效果类似但如果超时发生它会取消传入的协程。# Python 3.10 及以下的写法 async def main(): try: result await asyncio.wait_for(slow_operation(), timeout3) print(result) except asyncio.TimeoutError: print(超时了)我强烈建议你在每个有网络请求的协程上都套一层超时。实际生产环境中一个没有超时控制的爬虫早晚会遇到一次长时间无响应的状况届时整条任务链都会被拖垮。3.4 任务取消机制 Task.cancel 和 asyncio.shield你真的会正确使用吗除了超时取消也是异步编程里非常重要的工具。比如用户手动点击了“停止下载”你需要把正在跑的任务取消掉。import asyncio async def worker(): try: while True: print(工作中...) await asyncio.sleep(0.5) except asyncio.CancelledError: print(收到取消通知清理资源) raise # 重要CancelledError 必须继续抛出 async def main(): task asyncio.create_task(worker()) await asyncio.sleep(2) task.cancel() try: await task except asyncio.CancelledError: print(任务已取消) asyncio.run(main())这里有一个非常关键的细节协程内部捕获到asyncio.CancelledError后如果你想在取消前做清理工作比如关闭文件、释放连接可以处理但处理完必须继续raise把取消状态传给上层。如果你吞掉了这个异常任务不会正常取消行为会变得非常诡异。另外还有一种情况你不想让某个关键操作被外部取消可以用asyncio.shield()把协程保护起来外层取消时它不受影响但被保护的协程自己仍然可以正常被取消。3.5 混用同步代码的大坑time.sleep 会让整个事件循环卡住很多人写异步程序某处需要等待一下就顺手写了个time.sleep(1)。这是灾难性的错误。time.sleep是同步阻塞函数它会让当前线程整个睡过去事件循环也被迫暂停所有由它调度的协程全部停摆。正确做法是使用await asyncio.sleep()。如果协程里非要调用一个同步阻塞的库函数比如requests.get()注意 requests 是同步库你需要在调用前把它放到线程池里执行避免阻塞事件循环。import asyncio import requests def sync_fetch(url): # 这是个阻塞的同步函数 resp requests.get(url, timeout5) return resp.status_code async def main(): # 用 asyncio.to_thread 把同步函数丢进线程池避免阻塞事件循环 status await asyncio.to_thread(sync_fetch, https://httpbin.org/get) print(status) asyncio.run(main())Python 3.9 提供了asyncio.to_thread底层原理是把调用交给loop.run_in_executor()。如果你的运行环境是 3.8 及以下可以直接用loop.run_in_executor(None, func, *args)。这里送大家一个心法只要你发现事件循环卡住了先检查代码里有没有同步阻塞调用尤其是有没有time.sleep、requests.get、pymysql查询这类操作。它们每出现一次就会带走你一个事件循环里的协程。4. 常见问题与排查技巧实录我踩过的坑希望你不用再踩4.1 错误一多个 Event Loop 引发的 “Event loop is closed”这是一个非常经典的报错。很多人在 REPL 环境里运行asyncio.run()多次第二次就会报这个错或者在 Jupyter Notebook 里频繁运行异步代码也会遇到。原因是asyncio.run()每次都会创建新的事件循环并关闭它而在某些环境比如 Jupyter 的交互式环境里默认已经有一个事件循环在运行你再调用asyncio.run()就会冲突。解决办法有两种。第一已经有一个运行中的事件循环时使用await直接挂载协程第二调试代码时尽量避免反复运行asyncio.run()用nest_asyncio库可以解决大多数环境问题但生产环境建议还是把入口逻辑梳理清楚确保只有一个asyncio.run()。4.2 错误二忘记 await 导致 “coroutine was never awaited”这个报错说明你创建了协程对象却忘了使用await或包装成 Task它就永远不会被执行。Python 在 GC垃圾回收时发现这个协程没有被 await会打印这条警告。排查思路很简单找到报错提示里产生的协程位置检查是否漏写了await。我平时会开启 Python 的-W error模式让这类警告直接变成异常倒逼自己在开发阶段就把问题修掉。python -W error -m your_script.py4.3 错误三CPU 密集任务用了 async反而更慢了前面说了asyncio 是 IO 密集场景的利器但它对 CPU 密集场景基本没有帮助甚至更慢。原因是事件循环是单线程的一个 CPU 密集的协程一旦开始计算就会独占 CPU其他协程根本没有机会被调度等效于全都排队等待。真正的解决方案是把 CPU 密集任务丢到进程池里去执行用concurrent.futures.ProcessPoolExecutor。异步外部套一个await loop.run_in_executor(executor, func)来接收结果。记住IO 用异步CPU 密集用多进程两者混用时要特别注意数据传递的开销。4.4 调试技巧开启 asyncio 的调试模式Python 提供了内建的异步调试模式能帮你发现哪些协程运行时间过长、哪些 Future 没有被正确等待。设置环境变量PYTHONASYNCIODEBUG1 python your_script.py或者在代码里手动设置import asyncio import os os.environ[PYTHONASYNCIODEBUG] 1 async def main(): await asyncio.sleep(1) asyncio.run(main())在这个模式下Python 会显示所有仍在运行但未完成的 Task方便你定位泄漏的协程。就我个人的使用体验来说大多数“程序卡住不退出”的问题都能靠这个调试模式抓到元凶。4.5 领域经验Windows 环境下 asyncio 事件循环的特殊性如果你在Windows上开发可能会发现某些异步代码行为和Linux上不一样。原因在于 Windows 默认的 SelectorEventLoop 对某些 socket 支持有限信号处理也不一样。Python 3.8 起Windows 上默认使用 ProactorEventLoop它支持 subprocess 和 socket 异步操作更友好。如果遇到奇怪的 socket 问题可以显式设置import asyncio import sys if sys.platform win32: asyncio.set_event_loop_policy(asyncio.WindowsProactorEventLoopPolicy())生产部署如果目标是 Linux那么本地开发和线上环境的事件循环策略不一致可能引发“本地跑得好好的上线就出错”的经典问题。遇到这类问题先排查是不是事件循环差异导致的。我自己的建议是如果团队以 Linux 生产环境为主开发也尽量统一用 Linux 或 Docker 容器不然很多异步库的边界行为会反复折磨你。4.6 独门经验用 asyncio.Queue 搭建生产者-消费者模型真实项目里仅仅是并发抓取还不够你通常还需要一个任务队列来解耦生产者和消费者。asyncio.Queue是线程安全的异步版本用法和标准库queue.Queue很像但put和get都是异步方法。import asyncio import random async def producer(queue): for i in range(10): await queue.put(f任务 {i}) print(f已生产: 任务 {i}) await asyncio.sleep(random.uniform(0.1, 0.5)) async def consumer(queue): while True: item await queue.get() print(f已消费: {item}) queue.task_done() async def main(): queue asyncio.Queue(maxsize5) producers [asyncio.create_task(producer(queue))] consumers [asyncio.create_task(consumer(queue)) for _ in range(2)] await asyncio.gather(*producers) await queue.join() # 等待所有队列中的任务被处理完 for c in consumers: c.cancel() # 消费者需要手动取消否则永远不会退出 asyncio.run(main())这段代码里有几个细节值得注意。queue.join()会阻塞直到队列中所有项都被task_done()标记完成这是一个优雅的停止时机。消费者的while True循环不会自己结束所以需要在所有生产完成且队列清空后手动取消消费者协程。这是非常典型的异步生产消费处理模型无论你是做爬虫调度、消息消费、还是任务分发这套代码都能直接改改就用。5. 进阶技巧矩阵10 个关键点一次全部带走5.1 技巧总览一张表看清 10 个关键技巧为了让你以后翻阅方便我把这篇里提到的所有关键技巧整理成了速查表。这些不是教科书里的基础知识点而是我在实际开发中被坑过、验证过、沉淀下来的实战经验。序号关键技巧核心用途常见误区1asyncio.run 统一入口创建并运行事件循环在协程内部再次调用 asyncio.run2create_task 包装并发把协程变为可调度 Task创建后不 await 也不保存引用3gather return_exceptions批量执行并容忍单个失败忽略异常导致任务全部取消4Semaphore 限流控制并发请求数量不加限制导致资源耗尽5timeout 超时控制防止协程永久挂起忘记超时程序卡死6CancelledError 处理安全取消任务并清理资源吞掉取消异常7shield 保护关键操作避免关键协程被外部取消滥用导致任务无法取消8to_thread 规避阻塞同步阻塞函数放进线程池直接在协程里调用阻塞函数9Queue 生产消费模型解耦异步任务的生产和消费忘记 task_done 导致 join 死等10调试模式定位协程泄漏和未完成任务不知道有 PYTHONASYNCIODEBUG5.2 如何选择 gather、wait 还是 wait_for别再东猜西猜了除了表里的内容还有一个非常高频的选择困惑asyncio.gather()、asyncio.wait()、asyncio.wait_for()三者的区别。我给你一个最直白的总结。asyncio.gather()是你日常最常用的批量执行入口它把多个 awaitable 打包成一个 awaitable返回每个协程的返回值列表支持return_exceptions。如果你只关心“把一批任务跑完拿到结果”选它。asyncio.wait()更底层一点它接收 Task 集合或 Future 集合返回两个集合已完成的和未完成的。适合做“只要有一半任务完成就继续”的复杂编排但无法直接拿返回值需要自己去 Future 里取。asyncio.wait_for()是给单个协程加超时上限的包装器如果你只控制单个任务的时间用它。# 快速区分示例 import asyncio async def task(name, delay): await asyncio.sleep(delay) return f{name} done async def main(): # gather拿结果列表 results await asyncio.gather(task(A, 1), task(B, 2)) print(results) # wait拿 done / pending 集合 tasks {asyncio.ensure_future(task(C, 1)), asyncio.ensure_future(task(D, 2))} done, pending await asyncio.wait(tasks, timeout1.5) print(f已完成: {len(done)}, 未完成: {len(pending)}) # wait_for单任务超时 try: result await asyncio.wait_for(task(E, 5), timeout1) except asyncio.TimeoutError: print(E 任务超时被取消) asyncio.run(main())我建议从简化角度出发80% 的场景都能被gathertimeout覆盖其他两个工具了解原理即可真正用到复杂编排时再深入不迟。5.3 异步上下文管理器、异步迭代器让你的类更专业最后送你两个提升代码高级感的技巧异步上下文管理器和异步迭代器。异步上下文管理器就是async with核心是__aenter__和__aexit__用来管理需要异步初始化和异步清理的资源比如数据库连接、HTTP 会话。异步迭代器就是async for核心是__aiter__和__anext__用来异步地按需生成数据比如从数据库游标里逐条取记录。import asyncio class AsyncResource: async def __aenter__(self): print(异步初始化资源) await asyncio.sleep(0.1) return self async def __aexit__(self, exc_type, exc, tb): print(异步清理资源) await asyncio.sleep(0.1) async def __aiter__(self): for i in range(5): await asyncio.sleep(0.2) yield i async def main(): async with AsyncResource() as res: async for item in res: print(f拿到数据: {item}) asyncio.run(main())从我个人经验来看把这些协议实现好项目代码的复用性和可读性都会上一个台阶特别是做库、做框架或者团队协作时异步上下文管理器几乎是标准姿势。6. 踩坑之后的真实感受异步编程真正难在哪里写了这么多代码和技巧最后聊一点我自己的感悟。为什么很多人觉得异步编程难不是难在语法而是难在调试和心智模型。异步代码的执行顺序不再像同步代码那样从上到下打印出来的日志可能是交错出现的异常堆栈也不够直观出了问题你很难说清楚到底是哪个协程在什么时刻做了什么。我的建议是不要一开始就追求花哨的并发编排。先用最简单的await串行跑通业务再逐步引入create_task和gather来提速。每引入一个机制就加一次日志确认行为符合预期后再往前走。我在实际项目里踩过无数次坑之后养成了一个习惯给每个耗时超过 200ms 的异步任务包一层带超时和日志的包装函数。这样一旦线上出问题日志能告诉我到底是哪个环节超时了是哪个协程泄漏了。这是我从大量 run book 里提炼出的最朴素也最有效的手段。如果让我只留一条经验给后来的新手我会说把asyncio当成一个事件循环驱动的任务调度系统而不是“魔法并行工具”。理解调度、理解阻塞、理解取消比背十个库函数都管用。等你真正习惯了这种“让出、恢复”的思维方式再看异步代码就不会觉得乱反而会觉得它有一种单线程协作美学的流畅感。
返回列表