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

去年接手了一个商品详情接口,逻辑不复杂:

  1. 查 Redis 拿商品基础信息(缓存命中率约 92%)
  2. 缓存未命中时查 PostgreSQL
  3. 调两个下游服务:库存服务、促销服务
  4. 调推荐服务拿"猜你喜欢"
  5. 组装返回

原实现是 Flask + requests + psycopg2,典型的同步阻塞写法。上线后监控数据很难看:

  • 单进程 QPS:120 左右
  • P99 延迟:1200ms
  • 高峰期 CPU 利用率:只有 25%,但请求全堵在 IO 上
  • 部署了 8 台 4C8G 机器才扛住

问题很明显:一个请求要串行等 4~5 次网络 IO,每次 30~200ms,串起来就是几百毫秒。CPU 大部分时间在睡觉。

二、环境与版本

先说清楚版本,异步这块版本差异很大,抄代码前一定要对齐:

Python      3.11.6
Flask       2.3.3   (改造前)
requests    2.31.0  (改造前)
psycopg2    2.9.7   (改造前)
aiohttp     3.9.1   (改造后)
asyncpg     0.29.0  (改造后)
redis       5.0.1   (redis-py 的 asyncio 支持)
uvicorn     0.25.0
gunicorn    21.2.0  (用 UvicornWorker)

Python 3.11 的 asyncio 相比 3.8 有巨大的性能提升(Task 创建、Future 调度都优化过),如果还在 3.8/3.9,建议先升级。

三、方案设计:异步不是目的,并发才是

很多人对 asyncio 有个误解:以为把 def 改成 async def 就快了。不是的。 asyncio 的价值在于,当你要等 IO 时,把等待时间让给其他请求。

这个接口的调用链天然适合并发:

商品基础信息 ──┐
库存信息     ──┼──> 组装 ──> 推荐(依赖商品信息)──> 返回
促销信息     ──┘

库存、促销、基础信息之间没有依赖,可以并发;推荐依赖基础信息,所以放在后面。理论耗时从 t1+t2+t3+t4 降到 max(t1,t2,t3)+t4。

设计要点:

  1. Web 框架换成 FastAPI(基于 Starlette + asyncio)
  2. HTTP 客户端换 aiohttp,复用 ClientSession 连接池
  3. PostgreSQL 驱动换 asyncpg,注意它不走 SQLAlchemy 同步 session
  4. Redis 用 redis-py 的 asyncio 客户端 redis.asyncio
  5. 用 asyncio.gather 并发无依赖的调用,return_exceptions=True 做降级

四、核心实现

4.1 改造前:同步串行(Flask)

# before.py
import requests
import psycopg2
import redis
from flask import Flask, jsonify

app = Flask(__name__)
rds = redis.Redis(host="redis.internal", port=6379, decode_responses=True)
pg = psycopg2.connect("postgresql://app:pwd@pg.internal:5432/shop")
pg.autocommit = True

STOCK_URL = "http://stock.internal/api/stock/{}"
PROMO_URL = "http://promo.internal/api/promo/{}"
REC_URL = "http://rec.internal/api/recommend/{}"

@app.route("/product/")
def product_detail(pid):
    # 1. 缓存
    cached = rds.get(f"product:{pid}")
    if cached:
        import json
        base = json.loads(cached)
    else:
        with pg.cursor() as cur:
            cur.execute(
                "SELECT id, name, price, category_id FROM products WHERE id=%s",
                (pid,),
            )
            row = cur.fetchone()
            base = {"id": row[0], "name": row[1], "price": float(row[2]),
                    "category_id": row[3]}
            rds.setex(f"product:{pid}", 300, json.dumps(base))

    # 2. 串行调用下游
    stock = requests.get(STOCK_URL.format(pid), timeout=0.5).json()
    promo = requests.get(PROMO_URL.format(pid), timeout=0.5).json()
    rec   = requests.get(REC_URL.format(base["category_id"]), timeout=0.8).json()

    return jsonify({**base, "stock": stock, "promo": promo, "recommend": rec})

压测(wrk,4 线程 100 连接,跑 60s):

Requests/sec:   121.35
Latency P50:    780ms
Latency P99:    1210ms

4.2 改造后:asyncio 并发(FastAPI)

# after.py
import asyncio
import json
import asyncpg
import redis.asyncio as aioredis
import aiohttp
from fastapi import FastAPI, HTTPException
from contextlib import asynccontextmanager

STOCK_URL = "http://stock.internal/api/stock/{}"
PROMO_URL = "http://promo.internal/api/promo/{}"
REC_URL   = "http://rec.internal/api/recommend/{}"

# 全局资源,进程内复用
state = {}

@asynccontextmanager
async def lifespan(app: FastAPI):
    state["rd"] = aioredis.from_url(
        "redis://redis.internal:6379/0", decode_responses=True,
        max_connections=100,
    )
    state["pg"] = await asyncpg.create_pool(
        "postgresql://app:pwd@pg.internal:5432/shop",
        min_size=10, max_size=50, command_timeout=1.0,
    )
    # 关键:连接池 + 超时,避免下游慢打爆自己
    state["http"] = aiohttp.ClientSession(
        timeout=aiohttp.ClientTimeout(total=0.8, connect=0.2),
        connector=aiohttp.TCPConnector(
            limit=200, limit_per_host=100, ttl_dns_cache=60, keepalive_timeout=30,
        ),
    )
    yield
    await state["http"].close()
    await state["pg"].close()
    await state["rd"].close()

