Python异步编程:从async/await原理到高并发实战

📅 2026/8/1 13:13:35
Python异步编程:从async/await原理到高并发实战
1. 从“同步等待”到“异步协作”为什么我们需要 async/await如果你写过Python爬虫或者处理过需要等待网络请求、文件读写这类“慢操作”的程序大概率遇到过这样的场景程序卡在那里什么也做不了就为了等一个网页返回或者一个文件读完。CPU明明闲得发慌却只能干等着这就是典型的同步阻塞。在Python 3.5之前解决这类问题的主流方案是回调Callback和生成器协程如yield from代码写起来像“回调地狱”逻辑七拐八绕非常反直觉。async/await语法的引入就是为了用更符合人类线性思维的方式来编写高效的异步程序。你可以把它想象成组织一个高效的团队。同步编程就像让团队里最能干的小A去跑一个耗时很长的外勤在他回来之前整个团队都停下手里所有工作眼巴巴地等着。而异步编程则是小A出去后立刻给他挂上一个“任务进行中”的牌子然后团队转头就去处理其他任务比如让小B去打印文件让小C去整理数据。等小A办完事回来了系统会通知团队“小A的事办妥了结果在这里谁需要谁来取。”这样团队的总体效率吞吐量就大大提升了。async和await这两个关键字就是实现这套协作机制的核心。async用来声明一个函数是“可等待的”或者说它是一个协程函数它定义了一个可以暂停和恢复的任务单元。而await则是一个明确的“等待点”它表示“我知道这个操作比如网络请求需要时间我在这里主动让出控制权你去忙别的吧等这个操作有结果了再回到我这里继续执行。”这彻底改变了我们处理I/O密集型任务的编程模式。2. 核心概念再辨析协程、任务与事件循环在上一部分我们建立了基本认知现在需要深入理解驱动async/await运转的三个核心引擎协程Coroutine、任务Task和事件循环Event Loop。很多初学者混淆它们导致代码写出来跑不动或者达不到预期效果。2.1 协程Coroutine异步世界的基本执行单元协程本质上是一个可以暂停执行并在之后恢复的函数。用async def定义的函数就是一个协程函数。调用它并不会立即执行其中的代码而是返回一个协程对象。这个对象就是一个待执行的计划。import asyncio async def fetch_data(): print(开始获取数据...) await asyncio.sleep(2) # 模拟一个耗时2秒的I/O操作 print(数据获取完成) return {data: 123} # 这行代码不会打印任何东西它只是创建了一个协程对象 coro fetch_data() print(coro) # 输出coroutine object fetch_data at 0x...关键理解协程对象本身是“惰性”的。它就像一份写好的菜谱函数定义而协程对象是拿着这份菜谱、备好了料函数调用但还没开火烹饪的状态。你需要一个“厨师”事件循环来按照菜谱包含await的步骤执行它。2.2 任务Task被事件循环管理的协程任务是对协程的进一步封装。它将一个协程“提交”给事件循环进行调度和管理。创建任务后事件循环就会在合适的时机比如遇到await时自动驱动这个协程的执行。import asyncio async def main(): print(主函数开始) # 创建任务将协程 fetch_data() 包装成任务并加入事件循环的调度队列 task asyncio.create_task(fetch_data()) print(任务已创建主函数继续执行其他事情...) await asyncio.sleep(1) # 主函数也可以做其他“等待” print(主函数其他事情做完现在等待任务结果...) # 等待任务完成并获取结果 result await task print(f获取到的结果是{result}) asyncio.run(main())输出将会是主函数开始 任务已创建主函数继续执行其他事情... 开始获取数据... 此处主函数在 sleep(1)而 fetch_data 在 sleep(2) 主函数其他事情做完现在等待任务结果... 数据获取完成 获取到的结果是{data: 123}核心区别与联系协程是蓝图任务是工程。你直接await一个协程相当于现场按照蓝图施工并等待完工。而create_task则是把蓝图交给工程队事件循环去并行施工你可以继续做别的事稍后再来验收。任务提供了更强的控制力。你可以取消任务task.cancel()、检查任务是否完成task.done()、或设置回调函数。多个任务可以被并发地调度。一个常见的误区认为asyncio.create_task()是“启动”一个后台线程。不是的它仍然在同一个线程内由事件循环在单个线程上通过协作式多任务来切换执行。所有协程和任务都跑在同一个线程里。2.3 事件循环Event Loop异步程序的指挥中枢事件循环是asyncio应用的核心它负责在单个线程内调度和执行所有的协程任务并处理I/O事件如网络数据到达、文件操作完成。你可以把它想象成一个无限循环的调度器。它维护着两个队列一个是准备就绪的协程队列那些await的条件已满足可以继续执行的另一个是I/O等待队列那些正在等待I/O操作完成的协程比如socket.read()。它的工作流程从就绪队列中取出一个协程执行。协程执行直到遇到await表达式。await后面的对象通常是一个Future代表一个未来会完成的操作会被注册到事件循环的I/O等待队列中然后当前协程挂起。事件循环去检查I/O等待队列看看哪些操作已经完成例如一个网络请求收到了响应。将已完成操作对应的协程标记为就绪放回就绪队列。重复步骤1执行下一个就绪的协程。asyncio.run()的作用在Python 3.7以后asyncio.run(main())是我们运行异步程序的推荐入口。它帮我们做了三件重要的事1) 创建一个新的事件循环2) 将传入的协程通常是main()作为主任务运行3) 管理事件循环的关闭和清理。这避免了以前需要手动获取、运行、关闭循环的繁琐和潜在错误。注意在像Jupyter Notebook或某些已有事件循环的环境如GUI应用中直接调用asyncio.run()可能会报错RuntimeError: asyncio.run() cannot be called from a running event loop。此时你需要使用await main()或者使用环境特定的方式来运行协程。3. 实战用async/await重构一个经典场景理解了概念我们通过一个经典的例子——并发获取多个网页标题来对比同步和异步的写法并深入异步实现的细节。假设我们有10个URL需要获取。3.1 同步阻塞版本import requests import time def fetch_sync(url): 同步获取网页标题 try: resp requests.get(url, timeout5) resp.raise_for_status() # 简单提取title标签内容 title_start resp.text.find(title) title_end resp.text.find(/title) if title_start ! -1 and title_end ! -1: title resp.text[title_start7:title_end] else: title No title found return url, title[:50] # 截取前50字符 except Exception as e: return url, fError: {e} def main_sync(urls): start time.time() results [] for url in urls: result fetch_sync(url) results.append(result) print(fFetched: {result}) end time.time() print(f同步版本总耗时: {end - start:.2f}秒) return results if __name__ __main__: # 示例URL列表实际使用时请替换为有效的URL urls [https://httpbin.org/delay/1] * 5 # 模拟5个每个延迟1秒的请求 main_sync(urls)这个版本的问题是每个请求都必须等上一个完全结束后才能开始。如果每个请求耗时1秒5个请求就是5秒网络延迟被线性累加。3.2 异步并发版本我们使用aiohttp这个异步HTTP客户端库来重写。import asyncio import aiohttp import time async def fetch_async(session, url): 异步获取网页标题 try: async with session.get(url, timeout5) as resp: resp.raise_for_status() html await resp.text() # 同样简单提取标题 title_start html.find(title) title_end html.find(/title) if title_start ! -1 and title_end ! -1: title html[title_start7:title_end] else: title No title found return url, title[:50] except Exception as e: return url, fError: {e} async def main_async(urls): start time.time() results [] # 创建一个aiohttp客户端会话复用连接池提升性能 async with aiohttp.ClientSession() as session: # 为每个URL创建一个获取任务 tasks [asyncio.create_task(fetch_async(session, url)) for url in urls] # 并发执行所有任务并等待它们全部完成 completed_tasks await asyncio.gather(*tasks, return_exceptionsFalse) results.extend(completed_tasks) end time.time() print(f异步版本总耗时: {end - start:.2f}秒) for result in results: print(fFetched: {result}) return results if __name__ __main__: urls [https://httpbin.org/delay/1] * 5 asyncio.run(main_async(urls))关键解析与技巧async with与await resp.text()aiohttp.ClientSession()是一个异步上下文管理器使用async with来确保会话正确关闭。resp.text()也是一个协程方法它需要await来实际执行读取响应体的I/O操作。asyncio.create_task()与asyncio.gather()create_task将每个fetch_async协程包装成独立任务它们被“丢进”事件循环并发执行。asyncio.gather(*tasks)是并发执行的核心。它接收一系列可等待对象这里是任务并await它们全部完成。它不会按顺序等待而是同时“等待所有”。由于所有任务在遇到I/O网络请求时都会通过await让出控制权事件循环就可以在它们之间切换。当5个任务都在等待1秒的延迟时总耗时接近最慢的那个任务约1秒而不是5秒。return_exceptions参数默认为False。如果任何一个任务抛出未处理异常gather()会立即抛出该异常并取消其他所有任务。设置为True时任务中的异常会作为结果对象返回而不是抛出这样可以确保所有其他任务都能完成。根据你的错误处理策略谨慎选择。实测对比对于5个模拟延迟1秒的请求同步版本耗时约5秒而异步版本耗时仅略高于1秒。当请求数量增加到几十上百个且网络延迟不稳定时异步的优势将是数量级的。4. 高级模式与常见陷阱掌握了基础并发后我们会遇到更复杂的需求比如限制并发数、处理超时、取消任务等。同时一些常见的陷阱也值得警惕。4.1 控制并发度信号量Semaphore与任务组直接gather上百个任务同时发起请求可能会把目标服务器打垮或触发反爬机制也可能耗尽本地文件描述符。我们需要限制同时进行的任务数量。方案一使用asyncio.Semaphore信号量import asyncio async def worker(semaphore, url): async with semaphore: # 获取信号量如果已满则在此等待 print(f开始处理 {url}) await asyncio.sleep(1) # 模拟工作 print(f完成处理 {url}) return url async def main(): semaphore asyncio.Semaphore(3) # 限制并发数为3 urls [furl_{i} for i in range(10)] tasks [asyncio.create_task(worker(semaphore, url)) for url in urls] results await asyncio.gather(*tasks) print(results) asyncio.run(main())Semaphore(3)创建了一个计数为3的信号量。async with semaphore:上下文管理器会在进入时减少计数acquire离开时增加计数release。当计数为0时新的协程试图获取信号量就会被阻塞直到有其他协程释放。这就像只有3个工位的办公室确保了最多只有3个worker在同时“工作”。方案二使用asyncio.TaskGroupPython 3.11TaskGroup提供了更结构化、更安全的管理方式它自动处理了任务的取消和异常传播。import asyncio async def worker(url): print(f开始处理 {url}) await asyncio.sleep(1) if url url_5: raise ValueError(模拟任务失败) print(f完成处理 {url}) return url async def main(): urls [furl_{i} for i in range(10)] results [] # 使用TaskGroup管理一组任务 async with asyncio.TaskGroup() as tg: tasks [tg.create_task(worker(url)) for url in urls] # 离开async with块时会等待所有任务完成。 # 如果任何任务抛出异常所有其他任务会被取消异常会传播出来。 for task in tasks: if task.cancelled(): results.append(Cancelled) elif task.exception(): results.append(fError: {task.exception()}) else: results.append(task.result()) print(results) asyncio.run(main())TaskGroup的优点是1) 代码结构清晰2) 一个任务失败会自动取消同组所有任务避免资源泄漏3) 异常聚合所有异常信息都能被捕获。强烈推荐在Python 3.11中使用TaskGroup替代裸的gather进行任务分组管理。4.2 超时与取消网络请求不可靠必须设置超时。import asyncio async def slow_operation(): await asyncio.sleep(10) # 模拟一个很慢的操作 return Done async def main(): try: # 使用asyncio.wait_for设置超时 result await asyncio.wait_for(slow_operation(), timeout2.0) print(f结果: {result}) except asyncio.TimeoutError: print(操作超时) # 超时后wait_for会取消底层的任务。 # 但取消是协作式的需要任务内部能响应取消。 except asyncio.CancelledError: print(任务被取消。) asyncio.run(main())关于取消的深度理解asyncio.CancelledError是一个特殊的异常。当调用task.cancel()时事件循环会在下次该任务被调度时在它内部await的地方抛出这个异常。因此任务必须能在await点被中断。如果你的协程内部有一段纯CPU计算的密集循环没有await那么取消请求会一直被延迟直到下一个await出现。async def uncancellable(): try: i 0 while i 10_000_000: # 密集CPU计算没有await i 1 await asyncio.sleep(0) # 只有到这里取消信号才会被处理 print(循环结束) except asyncio.CancelledError: print(终于被取消了) raise async def main(): task asyncio.create_task(uncancellable()) await asyncio.sleep(0.001) # 稍等片刻 task.cancel() try: await task except asyncio.CancelledError: print(主函数捕获到取消) asyncio.run(main())这段代码中task.cancel()被调用后uncancellable函数会继续完成那个耗时的循环然后才会在await asyncio.sleep(0)处抛出CancelledError。要避免这个问题可以在长循环中定期使用await asyncio.sleep(0)来“让步”给事件循环一个处理取消或其他任务的机会。4.3 常见陷阱与避坑指南在同步函数中调用异步函数这是最常见的错误。你不能直接在一个普通的def函数里使用await。解决方案是使用asyncio.run()在程序顶层或者在另一个异步函数中调用。如果必须在同步上下文中运行异步代码例如在Web框架的同步视图函数里有些框架提供了工具如asgiref.sync.async_to_sync或者你可以使用loop.run_until_complete()但这需要你管理事件循环。忘记await调用一个async def函数会返回一个协程对象如果你忘记写await这个协程就不会被执行。有些静态类型检查工具如mypy或IDE可以帮你发现这个问题。错误示例result fetch_data()正确应为result await fetch_data()。阻塞事件循环在异步函数中执行了阻塞性的CPU密集型操作或同步I/O如time.sleep()、requests.get()、读写大文件而不使用异步版本会卡住整个事件循环导致所有其他任务“饿死”。解决方案CPU密集型操作用loop.run_in_executor()丢到线程池执行同步I/O操作寻找其异步替代库如aiofiles替代文件操作aiohttp替代requests。不正确的资源管理像数据库连接池、HTTP客户端会话如aiohttp.ClientSession这类资源必须在异步上下文async with中管理确保它们被正确创建和关闭。不要在每个协程里都创建一个新会话而应该在主函数或适当作用域内创建并复用。异常处理缺失异步任务中的异常默认不会立即抛出除非你await它或者检查task.exception()。使用asyncio.gather()时要注意return_exceptions参数。使用TaskGroup可以更好地聚合异常。5. 性能调优与最佳实践写出能跑的异步代码只是第一步写出高效、健壮的异步程序则需要遵循一些最佳实践。5.1 选择合适的并发模型I/O密集型asyncio的主战场。网络请求、数据库访问、磁盘I/O使用异步库等。CPU密集型asyncio不擅长。大量数学计算、图像处理、数据序列化等。此时应该使用concurrent.futures.ProcessPoolExecutor将计算任务分发到多个进程避免阻塞事件循环。asyncio可以通过loop.run_in_executor()与线程池/进程池配合。import asyncio import concurrent.futures import time def cpu_bound_task(n): 模拟一个CPU密集型任务 return sum(i * i for i in range(n)) async def main(): loop asyncio.get_running_loop() # 使用进程池执行CPU密集型任务 with concurrent.futures.ProcessPoolExecutor() as pool: start time.time() # 将函数提交到进程池并异步等待结果 result await loop.run_in_executor(pool, cpu_bound_task, 10_000_000) end time.time() print(f结果: {result}, 耗时: {end-start:.2f}秒) asyncio.run(main())5.2 监控与调试异步程序因为并发调试起来比同步程序更复杂。状态可能在你单步调试时改变。asyncio.run()的debug参数asyncio.run(main(), debugTrue)可以启用调试模式会输出更详细的日志例如慢回调警告默认超过100ms的回调会打印警告。结构化日志使用logging模块并在日志格式中包含asyncio.current_task().get_name()这样可以追踪是哪个任务在输出日志。可视化工具对于复杂应用可以考虑使用像tracy或专门的异步性能剖析工具。5.3 设计模式生产者-消费者与队列对于数据流处理场景asyncio.Queue是一个强大的工具它实现了协程安全的队列可用于构建生产者-消费者模型。import asyncio import random async def producer(queue, producer_id): for i in range(5): item f产品 {producer_id}-{i} await asyncio.sleep(random.random()) # 模拟生产时间 await queue.put(item) print(f生产者 {producer_id} 生产了 {item}) await queue.put(None) # 发送结束信号 async def consumer(queue, consumer_id): while True: item await queue.get() if item is None: # 把结束信号放回队列让其他消费者也能结束 await queue.put(None) break print(f消费者 {consumer_id} 消费了 {item}) await asyncio.sleep(random.random() * 2) # 模拟消费时间 queue.task_done() # 通知队列该项已被处理 print(f消费者 {consumer_id} 结束) async def main(): queue asyncio.Queue(maxsize3) # 设置队列容量 # 创建生产者和消费者任务 producers [asyncio.create_task(producer(queue, i)) for i in range(2)] consumers [asyncio.create_task(consumer(queue, i)) for i in range(3)] # 等待所有生产者完成 await asyncio.gather(*producers) # 等待队列中所有项目被消费完 await queue.join() # 取消消费者它们会在收到None后自行退出这里确保退出 for c in consumers: c.cancel() # 等待消费者任务正式结束处理取消 await asyncio.gather(*consumers, return_exceptionsTrue) print(所有任务完成) asyncio.run(main())这个模式可以很好地解耦生产速度和消费速度平衡负载是构建高效数据处理管道的基础。5.4 测试异步代码测试异步代码需要使用支持异步的测试框架如pytest配合pytest-asyncio插件。# test_async_code.py import pytest import asyncio from my_async_module import fetch_data pytest.mark.asyncio async def test_fetch_data_success(): 测试成功获取数据 # 可以使用unittest.mock的AsyncMock来模拟异步依赖 result await fetch_data(https://example.com) assert data in result assert result[data] 123 pytest.mark.asyncio async def test_fetch_data_timeout(): 测试超时情况 with pytest.raises(asyncio.TimeoutError): # 假设fetch_data内部使用了wait_for await fetch_data(https://httpbin.org/delay/5, timeout1.0)使用pytest.mark.asyncio装饰器告诉pytest这是一个异步测试函数。在测试中你可以像在普通异步代码中一样使用await。异步编程是Python处理高并发I/O的利器async/await语法让它变得清晰易读。核心在于理解事件循环驱动协程在单线程内协作这一模型。从简单的并发请求到使用信号量、队列控制流程再到处理异常、超时和调试每一步都需要对“协作式多任务”有深刻的理解。避免在协程中阻塞善用TaskGroup和Queue等高级原语你的异步程序就能既高效又健壮。记住异步不是银弹它完美解决了I/O等待问题但将CPU密集型任务交给它反而会适得其反。