ARTICLE DETAIL

资讯详情

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

Python异步编程核心:asyncio、协程与任务调度实战

Python异步编程核心:asyncio、协程与任务调度实战 如果你写过几段带网络请求或文件读写的 Python 代码大概率体会过这种场景一个爬虫循环请求 50 个页面90% 的时间都耗在那句requests.get()上。你以为自己在写代码实际却是在等网络。Python 异步编程这套东西就是专治这种“等”的问题。它不靠多线程而是用单线程里的事件循环把空闲时间榨干一个请求发出去之后先挂起CPU 立刻跑去处理另一个任务等数据回来再切换回来。今天从头到尾梳理一遍asyncio、async/await、协程、任务调度的完整用法把爬虫、接口批量调用、数据抓取这些 IO 密集场景的实践细节也拆开讲透。适合刚接触协程、或者已经写过async但还是经常踩坑的同学。先说一个比较反直觉的事实异步不是让你的代码“跑得更快”而是让你的代码“不空等”。无论用哪个版本的 Python这个心智模型建立不起来后面写多少await都是白搭。1. 先说清楚异步解决的是“等待”不是“计算”1.1 同步代码为什么会卡在 IO 上大部分 Python 新手写出来的网络请求代码是这样的import requests urls [fhttps://httpbin.org/delay/{i} for i in range(1, 6)] for url in urls: resp requests.get(url, timeout10) print(resp.status_code)这个脚本的行为非常直观一个一个请求上一个返回了才发下一个。httpbin.org/delay/1会故意拖 1 秒返回delay/2拖 2 秒所以整个脚本跑完的时间大约是 1234515 秒。问题出在哪儿出在requests.get()这个函数上。它在等待服务器响应的这 15 秒里Python 解释器什么都没干CPU 处于空闲状态。而我们只是为了让“发请求→等响应”这件事变得有序牺牲掉了本可以并行处理的时机。这里要理解一个关键点阻塞blocking不等于慢而是等于占用。即使每次请求只需要占用几十毫秒的计算量但一旦进入阻塞态整个线程就停止调度了后面排队的任务全被堵住。1.2 事件循环如何“单线程干多活”异步编程的底层基石是事件循环event loop。你可以把事件循环想象成一个餐厅前台只有一个服务员在岗但特别会统筹安排不需要站在每桌等客人吃完饭点完单就让厨房去做服务员立刻去接待下一桌。对应到asyncio里线程只有一个但“当前正在等待 IO 的任务”会被挂起事件循环会去执行其他还没被挂起的任务。当某个任务的 IO 事件完成了操作系统会通知事件循环再把它加回调度队列。这种非阻塞的实现依赖操作系统底层的 IO 多路复用机制比如 Linux 上的epoll、macOS 上的kqueue。Python 的asyncio默认在 Unix 系统上使用epoll/kqueue在 Windows 上使用IOCPPython 3.8 之后ProactorEventLoop成为 Windows 默认。这些机制的本质都是一样的让一个线程同时监控多个文件描述符有数据可读或可写时才去处理对应回调。拿刚才的 5 个请求举例用异步方式写总耗时和其中最慢的那个请求持平大约 5 秒而不是 15 秒。这不是魔术这就是合理地利用了原本被白白浪费的 IO 等待时间。1.3 异步适用的场景边界但异步也不是万能药。这里需要区分两种密集类型IO 密集程序大部分时间在等待网络、磁盘、数据库连接CPU 闲得发慌。CPU 密集程序一直在做大量计算比如图像渲染、数值模拟、复杂加密。asyncio最擅长的是 IO 密集。网络爬虫、API 批量调用、文件批量处理、RPC 查询、量化交易里的行情订阅这些都天然是 IO 密集适合上异步。但对 CPU 密集任务异步不仅没有好处还可能有反效果。事件循环是单线程的你的计算任务被分成了很多小段在切片交换的时候还要额外付出上下文切换的开销。遇到这种场景应该用multiprocessing开多进程或者直接交给 C 扩展去算而不是硬上asyncio。把这个边界定清楚你才不会被网上各种“异步性能提升 XX 倍”的文章带偏。异步不解决计算问题它只解决等待问题。2. 核心语法拆解async def 与 await 到底是什么2.1 调用 async def 不会执行代码Python 从 3.5 引入async/await语法3.7 正式把asyncio.run()作为标准入口。此后协程的写法基本稳定下来了。很多新手第一次写异步代码最容易犯的错import asyncio async def say_hello(): print(hello) await asyncio.sleep(1) print(world) say_hello() # 这行什么也不会输出控制台一片空白程序直接结束。因为调用一个async def函数并不会执行函数体它只是返回一个协程对象coroutine object。这个对象就像一个“待启动的任务说明书”要真正执行它你得把它交给事件循环。用asyncio.run()包一层asyncio.run(say_hello())这下才会输出hello等 1 秒再输出world。Python 会为没有被await的协程对象给出一个警告RuntimeWarning: coroutine say_hello was never awaited。这个警告本质上是在提醒你协程对象被创建了但你没有安排它执行。2.2 await 到底在等什么await后面可以跟的对象统称为可等待对象awaitable主要有三类原生协程对象async def函数调用后的返回值asyncio.Task被事件循环调度的任务对象asyncio.Future更低层的未来结果占位符await的行为需要好好理解它不是“站在这里傻等结果”而是“把控制权交还事件循环等结果准备好了再叫我回来”。这句话是异步的灵魂。import asyncio async def fetch_first(): await asyncio.sleep(2) return 第一个数据 async def fetch_second(): await asyncio.sleep(1) return 第二个数据 async def main(): # 串行等待:总耗时约3秒 a await fetch_first() b await fetch_second() print(a, b) asyncio.run(main())这里虽然是两个async函数但我用了两次await它们是顺序执行的总耗时 3 秒。await只负责让出控制权并不会魔改你代码的执行顺序。想要并发必须把任务包装成Task。2.3 把协程跑起来的三种正规姿势第一asyncio.run()只调一次作为异步程序的顶层入口asyncio.run(main())第二在协程内部await另一个协程async def main(): result await fetch_first()第三用asyncio.create_task()创建任务让它在后台并发执行async def main(): task asyncio.create_task(fetch_first()) # 此时 fetch_first() 已经进入调度队列,开始跑了 other await fetch_second() print(other) result await task # 最后再来拿第一个任务的结果create_task()是 3.7 之后推荐的写法它接收一个协程对象返回一个Task对象并将协程封装为事件循环中的一个调度单元。协程创建任务后立刻进入“待运行”状态事件循环会在合适的时机调度它。2.4 异步函数里常见的反例反例一在协程里用time.sleep()。async def bad_demo(): import time time.sleep(3) # 这是同步阻塞,会卡住整个事件循环如果协程里用了time.sleep()事件循环就会整个停摆所有其他任务都会被冻结。正确做法是await asyncio.sleep(3)。反例二把一个协程传入asyncio.run()之外的地方。# 错误 task asyncio.create_task(say_hello()) # 必须在一个运行中的事件循环里才能调用 # 正确 asyncio.run(say_hello())create_task()要求当前线程存在一个正在运行的事件循环所以它必须在协程内部调用或者通过asyncio.get_running_loop()获取事件循环之后再调用。反例三把async def函数当成普通函数直接 “调用后拿返回结果”。result fetch_first() # 这不是调用,是创建协程对象,result 不是字符串如果你拿到的是coroutine object而非具体值就说明你忘了await。3. 任务调度核心Task、gather、wait 怎么选3.1 用 create_task 实现真正的并发回到前面的例子要并发执行两个async函数需要这样写import asyncio async def worker(name, delay): await asyncio.sleep(delay) print(f{name} 完成了) return name async def main(): # 先创建任务,不等待 task_a asyncio.create_task(worker(A, 2)) task_b asyncio.create_task(worker(B, 1)) # 占用较少时间,后创建的 B 反而先完成 await task_a await task_b asyncio.run(main())输出顺序是B 完成了约 1 秒后、A 完成了约 2 秒后总耗时约 2 秒而不是 3 秒。因为两个协程都被封装成了 Task事件循环在两个任务之间来回切换B 的定时器 1 秒先到先打印A 的定时器 2 秒后到后打印。这里可以顺带体会一个调度细节await task_a写在前面不代表 task_a 先执行。Task 的调度是由事件循环决定的跟你代码里await的先后没有必然关系。3.2 gather 与 wait两个集合等待工具日常用得最多的是asyncio.gather()async def main(): results await asyncio.gather( worker(A, 2), worker(B, 1), ) print(results) # [A, B],保持传入顺序gather的返回结果保持传入协程的顺序即使它们完成的先后不一样。另外它还接收return_exceptionsTrue参数当某个任务抛异常时不会中断整个gather而是把异常对象放入返回列表。asyncio.wait()则更底层一些。它接收一组Task对象返回(done, pending)两个集合async def main(): task_a asyncio.create_task(worker(A, 2)) task_b asyncio.create_task(worker(B, 1)) done, pending await asyncio.wait( {task_a, task_b}, return_whenasyncio.FIRST_COMPLETED, ) for t in done: print(t.result()) # pending 里的任务可以让它继续,或取消掉 for t in pending: t.cancel()gather更适合“我要全部结果”wait更适合“我只需要最先完成的那批或者我需要在舞台上实时监控任务状态”。两者对应的心智模型不一样。Python 3.11 之后还引入了asyncio.TaskGroup它提供了比gather更干净的异常聚合逻辑async def main(): async with asyncio.TaskGroup() as tg: t1 tg.create_task(worker(A, 2)) t2 tg.create_task(worker(B, 1))如果任务组里有任务抛异常其他任务会被自动取消异常在退出async with块时统一抛出。这种“同生共死”的语义非常适合能容忍短事务失败、但整体必须原子化收尾的场景。3.3 超时控制与任务取消真实项目里网络请求永远可能超时。asyncio.wait_for()是最好用的超时工具import asyncio async def main(): try: result await asyncio.wait_for(worker(A, 10), timeout3) except asyncio.TimeoutError: print(任务超时了,已被取消)wait_for在超时后会自动取消内部任务并抛出asyncio.TimeoutError。这比在协程内部手动判断时间要稳健得多。如果你想取消一个任务用task.cancel()。取消后协程内部正在执行的await位置会抛出asyncio.CancelledError如果你在协程内部用 try/finally 清理资源finally 里的代码会执行但如果协程内部自行吞掉了CancelledError事件循环可能无法正常取消该任务导致后面的逻辑出现诡异问题。一个比较实用的资源清理模板async def worker(): try: await asyncio.sleep(10) except asyncio.CancelledError: print(任务被取消,正在清理资源...) raise # 必须重新抛出,否则取消流程会卡住有些场景你希望任务不被外部超时或取消影响可以用asyncio.shield()。它会创建一个“护盾”任务外层被取消了内层真正执行的任务不受影响这在新老接口替换、后台异步写日志的场景里很实用。但shield也容易造成任务泄漏新手可以先不用知道有这个东西就行。4. 实战搭一个可控并发的异步爬虫4.1 为什么异步爬虫选 aiohttprequests是同步阻塞库你没法在一个async def里直接await requests.get()。想在事件循环里发 HTTP 请求标准选择是aiohttp。aiohttp的核心组件是ClientSession。它和requests.Session类似内部维护连接池应当被发起请求的协程安全地共享复用。每次请求时直接session.get(url)用完不要急着关 session连接池还有用。先看最小可用版本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: return await resp.text() async def main(): async with aiohttp.ClientSession() as session: tasks [asyncio.create_task(fetch(session, url)) for url in URLS] pages await asyncio.gather(*tasks) print(len(pages), 个页面抓取完成) asyncio.run(main())这里有个新手容易犯的错把ClientSession()写在每个请求里面。每次新建 session 会重新走一遍 TCP 握手、TLS 协商连接池就形同虚设了。正确的做法是只创建一个 session把 20 个请求全部塞进去。4.2 给爬虫加上并发上限gather一次性把 20 个请求全发出去服务器可能扛不住你的本机连接数也可能爆。解决思路是用信号量import asyncio import aiohttp semaphore asyncio.Semaphore(5) async def fetch_with_limit(session, url): async with semaphore: async with session.get(url) as resp: return await resp.text() async def main(): async with aiohttp.ClientSession() as session: tasks [asyncio.create_task(fetch_with_limit(session, url)) for url in URLS] results await asyncio.gather(*tasks) print(f共 {len(results)} 个结果) asyncio.run(main())Semaphore(5)限制了同时最多 5 个协程进入async with semaphore区间。第 6 个协程会在信号量外侧等待等某个并发请求完成后再拿令牌进入。并发数设多少合理取决于你的目标和服务器承受能力。本地性能和延迟允许的场景10~50 是常见区间。如果目标是公开大站推荐 5~10 起步慢慢调别一上来就并发 200。4.3 加入超时、重试与状态码过滤到一个生产环境的爬虫还需要处理三件事超时session.get(url, timeoutaiohttp.ClientTimeout(total8))3xx/4xx/5xx 状态码判断后决定是否重试编码问题优先看响应头的charset没有就手动指定可以贴一个带重试的封装import asyncio import aiohttp RETRY_TIMEOUTS [1, 2, 4] async def fetch(url, session, semaphore): async with semaphore: for attempt, delay in enumerate(RETRY_TIMEOUTS): try: async with session.get(url, timeoutaiohttp.ClientTimeout(total8)) as resp: if resp.status 200: return await resp.text() if resp.status in (404, 403): return None await asyncio.sleep(delay) except (aiohttp.ClientError, asyncio.TimeoutError): if attempt len(RETRY_TIMEOUTS) - 1: print(f放弃 {url},重试全部失败) return None await asyncio.sleep(delay)重试间隔可以按指数退避加一点随机抖动避免多个爬虫进程同时重试把服务器打崩。RETRY_TIMEOUTS数组的取法是 1 秒、2 秒、4 秒第 4 次会放弃。这种数组写起来比backoff * factor ** attempt更直接也更容易在项目里根据接口特性微调。爬虫本身是一个典型的 IO 密集任务但如果你抓回来之后还要做正则匹配、HTML 解析、数据压缩这些属于计算密集操作塞在事件循环里会拖慢后面所有的请求。稳妥的做法是抓取阶段用异步并发解析阶段丢到线程池或者干脆开多进程。5. 踩坑词典异步代码最容易翻车的地方5.1 RuntimeError: asyncio.run() cannot be called from a running event loop这个报错几乎是所有异步新手都会遇到的。典型场景你在 Jupyter Notebook 里跑异步代码或者在一个协程里又调用了一次asyncio.run()。原因很简单事件循环是单例运行的同一线程内同一时刻只能存在一个正在运行的事件循环第二层asyncio.run()试图创建新的循环时被解释器拦截了。解决办法在 Jupyter 里改用await或者用nest_asyncio.apply()打补丁仅限临时环境别带到生产代码里。普通脚本里asyncio.run()只出现在顶层入口其他位置一律用await。5.2 同步阻塞代码混进事件循环异步并发最忌讳的就是在协程里用同步阻塞函数。常见的捣乱分子有requests.get()time.sleep()数据库 ORM 的同步查询subprocess.run()这些函数一旦被调用整个线程就卡住了事件循环里所有正在排队的任务都会被冻住。比如你用一个并发 50 的异步爬虫其中混了个requests.get()那一次执行期间其他 49 个任务全部原地罚站。处理办法是把这些阻塞操作放到线程池import asyncio import requests async def main(): loop asyncio.get_running_loop() resp await loop.run_in_executor(None, requests.get, https://httpbin.org/delay/1) return resp.status_coderun_in_executor(None, fn, ...)会把fn提交给默认的线程池执行器然后返回一个协程await它会拿到最终结果。等价写法还有asyncio.to_thread(fn, ...)3.9 之后提供更贴近普通人的直觉。换数据库也是同样的思路。异步 Web 框架里千万别直接调用同步 ORM 的查询要么用asyncpg这类异步驱动要么把查询丢给线程池。5.3 “coroutine was never awaited” 警告这个警告我见过太多次。要么是调用 async 函数时漏了await要么是忘了加async def要么是我们前文提到的“把协程对象挂在一棵树上却从不调度”。比较隐蔽的场景是装饰器或者魔法方法class Fetcher: async def run(self): await asyncio.sleep(1) return ok f Fetcher() asyncio.run(f.run())如果哪天你看到coroutine was never awaited第一反应就是去检查代码里有没有某个 async 函数被当作普通函数调用了该加await的地方是不是没有加。5.4 gather 里的异常处理别让一个小请求毁掉全部gather有个很迷惑人的行为默认情况下只要其中一个任务抛出异常gather整体就会中途中止异常会从await asyncio.gather(...)的位置抛出即使其他任务还在正常运行你也不会拿到它们的结果。解决方案很清晰加return_exceptionsTrueresults await asyncio.gather(*tasks, return_exceptionsTrue) for result in results: if isinstance(result, Exception): print(一个任务失败了:, result) continue print(result)拿到的结果列表里失败的位置会是异常对象实例成功的位置是正常返回值。逐个判断即可不会因为一个失败丢掉所有结果。Python 3.11 的TaskGroup则换了一种策略所有任务按正常流程并发跑一旦有任务异常其他任务自动取消然后在async with块退出后再统一抛出ExceptionGroup。需要一次性拿到所有异常列表的时候TaskGroup更符合直觉。6. 选型时刻asyncio、多线程、多进程到底怎么选6.1 并发和并行不是一回事并发concurrency是“同一时间段内处理很多任务”并行parallelism是“同一时刻同时执行很多任务”。事件循环只是并发不是并行它用单线程通过时间片轮转来表现“同时”。多线程也是并发同一个进程内多个线程轮流抢 GIL 执行字节码。Python 的 GIL 确实让纯 Python 多线程在 CPU 密集场景下没法利用多核但在 IO 等阻塞场景里线程还是会释放 GIL所以多线程做网络请求依然有效。多进程则是真并行每个进程独立解释器拥有独立的 GIL可以在多个 CPU 核心上同时推进计算任务。6.2 三套方案对比维度asyncio多线程多进程适用场景高并发的 IO 等待中小规模 IO 并发CPU 密集型计算资源开销极低协程几乎无空间成本每线程约 8MB 虚拟内存且有切换成本每进程开销最大并发上限上万任务轻松几百线程基本到顶根据核心数折中代码复杂度需要理解 async/await需要处理线程安全与锁需要处理进程间通信CPU 密集不擅长不擅长擅长一个有意思的细节是性能对比单纯做大量网络请求asyncio 的实现比多线程实现通常能承载更高的连接并发数但多线程的代码写起来更符合普通人的线性思维。6.3 混合使用在异步里跑线程池实际生产里没必要非要二选一。常见的混合形态是主流程用异步编排但个别的同步阻塞第三方库调用丢进线程池实现“异步框架同步线程池”的混搭这在 FastAPI、Sanic 这类异步 Web 框架里非常常见。另一个常见组合是 CPU 密集耗时任务用ProcessPoolExecutor跑这样既保住了异步主流程的高响应性又拿到了多核并行能力。唯一需要考虑的是进程间通信成本如果任务切分的粒度过小进程池的序列化开销会比计算本身还贵。写到这里想再强调一个观点不是所有代码都需要异步。如果你的程序只做 3 个请求、总量不到 1 秒“同步 for 循环”完全没问题。硬上 async 只会让代码更绕还要承担协程调度和异常处理的心智负担是用大炮打蚊子。判断准绳只有一个有没有大量 IO 等待且等待时间在整体耗时中占比够大。7. 最后分享几个实操体会真正把异步代码大规模落地之后我总结了一套比较稳的实践习惯。第一异步程序入口保持唯一。不管项目多大asyncio.run()只在最外层出现一次所有内部执行都用await或create_task。不要在多个地方循环里反复创建事件循环那既会出 RuntimeError也会让生命周期混乱。第二凡是涉及“批量发起任务”的都默认加上重试、加上并发上限。这两个缺一不可。没有并发上限就是给自己埋雷没有重试一个瞬时超时就能让整套任务失败。爬虫、批量接口调用、异步消息推送这三类场景我都在生产环境吃过亏。第三Python 版本建议直接用 3.11 或更新版本。TaskGroup、更快的asyncio实现、更好的异常报告都比在老版本里折腾省事得多。如果你在维护旧项目没有升级条件至少要把async with、wait_for、Semaphore这几个基础操作吃透它们是任何版本都通用的骨架。最后再分享一个小技巧排查异步死锁或者性能问题的时候开启asyncio的调试模式很有用。在启动时设置环境变量PYTHONASYNCIODEBUG1或者调用asyncio.run(main(), debugTrue)事件循环会额外检测超时未完成的任务、阻塞调用、未关闭的资源很多隐蔽问题在这个模式下几秒钟就暴露了。这个调试开关在本地开发时强烈建议开着它不改变代码行为只会多输出诊断日志等上生产再关掉不迟。
返回列表