写了个脚本要批量请求几百个URL,结果卡在原地干等网络响应,CPU几乎没动静,进度条半天走一格。相信不少朋友在Python里碰到过这种场景,于是找到了async、await和asyncio这套异步编程方案。Python异步编程从3.4引入asyncio库、3.5正式提供async/await语法,到现在已经是写高并发IO密集型任务的标配能力:它让你在单线程内同时管理成百上千个网络连接、文件读写或数据库查询,而不必为每个任务开线程、耗内存、拼切换。
这篇实战总结会从“异步到底解决了什么问题”开始,一步步拆解事件循环、协程、Task这几个核心概念,然后用完整可复现的代码做一次并发请求与批量处理实战,最后把我这些年踩过的坑、排查技巧和几个容易混淆的设计模式一并整理出来。适合刚接触异步的Python开发者,也适合写过一些async代码但总觉得“差一口气”的朋友。
1. 异步编程解决的核心问题:别让CPU空等IO
1.1 同步代码为什么慢:等待全部是死等
先看最传统的写法。假设用requests库逐个请求100个接口,每个接口平均响应200毫秒,整个流程大概需要20秒。问题不在于请求本身有多慢,而在于“发出请求后等待响应的这段时间”,程序什么都没干。
同步请求的过程可以理解成:你打电话给客服,拨号之后一直握着听筒不说话,直到对面接起来才继续交流。如果同时有10个电话要打,你只能一个一个来,每个电话的通话时间就是你的总时间。
网络IO、磁盘IO、数据库查询这类操作有一个共同特点:CPU发出指令后,设备开始工作,CPU就闲下来了,但它还在“傻等”结果返回。一个进程里哪怕只有10个并发请求,同步写法也会把它们变成10个串行等待,单个请求的延迟被反复叠加。
1.2 多线程方案的问题:线程不是越多越好
有人会说:那用多线程啊,每个线程处理一个请求不就行了?
早期我确实这么干过。concurrent.futures.ThreadPoolExecutor配合requests,代码改动小,效果也立竿见影。但随着并发量上来,问题就暴露了:
- 每个线程有自己的栈空间,动辄几十KB到几MB,几千个线程直接吃掉大量内存。
- 线程切换由操作系统调度,线程越多,上下文切换的CPU开销越大。
- Python的GIL(全局解释器锁)导致线程在CPU密集任务上无法并行,IO密集任务虽然可以释放GIL,但锁竞争依然存在。
- 线程池大小需要反复调参,调小了并发不够,调大了反而更慢。
我见过一个项目用ThreadPoolExecutor开200个线程去抢票,结果服务端没挂,本机CPU先飙到100%。这就是线程方案在超高并发下的典型困境。
1.3 协程的本质:在用户态自己安排“等待时间”
异步协程的思路完全不同。它保留单线程单进程,但是把“等待”变成“挂起”:遇到IO等待时,当前协程主动告诉事件循环“我先歇着,你有别的活儿就干别的”,等IO结果就绪了,再回来继续执行。
再拿打电话举例:同步是挨个电话握着听筒死等;协程是同时拨出100个电话,拨通后才接听讲话。话务员只有一个人(单线程),但他不需要在每个电话上干等,而是谁的声音来了就接谁的。
这个过程涉及两个关键词:
async def定义的是一个协程函数,调用它不会立即执行,而是返回一个协程对象。await就是“挂起点”,遇到它,协程让出控制权,事件循环去调度其他任务。
真正管理这些协程的,是asyncio的事件循环(Event Loop)。它是一个超级调度器,维护着一个就绪队列和等待队列,不停地在“检查IO状态、执行已就绪的协程、挂起等待中的协程”之间循环。
1.4 什么时候该用异步:IO密集是主场,CPU密集别凑热闹
用异步之前先做判断:你的任务是IO密集还是CPU密集?
- 网络请求、文件读写、数据库操作、消息队列消费,这类任务的特点是大部分时间花在等待上,适合异步。
- 图片处理、加解密、数据压缩、数值计算,这类任务吃CPU,协程帮不上忙,反而因为调度开销更慢。这时候该用多进程或直接上
numpy、Cython这类工具。
实际操作中我判断标准很简单:如果程序里大量时间都花在“等”上,就用asyncio;如果大量时间花在“算”上,就别凑热闹。
2. async/await与事件循环的工作原理
2.1 协程不是线程,也不是普通函数
初学者最容易混淆的概念就是:协程到底是什么?
从代码层面看,它仍然是一个函数,只不过用async def声明。但它的执行方式完全不同:
async def hello(): print("开始执行") await asyncio.sleep(1) print("执行结束")直接调用hello()不会打印任何内容,只会返回一个coroutine对象。要让里面代码真正跑起来,必须把它交给事件循环:
asyncio.run(hello())这个coroutine对象可以理解成一个“暂停的演员”:它知道自己从哪开始、到哪暂停、暂停后从哪继续。函数里每个await都是一个暂停标记,await右侧的表达式执行完毕之前,协程不会继续往下走。
协程最大的特征是“协作式调度”:协程自己决定何时让出,而不是被操作系统强占。这意味着协程只有在await处才会暂停,普通计算代码一旦开跑就停不下来,所以协程里千万不要写耗时很长的同步循环。
2.2 事件循环:核心调度器
事件循环可以理解成一个“老板”,手底下全是协程员工。员工们轮流汇报:“我在等网络请求”“我在等数据库返回”,老板就把他们挂到等待列表里,然后去推进其他可以继续工作的员工。
事件循环内部维护了多组数据结构:
_ready队列:存放已就绪、可以立即执行的协程。_scheduled:存放还没到时间的定时任务。- 各种IO事件监听器:底层通过
selectors模块监听文件描述符的可读可写状态。
当await asyncio.sleep(1)执行时,事件循环注册一个1秒后的定时回调,然后把协程挂起。这1秒内循环去执行其他任务,等时间到了再把协程放回就绪队列。
asyncio.run()本质上是做三件事:创建新的事件循环、把传入的协程作为第一个任务执行、最后关闭循环清理资源。它是Python 3.7引入的,之前用loop = asyncio.new_event_loop(); loop.run_until_complete(coro),现在统一推荐用asyncio.run。
2.3 await到底在等什么:可等待对象
await后面只能跟“可等待对象”(awaitable),主要分三类:
- 协程对象(
coroutine):来自async def函数的调用结果。 asyncio.Task:已经被事件循环调度的任务,比裸协程多一层调度管理。asyncio.Future:更底层的对象,表示一个“将来才会有结果”的操作,Task是Future的子类。
一个关键区别:直接await一个协程,是“等它执行完”;asyncio.create_task(coro)是把协程包装成任务、立刻丢给事件循环调度,然后你可以在之后某个时间点await这个任务。前者是串行等待,后者是并发安排。
async def main(): task1 = asyncio.create_task(hello()) task2 = asyncio.create_task(hello()) await task1 await task2这个写法两个hello()才真正并发执行。如果写成:
async def main(): await hello() await hello()那就是严格串行:第一个跑完,第二个才开始。我见过不少半懂不懂的代码把async函数一个接一个await,跑完发现性能毫无提升,问题就出在这。
2.4 Task的调度时机
create_task只是把协程“注册”进事件循环,不代表立即执行。事件循环会挑选合适的时机运行它。所以在create_task之后、await之前,任务可能还没真正跑起来,甚至在单任务场景下,await就是给事件循环机会去执行它。
这里必须理解:await作为“挂起点”,不只是等待结果,更是“让出CPU、让事件循环有机会调度其他任务”的关键动作。没有await,事件循环永远没有机会切换任务,再多的create_task也不会并发。
3. 完整实战:并发请求、超时控制与并发限制
3.1 场景设定
先设定一个实际场景:批量抓取100个商品详情页面的JSON数据,解析出关键字段后写入本地文件。这个场景非常典型,网络请求(IO)加文件写入(IO),全是异步的用武之地。
基础环境说明:
- Python版本3.10以上(3.7+都可运行,但新版本对
asyncio的API友好很多)。 - 需要安装
aiohttp:pip install aiohttp。注意requests是同步库,不能直接用在async函数里,它的阻塞式IO会卡住整个事件循环。 - 文件写入用
aiofiles:pip install aiofiles。普通open().write()虽然语法不报错,但它同步阻塞,文件一大整个事件循环就僵住了。
3.2 基础版:用asyncio.gather并发执行
先写一个最直接能跑通的基础版:
import asyncio import aiohttp import aiofiles import json async def fetch_one(session, url): async with session.get(url) as resp: resp.raise_for_status() return await resp.json() async def save_to_file(data, index): async with aiofiles.open(f"./data_{index}.json", "w", encoding="utf-8") as f: await f.write(json.dumps(data, ensure_ascii=False, indent=2)) async def main(): urls = [f"https://api.example.com/product/{i}" for i in range(1, 101)] async with aiohttp.ClientSession() as session: tasks = [fetch_one(session, url) for url in urls] results = await asyncio.gather(*tasks, return_exceptions=True) tasks = [] for i, result in enumerate(results): if isinstance(result, Exception): continue tasks.append(save_to_file(result, i)) await asyncio.gather(*tasks) asyncio.run(main())几个关键点解释一下:
aiohttp.ClientSession建议全局复用,不要每个请求新建一个。Session内部维护连接池,反复创建销毁会浪费大量资源和时间。return_exceptions=True非常关键。gather默认遇到第一个异常就立刻抛出,导致后面的任务被取消。生产环境里肯定不希望因为一个URL挂了就丢掉全部结果,所以让异常作为返回值保留,之后再逐个判断。- 这里的
results是按传入tasks的顺序返回的,不是按完成时间,这个特性在需要“结果与请求对应”的场景特别省心。
3.3 升级版:加超时控制、并发限制与重试
基础版能跑,但离生产可用还差两样东西:失控的并发量和慢请求拖垮整体。
100个请求同时发出去,本地可能承受得住,但对端服务器未必。互联网上很多接口对单IP并发数有限制,一口气打过去直接触发封禁。排查问题时要先把并发量压下来。
使用Semaphore是最简单的限流手段:
import asyncio import aiohttp semaphore = asyncio.Semaphore(10) # 最多同时10个请求 async def fetch_with_limit(session, url): async with semaphore: return await fetch_one(session, url)信号量的原理像商场限流:门口保安数人头,放进10个人,里面人出来一个,才放进下一个。这里的“人”就是协程,“保安”就是信号量内部的计数器。
超时控制同样不能少。一个API如果一直不返回,任务就会一直挂在那占着资源。用asyncio.wait_for包一层,给每个请求设一个最长时间:
async def fetch_with_timeout(session, url, timeout=5): try: return await asyncio.wait_for(fetch_with_limit(session, url), timeout=timeout) except asyncio.TimeoutError: return Nonewait_for的原理是给内部任务设置一个超时回调,超时就取消它并抛异常。取消动作本身也需要事件循环去执行,不会一次到位。
重试逻辑我习惯单独抽一个装饰器式函数,尽量让主流程保持干净:
async def fetch_with_retry(session, url, retries=3, timeout=5): for attempt in range(retries): try: return await asyncio.wait_for(fetch_with_limit(session, url), timeout=timeout) except Exception as e: if attempt == retries - 1: raise await asyncio.sleep(0.5 * (attempt + 1)) # 退避注意重试时await asyncio.sleep(0.5 * (attempt + 1))使用的是异步睡眠,不会阻塞事件循环。新手很容易在这里顺手写time.sleep(0.5),协同程序一睡,整个循环卡住,其他99个请求全部等死。
3.4 完整代码:一个可复制的并发抓取脚本
把限流、超时、重试、自动重命名文件、失败记录整合一下,写成一个可以直接拿去改的完整脚本:
import asyncio import aiohttp import aiofiles import json from datetime import datetime CONCURRENCY = 10 TIMEOUT = 5 RETRIES = 3 BASE_URL = "https://api.example.com/product/{}" OUTPUT_DIR = "./results" FAILED_LOG = "./failed.log" semaphore = asyncio.Semaphore(CONCURRENCY) async def fetch_one(session, url): async with semaphore: async with session.get(url) as resp: resp.raise_for_status() return await resp.json() async def fetch_with_retry(session, url): for attempt in range(RETRIES): try: return await asyncio.wait_for(fetch_one(session, url), timeout=TIMEOUT) except Exception: if attempt == RETRIES - 1: raise await asyncio.sleep(0.5 * (attempt + 1)) async def save_result(data, index): async with aiofiles.open(f"{OUTPUT_DIR}/data_{index}.json", "w", encoding="utf-8") as f: await f.write(json.dumps(data, ensure_ascii=False, indent=2)) async def main(): urls = [BASE_URL.format(i) for i in range(1, 101)] async with aiohttp.ClientSession() as session: tasks = [fetch_with_retry(session, url) for url in urls] results = await asyncio.gather(*tasks, return_exceptions=True) save_tasks = [] failed_count = 0 for i, result in enumerate(results): if isinstance(result, Exception): failed_count += 1 continue save_tasks.append(save_result(result, i)) await asyncio.gather(*save_tasks) print(f"完成,成功 {len(results) - failed_count}/{len(results)}") if __name__ == "__main__": asyncio.run(main())这个脚本设计思路是:限流压住对端压力,超时防止单点拖累,重试弥补临时故障,异步文件写入避免IO阻塞。每个环节单独拆开都能复用,组合起来就是一套可靠的小规模抓取骨架。
3.5 同步与异步版本对比
写完异步版本,我特意用同步requests写了个等价对照,用100个请求做实测对比(模拟接口,平均延迟200ms):
| 方案 | 总耗时 | 内存占用 | 并发能力 | 代码复杂度 |
|---|---|---|---|---|
| requests串行 | 约20秒 | 低 | 1 | 最低 |
| ThreadPoolExecutor(10线程) | 约2秒 | 中 | 10 | 中 |
| asyncio + aiohttp(限流10) | 约2.1秒 | 低 | 10 | 中高 |
| asyncio + aiohttp(不限流) | 约0.3秒 | 低 | 100 | 中高 |
结论很清楚:协程方案在并发量大的场景下,内存占用远低于线程方案,而且不需要手动调整线程池参数,并发上限由你写的信号量自行控制。如果对端服务完全不限流,100个并发请求在异步下几乎瞬间完成。
4. 事件循环的进阶操作与任务编排
4.1 用asyncio.wait做到“部分完成即返回”
gather的特性是所有任务都要等,但有些场景只需要“最快的一个结果”。比如同时访问多个可用性相同的API,谁先返回用谁的。此时asyncio.wait比gather更合适。
async def fetch_fastest(session, urls): tasks = [fetch_one(session, url) for url in urls] done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED) # 最快那个任务的结果 fastest_result = done.pop().result() # 取消还没完成的任务 for task in pending: task.cancel() return fastest_resultreturn_when参数支持三种模式:
FIRST_COMPLETED:第一个完成即返回。FIRST_EXCEPTION:出现第一个异常即返回,适合快速失败场景。ALL_COMPLETED:全部完成才返回,等价于gather默认行为。
wait返回两个集合:已完成的任务集合和未完成的任务集合。拿到done后,记得把pending里的任务取消掉,否则它们在后台继续执行,白白消耗资源。
4.2 用asyncio.Queue做生产者消费者
有些场景任务不是一次性生成完毕的,比如从分页接口持续拉数据、把数据分批写入数据库。这时候asyncio.Queue就派上用场了。
import asyncio import aiohttp async def producer(session, queue): for page in range(1, 20): data = await fetch_one(session, f"https://api.example.com/list?page={page}") for item in data["items"]: await queue.put(item) # 发送结束信号 await queue.put(None) async def consumer(session, queue): while True: item = await queue.get() if item is None: break await process_item(item) queue.task_done() async def main(): queue = asyncio.Queue(maxsize=50) async with aiohttp.ClientSession() as session: producer_task = asyncio.create_task(producer(session, queue)) consumer_task = asyncio.create_task(consumer(session, queue)) await asyncio.gather(producer_task, consumer_task) asyncio.run(main())Queue在这里相当于一条流水线:生产者从分页接口拿原始数据放到传送带上,消费者从传送带上取数据逐条处理。maxsize=50防止生产速度远超消费速度时内存暴涨。
值得说明的是Queue的get和put都是异步的:队列为空时get挂起等待;队列满时put挂起等待。这种天然的反压机制让生产消费节奏自动匹配。
4.3 协程之间的通信与结果传递
协程之间不能直接共享变量,应该通过返回值、队列、或者显式传参。初学者容易在协程函数内部修改一个全局字典,然后发现结果并发写入乱套。
正确的做法是在协程内部把结果返回,由上层统一收集。也就是我前面示例的写法:每个fetch_one返回自己的结果,gather汇总成列表。如果协程之间必须传递数据,用Queue或者asyncio.Event做同步。
asyncio.Event适合一个协程等待另一个协程完成某个前置动作的场景。比如先登录拿到token,再并发请求业务接口:
event = asyncio.Event() token = {} async def login_worker(): await asyncio.sleep(1) token["value"] = "fake-token" event.set() async def request_worker(): await event.wait() # 等登录完成 token = token["value"] # 继续请求4.4 多协程的优雅退出与超时兜底
一个常见的线上问题是:某个协程阻塞在第三方接口上,无论怎么设wait_for都不退出,整个进程无法正常关闭。这种情况通常是对端连接一直没有被关闭,wait_for超时后触发了cancel,但底层连接还在等待。
兜底方案是给整个主流程加一个总超时:
try: await asyncio.wait_for(main(), timeout=60) except asyncio.TimeoutError: print("主流程超时,强制退出")另外,事件循环关闭前应确保所有任务已完成或被取消。asyncio.run虽然会自动处理,但如果你自己管理循环,记得在退出前遍历所有任务:
def cleanup(loop): pending = asyncio.all_tasks(loop) for task in pending: task.cancel() loop.run_until_complete(asyncio.gather(*pending, return_exceptions=True))5. 常见问题与调试技巧实录
5.1 阻塞调用卡死事件循环:time.sleep是头号杀手
这是异步编程里最容易犯、后果也最严重的错误。在协程函数里使用time.sleep(1)代替asyncio.sleep(1),效果是整个事件循环冻结1秒,所有并发任务全部停摆。
为什么会这样?因为time.sleep是同步阻塞调用,它让当前线程休眠,而事件循环在这个线程上运行,自然跟着一起休眠。asyncio.sleep则是把“休眠”注册给事件循环,协程挂起后让出控制权,其他任务照常执行。
同理,在协程里调用requests.get()、open().read()、subprocess.run()都会阻塞事件循环。碰到这类需求要么换成异步库(aiohttp、aiofiles、asyncio.create_subprocess_shell),要么用loop.run_in_executor把同步操作丢到线程池执行。
import asyncio async def call_blocking(): loop = asyncio.get_running_loop() result = await loop.run_in_executor(None, requests.get, "https://api.example.com") return result5.2 为什么用了async却没有加速
排查思路从代码结构入手。最常见的两种原因:
第一,把所有async函数逐个await,没有用create_task或gather并发编排,写法变成了“异步的串行”。加速的前提是任务之间有独立的等待时间,并且这些等待被并发了。
第二,某个协程内部执行的是CPU密集操作,比如for循环里做了大量计算,没有await点,虽然函数声明了async,但其他协程根本无法在它运行期间插进来。这就是前面说的“协程是协作式调度,没有await就没有切换机会”。
5.3 调试工具:asyncio的调试模式与日志
asyncio自带了调试模式,开启后会在每个IO操作前记录耗时,帮你在任务卡顿时找出“是谁占着资源不放”。
在代码里设置:
import asyncio import sys if sys.flags.dev_mode: asyncio.run(main())或者运行脚本时加参数:
python -X dev script.py开启调试模式后,事件循环会警告那些执行时间过长的回调,打印出“Task was destroyed but it is pending”这类信息——后者通常是因为任务还没完成就被垃圾回收了,或者忘记await任务。
另一个实用技巧是给任务起名字,排查并发问题时一目了然:
task = asyncio.create_task(fetch_one(session, url), name=f"fetch-{i}") print(task.get_name())5.4 任务取消的坑:CancelledError怎么处理
task.cancel()会向任务内部抛入一个CancelledError,如果协程内部不处理,任务会静默终止。但如果在finally块里又执行了await,就会抛出另一个CancelledError,导致清理逻辑不完整。
推荐的处理方式是使用asyncio.shield保护关键清理操作,或者明确接收并吞掉取消异常:
async def worker(): try: while True: await do_work() except asyncio.CancelledError: # 做必要的清理 cleanup() raise # 重新抛出,保持取消语义注意:CancelledError在Python 3.8之后继承自BaseException而非Exception,所以普通except Exception是捕获不到它的。很多人在协程里写了兜底异常处理,但任务取消后却怎么都找不到原因,问题就在这里。
5.5 asyncio与多线程的边界
有些项目试图在异步代码里使用threading.Lock,或者在多线程环境里调用asyncio.get_event_loop(),结果各种报错。
几个明确的边界:
asyncio的Lock、Queue、Semaphore只能在协程里用,不能跨线程直接用。- 在另一个线程中调度协程,需要
asyncio.run_coroutine_threadsafe(coro, loop),它会返回一个concurrent.futures.Future供线程间同步。 - 多线程与多协程混合时,建议明确区分职责:线程负责阻塞式IO或CPU密集型任务,协程负责高并发网络请求,两者通过
run_in_executor或run_coroutine_threadsafe作为桥梁。
5.6 常用问题速查表
| 现象 | 可能原因 | 解决方案 |
|---|---|---|
| 协程代码“不执行” | 没有用asyncio.run或create_task调度 | 确认协程被交给事件循环 |
用async def但性能没提升 | 逐个await导致串行 | 用create_task+gather并发编排 |
| 所有协程卡住不动 | 协程内用了time.sleep、requests等阻塞调用 | 换成异步库或run_in_executor |
| 任务报“was destroyed but pending” | 任务未await就结束生命周期 | 保存任务引用,确保await完成或cancel |
except Exception捕获不到取消 | CancelledError是BaseException | 单独写except asyncio.CancelledError |
| 某个请求超时不生效 | wait_for没用或有嵌套阻塞逻辑 | 直接包住最终IO等待,逐层检查 |
| 对端服务器封禁IP | 并发量太高 | 用Semaphore限制并发,增加退避重试 |
| 回调函数里用异步库报错 | 回调是同步上下文 | 不要在回调里await,改用create_task包装 |
6. 异步项目里的代码组织与设计建议
6.1 按职责拆分:IO层、业务层、调度层
把异步代码写到一定规模后,比如几千行,就会觉得协程混在一起非常可怕。一个函数既发请求、又解析数据、又写文件、还管重试,排查时根本分不清是哪个环节出问题。
我习惯把代码拆成三层:
- IO层:只负责和外部系统打交道。
fetch_one、save_result这一层,参数是普通数据,返回也是普通数据,里面只做IO。 - 业务层:负责解析、清洗、校验数据,不涉及任何IO。这一层甚至不需要
async,普通同步函数就行。 - 调度层:负责创建任务、控制并发、处理异常、编排流程。
main函数和所有create_task、gather、Semaphore都在这一层。
这样的好处是:调换API、更换数据库驱动时只改IO层;调整并发策略、限流规则时只改调度层;业务逻辑变动完全不碰异步相关代码。
6.2 别把所有函数都写成async
一个常见误区:项目里到处都是async def,连加法运算都写成异步。这既难看又低效。
实际上,异步只有在“确实有IO等待”时才有价值。纯计算、纯字符串处理、纯数据转换,写普通函数即可。如果业务函数内没有await,却不小心把它写成了协程,那么每次调用都要包一层Task,调度开销反而更大。
判断标准很简单:函数内部是否出现await,没有就不该加async。
6.3 日志与监控:异步代码的调试线索
异步代码方便的地方在于单线程内极难看到“同时执行”的效果,也因此排查并发问题时如果没有日志,完全两眼一抹黑。我通常在关键节点打日志,记录任务开始、完成、失败和耗时:
import time import asyncio async def log_wrapper(coro, name): start = time.perf_counter() try: result = await coro print(f"[{name}] 完成,耗时 {time.perf_counter() - start:.2f}s") return result except Exception as e: print(f"[{name}] 失败:{e}") raise批量任务里,耗时明显偏长的名字往往是网络抖动或接口慢请求的线索。日志里附上任务名能快速定位到具体请求。
6.4 版本兼容性与API变迁
asyncio的API变化挺大,老项目升级Python版本时容易踩坑:
- Python 3.7之前:
asyncio.get_event_loop()在无当前循环时会创建并设置新循环,3.10之后在协程内部调用它会抛RuntimeError。 asyncio.run是3.7引入,之前版本只能手动管理循环。asyncio.current_task替代了老的asyncio.Task.current_task()。- Python 3.10之后,
asyncio.get_event_loop()的默认行为是如果当前没有正在运行的循环,抛出DeprecationWarning乃至RuntimeError。 - 3.11后协程任务异常处理、
TaskGroup也做了增强,asyncio.TaskGroup可以自动管理一组任务的错误传播和取消。
写跨版本兼容代码时,优先使用asyncio.run、asyncio.create_task、asyncio.current_task这些“新而稳”的API,避免直接操作底层loop。
6.5 与第三方异步库配合时的注意事项
异步生态里,每个第三方库都要有对应的异步实现,绝不能混用同步版本。常见的配套关系:
- HTTP:
aiohttp对应requests。 - 文件:
aiofiles对应open()。 - 数据库:
asyncpg、aiomysql对应psycopg2、pymysql。 - Redis:
redis.asyncio对应redis。 - ORM:
SQLAlchemy异步模式或Tortoise ORM。
选型时先看库是否实现了异步协议,没有就找替代。混用同步库到协程代码里,事件循环照样被阻塞。另一个注意事项是某些同步库内部实现C扩展并持有GIL,即使用run_in_executor也不一定能完全并行,需要重点关注。
7. 参考模式与个人实操心得
7.1 一个典型的生产消费编排模板
很多IO密集型系统都可以抽象成“获取数据、处理数据、落盘/入库”三段式。我总结出一个比较通用的模板,新项目起步时可以直接套用:
async def run_pipeline(concurrency=20): semaphore = asyncio.Semaphore(concurrency) queue = asyncio.Queue(maxsize=100) async def producer(): # 获取数据源 async for item in fetch_all_items(): await queue.put(item) await queue.put(None) # 结束标记 async def worker(): while True: item = await queue.get() if item is None: break async with semaphore: await process_one(item) queue.task_done() async with aiohttp.ClientSession() as session: # 每个worker实际上可以共用同一个session workers = [asyncio.create_task(worker()) for _ in range(concurrency)] await producer() await asyncio.gather(*workers)这个模板把获取、分发、执行、收尾都拆开,并发数通过concurrency一个参数控制,够用且清晰。
7.2 什么时候坚决别用异步
异步并非银弹,有几类场景我倾向直接拒绝异步方案:
- 项目规模很小,总共三五个请求串行也能接受,没必要引入异步框架增加理解成本和心智负担。
- 团队中没人懂事件循环,后续维护大概率出问题,写出来的“异步代码”可能比同步代码更慢、更乱。
- 对端服务完全同步且无法配合,比如只能串行交互的全双工协议,异步带来的收益有限。
- 任务本身是CPU密集计算,且没有调用外部服务,这时候异步约等于原地打转。
7.3 从同步到异步:我的迁移经验
如果手头有一套运行良好的同步代码,不要一次性全部重写成异步。我推荐步步为营的迁移路径:
第一步,先理清IO边界:哪些地方在等待外部系统,哪些在做计算,标注出来。
第二步,把耗时最长的IO操作改造成异步,单独跑一个异步入口,验证事件循环能正常工作。
第三步,逐步把其他同步IO替换为异步库,用Semaphore控制并发,压测对比指标。
第四步,清理平台化的问题:定时任务调度、信号处理、日志刷新等。
整个迁移过程要持续跑同一批回归测试,确保逻辑一致。
7.4 最后再分享一个排查卡死问题的方法
遇到协程“卡死”的疑难杂症,直接可以用faulthandler打印当前执行栈,一眼就能定位卡在哪一行:
python -X faulthandler script.py或者代码里主动开启:
import faulthandler import sys faulthandler.enable(sys.stderr)这套组合拳我救过几次场:上次生产环境一个采集任务突然卡住,就是靠它定位到有段代码在协程里偷偷用了time.sleep(10),三个worker全被卡死,整个队列堵了五分钟。从那以后我给自己定了一条规矩:协程函数里禁止出现任何同步阻塞调用,审查代码时看到time.sleep、requests.get、open().read()直接打回。这套经验帮助团队少踩了无数坑,希望也能帮到你。