一、问题背景:一个被同步IO拖死的聚合接口

去年接手了一个内部聚合接口 /api/v1/user/profile,逻辑不复杂:

  1. 根据user_id查MySQL拿基础信息
  2. 调用户中心HTTP接口拿标签
  3. 调风控HTTP接口拿风险等级
  4. 调订单服务HTTP接口拿最近订单数
  5. 合并返回

Flask写的,同步requests + pymysql。单机4核8G,gunicorn 4 worker + 8 thread,QPS峰值也就120出头,P99 1.8s。业务方天天投诉超时。

我算过一笔账:下游三个HTTP接口RT分别是80ms、120ms、60ms,MySQL 15ms。同步串行下来光IO就275ms,加上框架开销和GC,单请求300ms起步。4个worker × 8线程 = 32并发,理论QPS上限也就 32 / 0.3 ≈ 106。跟实测的120基本吻合——瓶颈根本不在CPU,全堵在IO等待上。

上asyncio是必然选择。

二、环境与版本

  • Python 3.11.6(3.11的asyncio性能比3.8好不少,尤其是Task创建开销)
  • Flask 2.3.3 → FastAPI 0.104.1(顺手换了,原生async支持更干净)
  • uvicorn 0.24.0 + uvloop 0.19.0
  • aiohttp 3.9.1(HTTP客户端)
  • aiomysql 0.2.0(连接池)
  • 压测:wrk 4.2.0,wrk -t8 -c200 -d60s

uvloop这个别省,实测能再压10%-15%的延迟。

三、方案设计

核心思路就一句话:把串行IO变成并发IO,把线程等待变成事件循环调度

原来4个下游调用是串行的,其实它们之间没有依赖关系(标签、风控、订单数都只依赖user_id),完全可以并发。理论RT从275ms降到 max(120ms) + 合并开销 ≈ 130ms。

设计要点:

  • FastAPI的async def视图,全程不阻塞事件循环
  • asyncio.gather 并发三个HTTP调用 + 一个DB查询
  • aiohttp用全局 ClientSession,禁用per-request session(这个坑后面细说)
  • aiomysql连接池 minsize=5, maxsize=20
  • 每个下游调用加 asyncio.wait_for 超时保护,避免一个慢下游拖垮整体
  • 下游失败降级:标签/订单失败返回空,风控失败直接拒绝(安全要求)

四、核心实现

Before:同步Flask版本

# before_app.py
import requests
import pymysql
from flask import Flask, jsonify, request

app = Flask(__name__)

DB = pymysql.connect(host="10.0.0.10", user="app", password="xxx", database="user")

@app.route("/api/v1/user/profile")
def profile():
    uid = request.args.get("user_id")
    if not uid:
        return jsonify({"code": 400, "msg": "missing user_id"}), 400

    # 1. MySQL 串行
    with DB.cursor() as cur:
        cur.execute("SELECT nickname, level FROM user WHERE id=%s", (uid,))
        row = cur.fetchone()
        if not row:
            return jsonify({"code": 404, "msg": "user not found"}), 404
        base = {"nickname": row[0], "level": row[1]}

    # 2. 三个HTTP串行
    tags = requests.get(f"http://user-center/tags?uid={uid}", timeout=1).json()
    risk = requests.get(f"http://risk/level?uid={uid}", timeout=1).json()
    orders = requests.get(f"http://order/recent?uid={uid}", timeout=1).json()

    return jsonify({
        "code": 0,
        "data": {**base, "tags": tags, "risk": risk, "orders": orders}
    })

if __name__ == "__main__":
    app.run(host="0.0.0.0", port=8000)

启动命令:gunicorn -w 4 -k gthread --threads 8 -b 0.0.0.0:8000 before_app:app

After:asyncio版本

# after_app.py
import asyncio
import aiohttp
import aiomysql
from fastapi import FastAPI, Query, HTTPException
from contextlib import asynccontextmanager

app = FastAPI()

HTTP_TIMEOUT = 0.8          # 下游HTTP超时
DB_TIMEOUT = 0.3            # DB超时
GATHER_TIMEOUT = 1.2        # 整体兜底超时

session: aiohttp.ClientSession = None
pool: aiomysql.Pool = None


@asynccontextmanager
async def lifespan(app: FastAPI):
    global session, pool
    # 全局session,连接池参数按下游数量调优
    connector = aiohttp.TCPConnector(limit=200, limit_per_host=50, ttl_dns_cache=300)
    session = aiohttp.ClientSession(connector=connector, timeout=aiohttp.ClientTimeout(total=HTTP_TIMEOUT))
    pool = await aiomysql.create_pool(
        host="10.0.0.10", user="app", password="xxx", db="user",
        minsize=5, maxsize=20, pool_recycle=1800, autocommit=True
    )
    yield
    await session.close()
    pool.close()
    await pool.wait_closed()


