第 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去包一层同步阻塞调用——那等于把整个事件循环交给一次数据库查询。
@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 秒,期间事件循环随便服务别人」。
❌ 不推荐:
@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)}✅ 推荐:
@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 库),正确做法是显式丢进线程池:
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):
@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:第一个异常立刻上抛,其他任务的结果全部丢弃。聚合接口通常不希望一个降级服务拖垮整个页面:
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:
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 NonePython 3.10(本教程最低版本)用 asyncio.wait_for,语义等价(只包单个 awaitable,抛 asyncio.TimeoutError):
return await asyncio.wait_for(_fetch(url), timeout=seconds)🔑 超时是必选项。给所有外部调用加超时,比加缓存、加机器都更有效。
13.8 用 httpx.AsyncClient 调用外部服务
最常见的性能坑:每个请求里 async with httpx.AsyncClient() 新建一次客户端,每次都重做 TCP 连接与 TLS 握手。正确做法是在 lifespan 里建单例客户端并复用连接池,退出时统一关闭:
@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()网络抖动导致的失败,用「指数退避 + 抖动 + 最大次数」重试即可,不必一上来就上熔断器:
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):
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 个独立进程,各自一份解释器和内存:
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 | 在事件循环中执行,函数体内任何阻塞调用都会拖垮全部并发 |
def | FastAPI 自动用 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 |
练习题
- 写
async def bad()与def good()两个端点,内部都调用time.sleep(1);用ab -n 20 -c 20压测,记录总耗时差异并解释原因。 - 把 13.6 的并发聚合改成「任一子服务超时 300 ms 就返回默认值」的版本,并在响应中标记
degraded。 - 给 13.8 的
request_with_retry增加retry_on_status: set[int]参数,并在 429 场景下读取Retry-After响应头决定等待时长。 - 写
/render端点用ProcessPoolExecutor执行一段纯 CPU 计算,同时并发请求/health,验证事件循环未被阻塞。
下一章预告
并发模型解决的是「同时服务更多请求」,但请求里除了 JSON,还可能是文件、表单和需要渲染的页面。下一章补齐
UploadFile、静态资源挂载与 Jinja2 模板。