Skip to content

第 13 章:异步编程与并发模型 ​

第 10 章用了异步 SQLAlchemy,第 12 章在依赖里查了用户。但「这个路径操作函数到底该写 async def 还是 def」,决定了你的服务能扛住上千并发,还是被一个慢请求拖死。


学习目标 ​

  • 理解「async def 不会让单个请求变快」,明确异步解决的是 I/O 等待的并发问题
  • 掌握 FastAPI 对 async def 与 def 的不同调度方式:事件循环 vs 线程池
  • 能识别并修复「在 async def 里做阻塞调用」这类生产事故级写法
  • 会用 asyncio.gather 并发聚合并配合超时、退避重试
  • 掌握 httpx.AsyncClient 单例 + lifespan 复用连接池的封装,理解 CPU 密集任务的边界

13.1 先破除误解:async def 不等于更快 ​

  • 协程(Coroutine)是协作式并发,不是并行,不能把 CPU 计算拆到多个核心。
  • async def 唯一的收益是:代码在 await 处等待 I/O 时,事件循环把 CPU 让给其他请求。
  • 若函数体内没有 await,或 await 的其实是同步阻塞调用,事件循环在此期间完全停摆,吞吐比 def + 线程池还差。

13.2 ASGI 事件循环模型 ​

uvicorn app.main:app 默认启动 1 个进程、1 个事件循环、1 个线程,所有 async def 路径操作函数都在这一个线程上排队。谁阻塞了它,谁就阻塞了全部并发。


13.3 关键规则:async def 与 def ​

函数定义执行位置阻塞后果该配什么库
async def事件循环线程内直接 await一个阻塞调用拖垮全部并发异步库:AsyncSession / httpx.AsyncClient
def自动丢进线程池(run_in_threadpool)只占一个线程,其他请求不受影响同步库:requests / 同步 Session

线程池由 AnyIO 管理,默认上限 40 个线程;40 个同步任务在跑时,第 41 个请求只能排队。

🔑 用什么库,就写什么签名。 同步库 → def;异步库 → async def。绝不要用 async def 去包一层同步阻塞调用——那等于把整个事件循环交给一次数据库查询。

python
@app.get("/users/")
def list_users(db: Session = Depends(get_sync_db)) -> list[dict[str, object]]:
    rows = db.execute(select(User)).scalars().all()      # 同步库 → 线程池
    return [{"id": u.id, "name": u.name} for u in rows]


@app.get("/users/async")
async def list_users_async(db: AsyncSession = Depends(get_async_db)) -> list[dict[str, object]]:
    result = await db.execute(select(User))              # 异步库 → 事件循环
    return [{"id": u.id, "name": u.name} for u in result.scalars().all()]

💡 两种写法可在同一应用共存,但同一条调用链上不要混:async def 里别依赖同步 Session,同步上下文里也没法 await 异步 Session。


13.4 阻塞事件循环的灾难现场 ​

time.sleep(3) 只把线程挂起,不交出控制权;asyncio.sleep(3) 才是「睡 3 秒,期间事件循环随便服务别人」。

❌ 不推荐:

python
@app.get("/slow")
async def slow() -> dict[str, str]:
    time.sleep(3)                                 # ❌ 阻塞事件循环 3 秒
    r = requests.get("https://httpbin.org/get")   # ❌ 同步 HTTP,再阻塞一次
    return {"status": str(r.status_code)}

✅ 推荐:

python
@app.get("/fast")
async def fast() -> dict[str, str]:
    await asyncio.sleep(3)                              # ✅ 让出控制权
    async with httpx.AsyncClient() as client:
        r = await client.get("https://httpbin.org/get")  # ✅ 非阻塞
    return {"status": str(r.status_code)}

单请求耗时几乎一样,差别在并发下(单进程单 worker,10 个并发请求):

场景总耗时说明
async def + time.sleep(3)约 30 s串行,事件循环被完全占据
def + time.sleep(3)约 3 s线程池并行,40 线程够用
async def + asyncio.sleep(3)约 3 s事件循环持续服务其他请求

这个差异在开发环境用 curl 单发请求时完全看不出来,上线后表现为「QPS 上不去、CPU 却很闲」,排查成本极高。


13.5 必须调用同步阻塞函数时 ​

有些库没有异步版本(老 SDK、boto3、部分图像/PDF 库),正确做法是显式丢进线程池:

python
from fastapi.concurrency import run_in_threadpool


@app.post("/legacy/")
async def call_legacy(payload: str) -> dict[str, str]:
    # ✅ legacy_client.process 是同步阻塞的,丢进线程池执行
    result = await run_in_threadpool(legacy_client.process, payload)
    return {"result": result}

