一、问题背景:一个慢得离谱的用户中心API

去年接手了一个内部用户中心服务,技术栈是 FastAPI + SQLAlchemy + Redis + 下游HTTP调用。接口逻辑不复杂,/user/profile/{uid} 大概做四件事:

  1. 查 Redis 缓存(miss 就查 MySQL)
  2. 查用户基础信息(MySQL)
  3. 调下游风控服务拿风险标签(HTTP)
  4. 调下游积分服务拿积分余额(HTTP)

单看代码没问题,但压测数据很难看。4核8G的容器,wrk -t4 -c200 -d30s 跑下来:

Requests/sec:    321.47
Latency  P50:    412ms
Latency  P99:    1.21s

CPU 利用率只有 18%,但 RT 高得离谱。原因很典型:接口是 async def 定义的,但里面全是同步阻塞调用。FastAPI 把它扔进线程池跑,线程池默认 40 个 worker,200 并发直接把线程池打满,请求排队。

更坑的是,SQLAlchemy 用的是同步 engine,requests 发 HTTP 请求,Redis 用的是同步的 redis-py。整个链路一个 await 都没有,异步框架被当同步框架用。

二、环境与版本

先把环境列清楚,避免版本差异带来的误导:

  • Python: 3.11.6
  • FastAPI: 0.109.0
  • uvicorn: 0.27.0(启动参数 --workers 4 --loop uvloop
  • SQLAlchemy: 2.0.25(异步用 asyncpg 驱动)
  • asyncpg: 0.29.0
  • aiohttp: 3.9.1
  • redis-py: 5.0.1(redis.asyncio
  • wrk: 4.2.0

数据库连接池配置:pool_size=20, max_overflow=10。下游 HTTP 连接池:aiohttp.TCPConnector(limit=200, limit_per_host=100)

三、方案设计:把每个阻塞点都换掉

改造思路就一句话:这条链路上所有 I/O 都必须是非阻塞的。逐个拆:

组件 改造前 改造后
MySQL SQLAlchemy sync + pymysql SQLAlchemy async + asyncpg
Redis redis-py sync redis.asyncio
HTTP requests aiohttp
并发编排 串行 await asyncio.gather 并发

关键设计点有两个:

第一,下游的两个 HTTP 调用必须并发。 风控和积分互不依赖,串行调用等于把两个 RTT 相加。用 asyncio.gather 并发后,这部分耗时直接砍半。

第二,注意缓存和 DB 的依赖关系。 Redis miss 才查 DB,这是串行依赖,不能并发。但可以在 cache miss 时用 asyncio.shield 保护回填任务,避免请求取消时缓存没写进去。

四、核心实现

先看改造前的代码(简化版):

# before: 看起来是 async,实际全是阻塞
@app.get("/user/profile/{uid}")
async def get_profile(uid: int):
    # 1. 查缓存(同步 redis)
    cache_key = f"user:profile:{uid}"
    cached = redis_client.get(cache_key)   # 阻塞!
    if cached:
        return json.loads(cached)

    # 2. 查 DB(同步 SQLAlchemy)
    user = db.query(User).filter(User.uid == uid).first()  # 阻塞!

    # 3. 串行调两个下游
    risk = requests.get(f"{RISK_URL}/risk/{uid}", timeout=2).json()   # 阻塞!
    points = requests.get(f"{POINTS_URL}/points/{uid}", timeout=2).json()  # 阻塞!

    result = {"uid": uid, "name": user.name, "risk": risk, "points": points}
    redis_client.setex(cache_key, 300, json.dumps(result))
    return result

改造后:

# after: 全链路异步 + 并发编排
import asyncio
import aiohttp
from redis.asyncio import Redis
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine, async_sessionmaker

engine = create_async_engine(
    "postgresql+asyncpg://user:pwd@db:5432/usercenter",
    pool_size=20, max_overflow=10, pool_pre_ping=True,
)
SessionLocal = async_sessionmaker(engine, expire_on_commit=False)
redis = Redis(host="redis", port=6379, decode_responses=True)

# 全局复用一个 session,避免每次建连接
http_session: aiohttp.ClientSession | None = None

@app.on_event("startup")
async def startup():
    global http_session
    connector = aiohttp.TCPConnector(limit=200, limit_per_host=100, ttl_dns_cache=300)
    timeout = aiohttp.ClientTimeout(total=2, connect=0.5)
    http_session = aiohttp.ClientSession(connector=connector, timeout=timeout)

@app.on_event("shutdown")
async def shutdown():
    await http_session.close()
    await redis.close()
    await engine.dispose()


async def fetch_json(url: str) -> dict:
    async with http_session.get(url) as resp:
        resp.raise_for_status()
        return await resp.json()


@app.get("/user/profile/{uid}")
async def get_profile(uid: int):
    cache_key = f"user:profile:{uid}"
    cached = await redis.get(cache_key)          # 非阻塞
    if cached:
        return json.loads(cached)

    # DB 查询和下游调用可以并发(下游不依赖 DB 结果)
    async def load_user():
        async with SessionLocal() as session:
            return await session.get(User, uid)

    user_task = asyncio.create_task(load_user())
    risk_task = asyncio.create_task(fetch_json(f"{RISK_URL}/risk/{uid}"))
    points_task = asyncio.create_task(fetch_json(f"{POINTS_URL}/points/{uid}"))

    user, risk, points = await asyncio.gather(user_task, risk_task, points_task)

    result = {"uid": uid, "name": user.name, "risk": risk, "points": points}

    # 回填缓存,用 shield 防止请求取消导致写不进去
    asyncio.create_task(asyncio.shield(
        redis.setex(cache_key, 300, json.dumps(result, ensure_ascii=False))
    ))
    return result

这里有几个细节值得单独说:

  • asyncio.create_task 把三个 I/O 同时启动,gather 等全部完成。总耗时 ≈ max(DB, 风控, 积分),而不是三者之和。
  • aiohttp 的 ClientSession 必须全局复用,每次请求新建 session 会导致连接池失效,性能比 requests 还差。
  • 缓存回填用后台任务,不阻塞响应返回。但要用 shield 包一层,否则请求被客户端取消时回填任务会被一起取消。

五、踩坑与优化

坑1:uvloop 没装,白等。 一开始 uvicorn 启动没加 --loop uvloop,默认用 asyncio 原生事件循环。装上 uvloop 后同样压测 QPS 从 1800 涨到 2400,提升 33%。这个参数一定要加。

坑2:连接池太小。 asyncpg 的 pool_size 默认 5,200 并发下大量请求在等连接。调到 20 后 P99 从 340ms 降到 180ms。经验值:pool_size ≈ worker数 * 5,但别超过 DB 的 max_connections

坑3:asyncio.gather 默认吞异常。 下游风控服务偶尔超时抛异常,gather 会把整个请求打挂。改成 return_exceptions=True,然后对失败的子任务做降级:

results = await asyncio.gather(user_task, risk_task, points_task, return_exceptions=True)
user, risk, points = results
if isinstance(risk, Exception):
    risk = {"tags": [], "degraded": True}

坑4:同步日志库拖后腿。 用了 logging 的 FileHandler,磁盘 I/O 在高并发下变成瓶颈。换成 QueueHandler + QueueListener 异步写日志后,QPS 又涨了约 8%。

坑5:time.sleep 混进来了。 代码里有一处重试逻辑用了 time.sleep(0.1),直接阻塞整个 event loop。改成 await asyncio.sleep(0.1)。这类问题很难发现,建议用 asyncio.debug=True 跑一遍,慢回调会打日志。

六、效果数据

同样的 4核8G 容器,同样的压测命令 wrk -t4 -c200 -d30s

指标 改造前 改造后 提升
QPS 321 2418 7.5x
P50 412ms 62ms 6.6x
P99 1210ms 183ms 6.6x
CPU 使用率 18% 71% -
平均内存 380MB 520MB +37%

内存涨了 140MB,主要是 aiohttp 连接池和 asyncpg 连接池的开销,这个代价可以接受。CPU 从 18% 涨到 71%,说明之前根本没用满,瓶颈全在阻塞等待上。

再补一组单接口拆解数据,看各段耗时:

阶段 改造前 改造后
Redis 查询 3ms 2ms
MySQL 查询 45ms 38ms
风控 HTTP 180ms 165ms
积分 HTTP 210ms 190ms
合计 438ms 205ms(并发后)

单请求耗时从 438ms 降到 205ms,主要是两个下游 HTTP 从串行变并行的功劳。

七、总结

这次改造最大的感受是:异步不是加几个 async/await 就完事,而是要保证整条链路没有阻塞点。一个 requests.get 就能让整个 event loop 卡住,前面的异步全白做。

几条可以直接抄的经验:

  1. asyncio 前先审计依赖库,requestspymysqlredis-py(同步版)都是雷。
  2. uvicorn 一定加 --loop uvloop,白捡 30% 性能。
  3. 连接池大小按 worker数 * 5 起步,压测调优。
  4. 独立的下游调用一律 gather 并发,注意 return_exceptions=True 做降级。
  5. 后台任务用 asyncio.shield 保护,避免被请求取消牵连。
  6. 上线前用 asyncio.debug=True 扫一遍慢回调和未 await 的协程。

异步编程的收益很高,但坑也很隐蔽。希望这篇记录能帮你少走点弯路。