Python异步革命收官,结构化并发高并发场景实战落地方案
文章目录
你说这年头写个高并发程序,咋就跟开饭店似的?以前咱们用asyncio,那简直就是雇了一群散工——你喊一嗓子"都去干活",结果有的摸鱼去了,有的偷偷跑路,还有的炸锅了(报错)你都不知道,等客人骂上门才发现菜根本没做。蚌埠住了,真的。
好在Python 3.11之后,官方终于看不下去这种"野生并发"的混乱局面,反手掏出了TaskGroup(任务组)和ExceptionGroup(异常组)这两张王炸。说白了,这就是给咱的异步代码配了个"托儿所园长"——所有任务都得签到,放学一个都不能少,谁尿裤子(抛异常)了立马全园通报。今天咱就聊聊这"结构化并发"到底怎么落地,让你在生产环境也能放心大胆地把并发拉到满。
一、先整明白:啥是"结构化并发"?
说白了就是"任务的生命周期得有个爹"。以前咱们写代码,用asyncio.create_task()创建的任务就像放飞的气球,飘哪儿去全看缘分。你要是没在finally里手动await,那异常可能就真·石沉大海了,程序"死"得不明不白。
结构化并发的核心就一句话:任务得抱团,生一起生,死一起死。这跟函数调用的栈帧一个理儿——你调用一个函数,它里面的子任务没跑完,父函数就别想退出去。TaskGroup就是这个"爹",它用async with圈出一块地,里面所有的create_task都是它的崽,出了这个门,要么全活着,要么全埋了,不存在"任务泄露"这一说。
对比之前的asyncio.gather(),TaskGroup牛逼在哪?Gather就像让几个快递员同时出发送外卖,有一个摔车了(异常),其他的还在傻乎乎继续跑,等你发现时超时都超了八回了。TaskGroup则是带了团队保险——一个出事,全组熔断,剩下的任务立马取消,省得浪费电(资源)。
二、Python 3.11+的新玩具:TaskGroup实战
先装个环境(假设你已经上了Python 3.11或更新的,没上的赶紧升,3.10及以下在这篇文章里就是"古代"):
python --version # 确认 >= 3.11
pip install aiohttp # 后面爬虫要用
来看个"Hello World"级别的对比。以前咱们这么写,心里直打鼓:
import asyncio
import random
async def fetch_data(name):
delay = random.uniform(0.5, 2.0)
await asyncio.sleep(delay)
if delay > 1.5:
raise ValueError(f"{name} 超时蚌埠住了")
print(f"{name} 搞定,耗时 {delay:.2f}s")
async def main_old():
# 上古写法,危险!
tasks = [
asyncio.create_task(fetch_data("任务A")),
asyncio.create_task(fetch_data("任务B")),
asyncio.create_task(fetch_data("任务C")),
]
try:
await asyncio.gather(*tasks)
except Exception as e:
print(f"抓到异常:{e}")
# 问题来了:另外两个任务还在跑吗?大概率是,资源泄露了兄弟
asyncio.run(main_old())
现在用TaskGroup,稳得一批:
async def main_modern():
try:
async with asyncio.TaskGroup() as tg:
tg.create_task(fetch_data("任务A"))
tg.create_task(fetch_data("任务B"))
tg.create_task(fetch_data("任务C"))
except* ValueError as eg:
# 注意这个except*,新语法,后面细讲
print(f"卧槽,有{len(eg.exceptions)}个任务炸锅了")
for e in eg.exceptions:
print(f" - {e}")
asyncio.run(main_modern())
看到async with没?这就是结构化并发的"结界"。一旦进了这个门,所有tg.create_task()创造的协程都被标记为"自己人"。只要有一个任务抛了没被处理的异常,TaskGroup立马启动取消连锁反应——给其他所有活着的任务发CancelledError,等它们都咽气了(或者处理完取消信号),才允许代码继续往下走。绝绝子,这安全感拉满。
三、ExceptionGroup与except*:一网打尽所有异常
以前最头疼的是啥?是并发任务里好几个同时报错,你只能抓到第一个,后面的异常全被吞了,日志里就留一行,排查问题跟猜灯谜似的。Python 3.11引入了ExceptionGroup,这玩意儿就像个异常收纳袋,把同时发生的多个错误打包一起扔给你。
配合except*语法(注意有个星号),你可以精准捕捞特定类型的异常,不管它们在不在一个包里:
import asyncio
import random
async def flaky_task(tag):
delay = random.uniform(0.1, 0.5)
await asyncio.sleep(delay)
if tag == "A" and delay > 0.3:
raise ValueError(f"{tag} 数值异常")
if tag == "B":
raise TypeError(f"{tag} 类型对不上")
if tag == "C":
raise ValueError(f"{tag} 又一个数值异常")
return f"{tag} 成功"
async def main():
try:
async with asyncio.TaskGroup() as tg:
tg.create_task(flaky_task("A"))
tg.create_task(flaky_task("B"))
tg.create_task(flaky_task("C"))
except* ValueError as eg:
# 这能抓住A和C的ValueError,哪怕它们是一起爆发的
print(f"数值错误收集器:共{len(eg.exceptions)}个")
for e in eg.exceptions:
print(f" 详情:{e}")
except* TypeError as eg:
# 抓住B的TypeError
print(f"类型错误:{eg.exceptions[0]}")
asyncio.run(main())
说白了,except ValueError只能抓"单个裸异常",而except* ValueError能拆开ExceptionGroup的包装,把里面所有匹配的类型都拎出来批处理。这在生产环境简直是救命神器——以前你可能只知道"有个任务错了",现在你能拿到完整的错误清单,写日志、发告警、重试策略都能做得明明白白。
四、生产级实战:千级并发API聚合器
光说不练假把式。假设你现在要搞个数据聚合服务:同时去10个不同的API抓数据,整合后返回,而且必须3秒内完成,谁超时谁滚蛋,还得限流防止把对方服务器打挂。这在微服务架构里太常见了。
直接上硬菜,这是能跑的生产代码骨架:
import asyncio
import aiohttp
import time
from typing import List, Dict
async def fetch_one(
session: aiohttp.ClientSession,
url: str,
semaphore: asyncio.Semaphore,
timeout: float
) -> Dict:
"""
单个请求:带限流 + 超时 + 错误处理
"""
async with semaphore: # 限制并发数,做个人,别DDOS别人
try:
# asyncio.timeout是3.11+的新API,比wait_for语义更清晰
async with asyncio.timeout(timeout):
async with session.get(url) as resp:
if resp.status == 200:
data = await resp.json()
return {"url": url, "status": "ok", "data": data}
else:
return {"url": url, "status": "error", "msg": f"HTTP {resp.status}"}
except asyncio.TimeoutError:
return {"url": url, "status": "timeout", "msg": f"超过{timeout}秒"}
except Exception as e:
return {"url": url, "status": "exception", "msg": str(e)}
async def batch_fetch(urls: List[str], max_concurrent: int = 50, timeout: float = 3.0) -> List[Dict]:
"""
批量抓取:结构化并发确保一个不落
"""
semaphore = asyncio.Semaphore(max_concurrent)
results = []
# 使用connector限制连接池大小,防止TCP端口耗尽
conn = aiohttp.TCPConnector(limit=100, limit_per_host=30)
async with aiohttp.ClientSession(connector=conn) as session:
try:
async with asyncio.TaskGroup() as tg:
# 技巧:用字典把task映射到url,方便后面取结果
task_to_url = {}
for url in urls:
task = tg.create_task(
fetch_one(session, url, semaphore, timeout)
)
task_to_url[task] = url
# TaskGroup退出时,所有任务已完成(或已取消/异常)
# 现在安全地收集结果
for task, url in task_to_url.items():
try:
results.append(task.result())
except Exception as e:
# 理论上fetch_one已经吞了异常,这里兜底
results.append({"url": url, "status": "fatal", "msg": str(e)})
except* asyncio.TimeoutError as eg:
# 虽然fetch_one里处理了超时,但如果TaskGroup层面熔断,会进这里
print(f"批量超时,{len(eg.exceptions)}个任务被强制取消")
except Exception as eg:
print(f"未预料的异常组:{eg}")
return results
# 模拟测试
test_urls = [f"https://httpbin.org/delay/{i%5}" for i in range(20)] # 模拟20个请求
async def main():
start = time.time()
results = await batch_fetch(test_urls, max_concurrent=5, timeout=2.0)
cost = time.time() - start
success = sum(1 for r in results if r["status"] == "ok")
print(f"\n耗时:{cost:.2f}s | 成功:{success}/{len(results)}")
# 打印失败的
for r in results:
if r["status"] != "ok":
print(f" 失败详情:{r['url']} -> {r['status']}: {r.get('msg', '')}")
if __name__ == "__main__":
asyncio.run(main())
这段代码里藏了几个保命技巧:
- Semaphore限流:就像火锅店门口放号的黄牛,一次只让50个号进去,防止后厨(对方服务器)炸了。网络IO可不是越多越好,TCP连接数、对方限流、本地端口(65535上限)都是坑。
- asyncio.timeout:3.11新API,替代了以前那个被诟病已久的
asyncio.wait_for。它就是个上下文管理器,语义特清楚——“这坨代码最多跑X秒”,超了直接抛TimeoutError。 - TaskGroup兜底:就算
fetch_one里写了try-except,万一哪天手滑忘写了,TaskGroup也能确保异常被捕获并触发其他任务取消,不至于一个线程泄露导致整个进程 hang 住。
五、进阶骚操作:嵌套任务与优雅退出
有时候你的任务不是平铺的,而是套娃——比如先并发抓10个页面,每个页面里再并发抓5个子链接。TaskGroup支持嵌套,而且父子之间互不干扰(当然异常会向上冒泡)。
async def process_sub_links(parent_url: str, links: List[str]):
"""处理子链接,又是一个TaskGroup"""
results = []
async with asyncio.TaskGroup() as tg:
for link in links:
tg.create_task(fetch_sub(link)) # 假设的抓取函数
return results
async def main_nested():
async with asyncio.TaskGroup() as tg:
for url in top_level_urls:
# 把子任务组作为大任务组的一员
tg.create_task(process_sub_links(url, get_links(url)))
优雅退出(Graceful Shutdown)也是生产环境必考项。收到SIGTERM信号时,你得让正在跑的任务保存状态、关闭连接,而不是直接kill -9。
import signal
import asyncio
async def cancellable_worker(name):
try:
while True:
print(f"{name} 正在摸鱼...")
await asyncio.sleep(1)
except asyncio.CancelledError:
print(f"{name} 收到取消信号,开始保存数据...")
# 这里做清理工作,比如刷盘、关数据库连接
await asyncio.sleep(0.5) # 模拟清理耗时
print(f"{name} 已安全退出")
raise # 必须重新抛出,让上层知道你真的退了
async def main():
loop = asyncio.get_running_loop()
# 创建任务组
async with asyncio.TaskGroup() as tg:
tasks = [
tg.create_task(cancellable_worker(f"Worker-{i}"))
for i in range(3)
]
# 模拟外部信号触发取消(比如Docker容器停止)
async def mock_signal():
await asyncio.sleep(3)
print("\n模拟收到关闭信号...")
# TaskGroup在退出块时会自动取消所有任务
# 这里我们手动触发退出(实际中可以用signal处理)
tg.create_task(mock_signal())
if __name__ == "__main__":
asyncio.run(main())
关键点在于:每个任务必须处理CancelledError,做清理,然后重新抛出异常。TaskGroup会等待所有子任务都处理完取消信号后,才会退出async with块。这比你以前手动task.cancel()然后await task要安全多了,毕竟手动管理容易漏。
六、这些坑,千万别踩
- TaskGroup不是万能药:它是为了"同生共死"的场景设计的。如果你希望任务独立——比如后台跑个日志上传,主请求返回了它还在后台慢慢传——那你还是得用
create_task+手动管理,或者上asyncio.Supervisor(3.13+实验性功能)。 - 别在TaskGroup里混同步代码:有些兄弟喜欢
time.sleep()往里面塞,直接阻塞整个事件循环,别的任务全卡死。必须用await asyncio.sleep()或者await asyncio.to_thread()把同步逻辑扔到线程池里。 - ExceptionGroup的粒度:如果你只关心"有没有错",用裸的
except ExceptionGroup as eg也能抓,但你就得手动遍历eg.exceptions。用except* SpecificError更精准,但记得3.11以下版本没有这语法,会直接语法错误。 - 内存爆炸:TaskGroup会持有所有任务的引用直到退出,如果你一口气塞进去10万个任务,内存可能直接爆掉。配合Semaphore做流控是必须的,或者分批处理。
七、总结:这波可以冲
说实话,Python的异步编程从3.4走到现在3.12/3.13,确实算"革命收官"了。以前写并发像走钢丝,现在有了TaskGroup+ExceptionGroup+asyncio.timeout这套组合拳,终于能像写同步代码那样有确定性了——你知道任务从哪开始、到哪结束,知道异常一定会被兜住,知道资源不会泄露。
对于搞爬虫、微服务、实时数据处理的兄弟,现在就是最佳上车时机。把gather和裸create_task的祖传代码 refactor 一波,用上结构化并发,你会发现bug少了,睡得香了,上线也不手抖了。这波真不亏,建议直接冲。
更多推荐

所有评论(0)