一、问题背景:一个被同步IO拖死的聚合接口
去年接手了一个内部聚合接口 /api/v1/user/profile,逻辑不复杂:
- 根据user_id查MySQL拿基础信息
- 调用户中心HTTP接口拿标签
- 调风控HTTP接口拿风险等级
- 调订单服务HTTP接口拿最近订单数
- 合并返回
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不但没收益,还会因为事件循环调度带来额外开销。
几个实操建议:
- 全局ClientSession + 连接池,别每次请求新建
- 所有IO都加超时,
asyncio.wait_for是保命符 - 并发任务的异常要单独处理,别让一个失败拖垮整批
- 坚决不写同步IO,code review重点盯这个
- uvloop必装,白捡的性能
最后一句:换异步之前先想清楚瓶颈在哪。如果是DB慢、下游慢,先优化那些,asyncio只能帮你把"等待"重叠起来,不能让单个等待变快。