你说这年头写个高并发程序,咋就跟开饭店似的?以前咱们用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())

这段代码里藏了几个保命技巧:

  1. Semaphore限流:就像火锅店门口放号的黄牛,一次只让50个号进去,防止后厨(对方服务器)炸了。网络IO可不是越多越好,TCP连接数、对方限流、本地端口(65535上限)都是坑。
  2. asyncio.timeout:3.11新API,替代了以前那个被诟病已久的asyncio.wait_for。它就是个上下文管理器,语义特清楚——“这坨代码最多跑X秒”,超了直接抛TimeoutError
  3. 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要安全多了,毕竟手动管理容易漏。

六、这些坑,千万别踩

  1. TaskGroup不是万能药:它是为了"同生共死"的场景设计的。如果你希望任务独立——比如后台跑个日志上传,主请求返回了它还在后台慢慢传——那你还是得用create_task+手动管理,或者上asyncio.Supervisor(3.13+实验性功能)。
  2. 别在TaskGroup里混同步代码:有些兄弟喜欢time.sleep()往里面塞,直接阻塞整个事件循环,别的任务全卡死。必须用await asyncio.sleep()或者await asyncio.to_thread()把同步逻辑扔到线程池里。
  3. ExceptionGroup的粒度:如果你只关心"有没有错",用裸的except ExceptionGroup as eg也能抓,但你就得手动遍历eg.exceptions。用except* SpecificError更精准,但记得3.11以下版本没有这语法,会直接语法错误。
  4. 内存爆炸:TaskGroup会持有所有任务的引用直到退出,如果你一口气塞进去10万个任务,内存可能直接爆掉。配合Semaphore做流控是必须的,或者分批处理。

七、总结:这波可以冲

说实话,Python的异步编程从3.4走到现在3.12/3.13,确实算"革命收官"了。以前写并发像走钢丝,现在有了TaskGroup+ExceptionGroup+asyncio.timeout这套组合拳,终于能像写同步代码那样有确定性了——你知道任务从哪开始、到哪结束,知道异常一定会被兜住,知道资源不会泄露。

对于搞爬虫、微服务、实时数据处理的兄弟,现在就是最佳上车时机。把gather和裸create_task的祖传代码 refactor 一波,用上结构化并发,你会发现bug少了,睡得香了,上线也不手抖了。这波真不亏,建议直接冲。

更多推荐