run_in_threadpool 是 anyio.to_thread.run_sync 的薄封装,两者等价。注意线程池容量只有 40,别把重 CPU 任务也塞进来,那会占满线程池、让同进程内所有 def 路由一起排队。


13.6 并发聚合:asyncio.gather ​

一个接口要同时调 3 个下游服务,串行 await 是最常见的性能浪费:三个各 200 / 300 / 150 ms 的调用串起来要 650 ms,gather 并发后只取决于最慢的那个。

❌ 不推荐:profile = await ...,再 orders = await ...,再 coupons = await ...,三次等待首尾相接,总耗时 650 ms。

✅ 推荐(并发,总耗时取决于最慢的那个,约 300 ms):

python
@app.get("/dashboard/parallel")
async def dashboard_parallel(user_id: int) -> dict[str, object]:
    profile, orders, coupons = await asyncio.gather(
        user_service.get_profile(user_id),
        order_service.list_orders(user_id),
        coupon_service.list_coupons(user_id),
    )
    return {"profile": profile, "orders": orders, "coupons": coupons}

gather 默认 return_exceptions=False:第一个异常立刻上抛,其他任务的结果全部丢弃。聚合接口通常不希望一个降级服务拖垮整个页面:

python
results = await asyncio.gather(
    user_service.get_profile(user_id),
    order_service.list_orders(user_id),
    return_exceptions=True,   # ✅ 异常作为结果返回,不中断其他任务
)
profile, orders = results
# ⚠️ 必须逐个判类型,否则异常对象会被当成正常数据塞进 response_model
return {"profile": None if isinstance(profile, BaseException) else profile,
        "orders": [] if isinstance(orders, BaseException) else orders}

13.7 超时控制:别让下游拖死你的接口 ​

下游不设超时是另一个高频事故:对方 hang 住,你的连接和协程全部堆积。Python 3.11+ 用 asyncio.timeout:

python
async def fetch_with_timeout(url: str, seconds: float = 2.0) -> str:
    try:
        async with asyncio.timeout(seconds):       # 上下文管理器,可包住多步逻辑
            async with httpx.AsyncClient() as client:
                return (await client.get(url)).text
    except TimeoutError:                           # 3.11+ 抛内置 TimeoutError
        raise HTTPException(status_code=504, detail=f"下游超时: {url}") from None

Python 3.10(本教程最低版本)用 asyncio.wait_for,语义等价(只包单个 awaitable,抛 asyncio.TimeoutError):

python
return await asyncio.wait_for(_fetch(url), timeout=seconds)

🔑 超时是必选项。给所有外部调用加超时,比加缓存、加机器都更有效。


13.8 用 httpx.AsyncClient 调用外部服务 ​

最常见的性能坑:每个请求里 async with httpx.AsyncClient() 新建一次客户端,每次都重做 TCP 连接与 TLS 握手。正确做法是在 lifespan 里建单例客户端并复用连接池,退出时统一关闭:

python
@asynccontextmanager
async def lifespan(app: FastAPI):
    # ✅ 应用启动时建一次,复用连接池
    app.state.http_client = httpx.AsyncClient(
        base_url="https://api.example.com",
        timeout=httpx.Timeout(connect=2.0, read=5.0, write=5.0, pool=2.0),
        limits=httpx.Limits(max_connections=100, max_keepalive_connections=20),
    )
    try:
        yield
    finally:
        await app.state.http_client.aclose()   # ✅ 优雅关闭,释放连接


app = FastAPI(lifespan=lifespan)


def get_http_client(request: Request) -> httpx.AsyncClient:
    return request.app.state.http_client       # 从 app.state 取单例


HttpClientDep = Annotated[httpx.AsyncClient, Depends(get_http_client)]


@app.get("/remote/items")
async def remote_items(client: HttpClientDep) -> list[dict[str, object]]:
    return (await client.get("/items", params={"limit": 20})).json()

网络抖动导致的失败,用「指数退避 + 抖动 + 最大次数」重试即可,不必一上来就上熔断器:

python
async def request_with_retry(
    client: httpx.AsyncClient, url: str, *, max_attempts: int = 3, base_delay: float = 0.2
) -> httpx.Response:
    for attempt in range(max_attempts):
        try:
            response = await client.get(url)
            if response.status_code < 500 and response.status_code != 429:
                return response                  # 4xx 是客户端错误,不重试
        except (httpx.TimeoutException, httpx.TransportError):
            pass
        if attempt < max_attempts - 1:
            # ✅ 退避 + 抖动,避免失败请求同时重试把下游打成雪崩
            await asyncio.sleep(base_delay * (2**attempt) + random.uniform(0, 0.1))
    raise HTTPException(status_code=502, detail=f"下游重试 {max_attempts} 次仍失败")

