一、问题背景:一个被同步IO拖垮的接口
去年接手了一个商品详情接口,逻辑不复杂:
- 查 Redis 拿商品基础信息(缓存命中率约 92%)
- 缓存未命中时查 PostgreSQL
- 调两个下游服务:库存服务、促销服务
- 调推荐服务拿"猜你喜欢"
- 组装返回
原实现是 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。
设计要点:
- Web 框架换成 FastAPI(基于 Starlette + asyncio)
- HTTP 客户端换 aiohttp,复用
ClientSession连接池 - PostgreSQL 驱动换 asyncpg,注意它不走 SQLAlchemy 同步 session
- Redis 用 redis-py 的 asyncio 客户端
redis.asyncio - 用
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+,异步改造意义有限。
几个关键结论:
- 异步客户端必须配套使用(aiohttp/httpx/asyncpg),混入同步库会毁掉整个事件循环
- 连接池、超时、降级是异步服务的三件套,缺一不可
- worker 数量不要按同步经验给,2~4 个 UvicornWorker 足够
asyncio.gather(return_exceptions=True)是最实用的降级工具- Python 3.11+ 的 asyncio 性能比 3.8 好太多,升级收益明显
如果你的接口也是"一个请求串行等 N 个下游"的结构,别犹豫,改异步。改完之后你会发现,性能瓶颈终于变成了数据库,而不是你的代码。