app = FastAPI(lifespan=lifespan)

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

async def get_base(pid: int):
    rd = state["rd"]
    key = f"product:{pid}"
    cached = await rd.get(key)
    if cached:
        return json.loads(cached)

    pg = state["pg"]
    async with pg.acquire() as conn:
        row = await conn.fetchrow(
            "SELECT id, name, price, category_id FROM products WHERE id=$1", pid
        )
    if not row:
        raise HTTPException(404, "product not found")
    base = {"id": row["id"], "name": row["name"],
            "price": float(row["price"]), "category_id": row["category_id"]}
    await rd.setex(key, 300, json.dumps(base))
    return base

@app.get("/product/{pid}")
async def product_detail(pid: int):
    base = await get_base(pid)
    session = state["http"]

    # 库存、促销、推荐并发拉取;推荐用 category_id,不依赖前两者
    stock_t = fetch_json(session, STOCK_URL.format(pid))
    promo_t = fetch_json(session, PROMO_URL.format(pid))
    rec_t   = fetch_json(session, REC_URL.format(base["category_id"]))

    results = await asyncio.gather(
        stock_t, promo_t, rec_t, return_exceptions=True
    )

    # 降级:任何一路失败都不影响主流程
    stock, promo, rec = (
        r if not isinstance(r, Exception) else None for r in results
    )
    return {**base, "stock": stock, "promo": promo, "recommend": rec}

启动:

gunicorn after:app \
  -k uvicorn.workers.UvicornWorker \
  -w 2 --threads 1 \
  -b 0.0.0.0:8000 \
  --timeout 30 --graceful-timeout 20

五、踩坑与优化

坑 1:以为 async def 就快。 第一版我把 requests.get 直接放进 async def,结果 QPS 反而降到 90。原因是 requests 是同步库,会阻塞整个事件循环,比多线程还惨。必须换成 aiohttp/httpx 的异步客户端。

坑 2:忘记复用 ClientSession。 一开始每次请求 async with aiohttp.ClientSession() as s:,QPS 只有 400。因为每次都要建 TCP 连接、TLS 握手。改成全局单例 + TCPConnector 连接池后,QPS 直接到 1400。

坑 3:下游抖动打爆事件循环。 促销服务有一次 P99 冲到 3s,导致所有请求都卡住。加了 ClientTimeout(total=0.8, connect=0.2) 和 return_exceptions=True 后,慢请求被超时切断,主流程走降级分支,接口稳定性大幅提升。

坑 4:asyncpg 和 SQLAlchemy 混合用。 项目里有些老代码用 SQLAlchemy 同步 Session,在 async 函数里直接调用会阻塞事件循环。后来统一走 asyncpg 原生 SQL,或者用 run_in_executor 包一层(但性能差很多,只适合低频接口)。

坑 5:Gunicorn worker 数量。 异步框架不需要开很多进程,IO 密集场景 2~4 个 UvicornWorker 就够。我一开始按同步经验开了 8 个 worker,结果每个 worker 的连接池都建满,PostgreSQL 连接被打爆。改成 -w 2 + asyncpg max_size=50 之后最稳。

坑 6:CPU 密集任务混进来。 有个 JSON 序列化 + 复杂计算用了 30ms CPU,在异步里会阻塞所有协程。这类任务拆到独立线程池:await asyncio.to_thread(compute, data)。

六、效果数据

同一份压测脚本(wrk -t4 -c100 -d60s),同一批机器(4C8G):

指标 改造前 (Flask) 改造后 (FastAPI+asyncio) 提升
单进程 QPS 121 1820 15.0x
P50 延迟 780ms 32ms 24.4x
P99 延迟 1210ms 85ms 14.2x
CPU 利用率 25% 68% —
内存占用 210MB 320MB +52%
部署机器数 8 台 2 台 -75%
下游错误率 0.8% 0.1%(含降级) —

内存涨了 100MB 左右,主要是 aiohttp 连接池和 asyncpg 连接池的开销,跟省下的机器成本比可以忽略。

线上灰度一周,核心接口平均耗时从 420ms 降到 62ms,超时告警从每天 30+ 条降到 0。

七、总结

asyncio 不是银弹,它的收益完全取决于你的场景:IO 密集、调用链有并发空间、下游稳定可控,收益巨大;如果你的接口是 CPU 密集,或者下游本身就慢到 1s+,异步改造意义有限。

几个关键结论:

  1. 异步客户端必须配套使用(aiohttp/httpx/asyncpg),混入同步库会毁掉整个事件循环
  2. 连接池、超时、降级是异步服务的三件套,缺一不可
  3. worker 数量不要按同步经验给,2~4 个 UvicornWorker 足够
  4. asyncio.gather(return_exceptions=True) 是最实用的降级工具
  5. Python 3.11+ 的 asyncio 性能比 3.8 好太多,升级收益明显

如果你的接口也是"一个请求串行等 N 个下游"的结构,别犹豫,改异步。改完之后你会发现,性能瓶颈终于变成了数据库,而不是你的代码。