1. 问题背景:一个慢得像蜗牛的聚合API

上个月接手一个内部数据看板服务,核心接口/api/dashboard需要同时从用户服务、订单服务、库存服务和推荐服务拉取数据,然后组装返回。

原始代码是典型的串行写法:

# before.py - 串行版本
def get_dashboard(user_id: int):
    user = requests.get(f"http://user-svc/users/{user_id}", timeout=2).json()
    orders = requests.get(f"http://order-svc/orders?uid={user_id}", timeout=2).json()
    stock = requests.get(f"http://stock-svc/stock?uid={user_id}", timeout=2).json()
    recs = requests.get(f"http://rec-svc/recommend?uid={user_id}", timeout=2).json()
    return {"user": user, "orders": orders, "stock": stock, "recs": recs}

每个上游服务平均响应120ms,四个串行就是480ms。高峰期并发上来,Gunicorn worker线程被requests阻塞,CPU闲着呢但请求排队。监控数据显示:P95延迟612ms,QPS只有92,经常触发网关超时告警(我们网关阈值1s)。

2. 环境与版本:Python 3.10 + FastAPI + aiohttp

先说环境,避免版本坑:

  • Python 3.10.12(注意:3.10以下没有asyncio.run()的完整语义,3.8的loop.run_until_complete会有事件循环关闭警告)
  • FastAPI 0.104.1(用async def定义路由,直接支持asyncio)
  • aiohttp 3.9.1(推荐用ClientSession,不要自己搞连接池)
  • uvloop 0.19.0(可选项,生产环境建议装,能再压5-8%吞吐)
  • 部署:Gunicorn + uvicorn worker,4进程

3. 方案设计:asyncio.gather + Semaphore限流

核心思路很简单:把四个独立请求从串行改成并发。

但有几个设计决策要说明:

第一,用asyncio.gather还是asyncio.wait?
gather,因为要拿到所有结果统一返回。wait适合只关心完成状态不需要结果的场景。

第二,必须加信号量限流。
如果没有Semaphore,当上游服务抖动时,asyncio会瞬间发出成百上千个并发连接,把上游打挂。我用asyncio.Semaphore(20)限制同一时间最多20个出站请求。

第三,aiohttp连接池复用。
不要每次请求都创建ClientSession,那会重复建立TCP连接和TLS握手。应该在应用启动时创建全局session,设置connectorlimit参数。

4. 核心实现:从requests到aiohttp的重构

直接上重构后的代码:

# after.py - asyncio并发版本
import asyncio
import aiohttp
from fastapi import FastAPI

app = FastAPI()

# 全局session,应用启动时创建
session = None

@app.on_event("startup")
async def create_session():
    global session
    conn = aiohttp.TCPConnector(limit=50, ttl_dns_cache=300)
    session = aiohttp.ClientSession(connector=conn, timeout=aiohttp.ClientTimeout(total=2.5))

@app.on_event("shutdown")
async def close_session():
    await session.close()

# 信号量:限制并发出站请求数为20
sem = asyncio.Semaphore(20)

async def fetch_json(url: str):
    async with sem:
        async with session.get(url) as resp:
            return await resp.json()

@app.get("/api/dashboard")
async def get_dashboard(user_id: int):
    urls = [
        f"http://user-svc/users/{user_id}",
        f"http://order-svc/orders?uid={user_id}",
        f"http://stock-svc/stock?uid={user_id}",
        f"http://rec-svc/recommend?uid={user_id}"
    ]

    # gather并发跑四个请求,return_exceptions=True防止单个失败拖垮全部
    results = await asyncio.gather(
        *[fetch_json(url) for url in urls],
        return_exceptions=True
    )

    # 处理异常情况
    user, orders, stock, recs = results
    if isinstance(user, Exception):
        user = {"error": "user service unavailable"}
    # ... 其他类似处理

    return {"user": user, "orders": orders, "stock": stock, "recs": recs}

关键点说明:

  • TCPConnector(limit=50):连接池最大50个连接,超出会等待释放
  • ttl_dns_cache=300:DNS缓存5分钟,避免每次解析
  • ClientTimeout(total=2.5):总超时2.5秒,比requests的2秒略宽松,但比网关1s要短(实际上P95只有189ms,完全够用)
  • asyncio.Semaphore(20):20是压测得出的最佳值,低于10吞吐不足,高于30上游会偶尔超时

5. 踩坑与优化:三个真实教训

坑1:async def函数里用了阻塞调用

第一次重构时,我在fetch_json里用了json.loads(resp.text),这是CPU阻塞操作,会卡住事件循环。正确做法是直接用await resp.json(),它内部是异步解析JSON。

坑2:Gunicorn worker类型没换

我一开始用gunicorn -k sync跑FastAPI,async def路由根本不生效,全部退化到同步执行。必须用-k uvicorn.workers.UvicornWorker或者直接uvicorn app:app --workers 4。最终命令:

gunicorn app:app -k uvicorn.workers.UvicornWorker -w 4 -b 0.0.0.0:8000

坑3:asyncio.run()在FastAPI里的误用

有人会在async def路由里调用asyncio.run(fetch_json(url)),这会导致RuntimeError因为事件循环已存在。正确做法是用await直接调用,asyncio.run只在顶层入口用(比如独立脚本)。

6. 效果数据:从92 QPS到296 QPS

压测工具:wrk -t4 -c100 -d30s

指标 串行版本 asyncio版本 提升
平均延迟 482ms 172ms 2.8x
P95延迟 612ms 189ms 3.2x
QPS 92 296 3.2x
上游连接数峰值 1(串行) 20(信号量控制) -

补充说明:uvloop开启后,QPS从296涨到318,提升7.4%。另外把日志中的日志序列化从json.dumps换成orjson,延迟又降了5ms——这是细节优化,不展开。

生产环境运行两周,上游服务偶发抖动时,asyncio版本表现稳定:单个上游500ms超时,不影响其他三个请求结果,只是该字段返回错误占位符。而串行版本遇到超时,整个接口直接504。

7. 总结:什么时候该上asyncio

结论很明确:如果你的接口是IO密集型(HTTP调用、数据库查询、文件读写),且并发量要求高,asyncio是必选项,不是可选项。

但要说清楚,asyncio不是万能的:

  • CPU密集场景(图像处理、复杂计算)用多进程,别用asyncio
  • 简单接口只有一个上游调用,没必要上asyncio,线性代码更清晰
  • 团队不熟悉async/await语法,容易写出阻塞代码,需要code review把关

最后留个思考题:如果你的上游服务有依赖关系(比如B需要A的返回值),怎么用asyncio实现?我的方案是用asyncio.create_task + await分阶段处理,但注意不要嵌套过深导致代码难以维护。评论区聊。