app.router.lifespan_context = lifespan


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


async def get_base(uid: str):
    async with pool.acquire() as conn:
        async with conn.cursor() as cur:
            await cur.execute("SELECT nickname, level FROM user WHERE id=%s", (uid,))
            row = await cur.fetchone()
            return {"nickname": row[0], "level": row[1]} if row else None


async def safe_fetch(url: str, default):
    """单个下游失败不影响整体,返回默认值"""
    try:
        return await asyncio.wait_for(fetch_json(url), timeout=HTTP_TIMEOUT)
    except Exception:
        return default


@app.get("/api/v1/user/profile")
async def profile(user_id: str = Query(...)):
    # 风控必须成功,单独拿
    try:
        risk = await asyncio.wait_for(
            fetch_json(f"http://risk/level?uid={user_id}"), timeout=HTTP_TIMEOUT
        )
    except Exception:
        raise HTTPException(status_code=503, detail="risk service unavailable")

    # 其余三个并发
    base_task = asyncio.wait_for(get_base(user_id), timeout=DB_TIMEOUT)
    tags_task = safe_fetch(f"http://user-center/tags?uid={user_id}", {})
    orders_task = safe_fetch(f"http://order/recent?uid={user_id}", {"count": 0})

    try:
        base, tags, orders = await asyncio.wait_for(
            asyncio.gather(base_task, tags_task, orders_task),
            timeout=GATHER_TIMEOUT
        )
    except asyncio.TimeoutError:
        raise HTTPException(status_code=504, detail="upstream timeout")

    if not base:
        raise HTTPException(status_code=404, detail="user not found")

    return {"code": 0, "data": {**base, "tags": tags, "risk": risk, "orders": orders}}

启动:uvicorn after_app:app --host 0.0.0.0 --port 8000 --workers 4 --loop uvloop --http httptools

注意风控我故意串行放在前面——它失败要拒绝请求,没必要并发浪费下游资源。如果追求极致RT,也可以并发后判断,看业务取舍。

五、踩坑与优化

坑1:每次请求新建ClientSession。
第一版我在视图里 async with aiohttp.ClientSession() as s,QPS只有400。原因是每次建session都要新建TCP连接、TLS握手、DNS解析。改成全局session + TCPConnector 连接池后,直接翻倍到900+。

坑2:aiomysql连接池不够。
默认maxsize=10,压测到1500 QPS时大量请求卡在 pool.acquire()。改成20后缓解,但CPU开始飙。最后定位是DB侧max_connections不够,协调DBA调到500才彻底解决。

坑3:uvicorn workers和asyncio的关系。
uvicorn的 --workers 4 是4个进程,每个进程独立事件循环。别指望单进程多线程,asyncio本身就是单线程模型。4核机器就4个worker,多了反而上下文切换开销大。

坑4:gather里某个task抛异常会"吞"掉其他结果。
第一版风控超时时整个gather直接raise,标签和订单的结果全丢了。用 return_exceptions=True 或者在task内部就catch掉。

坑5:别在async函数里写同步代码。
有次同事在视图里加了一行 requests.get(...) 做埋点,QPS直接从2000掉到300。同步IO会阻塞整个事件循环,这是asyncio最致命的坑。埋点后来换成aiohttp异步上报。

六、效果数据

压测环境:4核8G,下游服务本地mock相同RT(80/120/60ms),wrk 200并发60秒。

指标 Before (Flask) After (FastAPI+asyncio) 提升
QPS 121 2137 17.6x
P50 420ms 38ms -91%
P95 1.1s 78ms -93%
P99 1.83s 95ms -95%
CPU峰值 65% 52% -
内存 380MB 210MB -45%
线上机器数 8台 2台 -75%

线上灰度一周,错误率从0.8%降到0.05%(主要是超时错误消失),P99稳定在100ms以内。

七、总结

asyncio不是什么银弹,它的收益完全取决于你的场景是不是IO密集。像这个接口,90%时间都在等下游,换成asyncio就是降维打击。但如果你的接口是CPU密集(比如图像处理、加密计算),上asyncio不但没收益,还会因为事件循环调度带来额外开销。

几个实操建议:

  1. 全局ClientSession + 连接池,别每次请求新建
  2. 所有IO都加超时asyncio.wait_for 是保命符
  3. 并发任务的异常要单独处理,别让一个失败拖垮整批
  4. 坚决不写同步IO,code review重点盯这个
  5. uvloop必装,白捡的性能

最后一句:换异步之前先想清楚瓶颈在哪。如果是DB慢、下游慢,先优化那些,asyncio只能帮你把"等待"重叠起来,不能让单个等待变快。