13.9 CPU 密集任务的处理 ​

图像缩放、PDF 生成、大文件加解密都没有 await 点,放在事件循环里就是纯阻塞。偶发的秒级计算丢给进程池,长任务则交给外部任务队列(Celery / arq / RQ):

python
process_pool = ProcessPoolExecutor(max_workers=4)


@app.post("/render/")
async def render(data: bytes) -> dict[str, int]:
    loop = asyncio.get_running_loop()
    # ✅ 丢进进程池,绕过 GIL,事件循环不受影响
    pdf = await loop.run_in_executor(process_pool, render_pdf, data)
    return {"size": len(pdf)}

⚠️ 不要用线程池跑 CPU 密集任务:GIL 决定线程无法并行计算,只会把 40 个线程占满;进程池的参数与返回值还必须可 pickle。


13.10 多 worker 下的共享状态 ​

uvicorn --workers 4 启动的是 4 个独立进程,各自一份解释器和内存:

python
CACHE: dict[str, str] = {}          # ❌ 模块级全局变量只在单个进程内有效


@app.get("/stats")
async def stats(redis: RedisDep) -> dict[str, int]:
    # ✅ 换 Redis 后,所有 worker 看到同一份数据
    return {"requests": int(await redis.incr("app:requests"))}

app.state 里的对象(含 13.8 的 httpx.AsyncClient)同样是每进程一份——连接池本就该按进程隔离;但拿它做业务缓存或计数器就是错的。多进程部署与 worker 数量详见第 17 章。


常见坑与排查 ​

现象原因解决
QPS 上不去,CPU 使用率却很低async def 里调了 time.sleep / requests / 同步 DB 驱动换成 asyncio.sleep / httpx.AsyncClient / 异步驱动;同步库改回 def
压测时前几个请求正常,之后全部超时事件循环被阻塞,请求在队列里堆积同上;用 py-spy dump 定位阻塞点
gather 中一个任务失败,其他结果全丢默认 return_exceptions=False加 return_exceptions=True 并逐个判 isinstance(r, BaseException)
异步函数里用 httpx.Client(同步版)同步客户端内部使用阻塞 socket改用 httpx.AsyncClient,在 lifespan 中做单例
asyncio.run() cannot be called from a running event loop在协程内部又调了 asyncio.run()协程里直接 await;asyncio.run() 只能出现在同步入口
大量 def 路由同时变慢线程池 40 个线程被占满,请求排队减少同步路由或改异步;必要时调大 AnyIO 线程上限

本章小结 ​

要点说明
async def在事件循环中执行,函数体内任何阻塞调用都会拖垮全部并发
defFastAPI 自动用 run_in_threadpool 调度,线程池默认 40 线程
选型规则同步库写 def,异步库写 async def,不要用 async def 包同步调用
应急手段run_in_threadpool / anyio.to_thread.run_sync 把同步函数移出事件循环
并发聚合asyncio.gather,异常场景加 return_exceptions=True
超时3.11+ 用 asyncio.timeout,3.10 用 asyncio.wait_for,外部调用必须设超时
HTTP 客户端httpx.AsyncClient 在 lifespan 中建单例,退出时 aclose()
CPU 密集与共享状态进程池或外部任务队列;多 worker 内存不共享,计数与缓存放 Redis

练习题 ​

  1. 写 async def bad() 与 def good() 两个端点,内部都调用 time.sleep(1);用 ab -n 20 -c 20 压测,记录总耗时差异并解释原因。
  2. 把 13.6 的并发聚合改成「任一子服务超时 300 ms 就返回默认值」的版本,并在响应中标记 degraded。
  3. 给 13.8 的 request_with_retry 增加 retry_on_status: set[int] 参数,并在 429 场景下读取 Retry-After 响应头决定等待时长。
  4. 写 /render 端点用 ProcessPoolExecutor 执行一段纯 CPU 计算,同时并发请求 /health,验证事件循环未被阻塞。

下一章预告 ​

并发模型解决的是「同时服务更多请求」,但请求里除了 JSON,还可能是文件、表单和需要渲染的页面。下一章补齐 UploadFile、静态资源挂载与 Jinja2 模板。

👉 第 14 章:文件上传、静态资源与模板

📖本文阅读--次|📊全站访问--次|👥访客--人