一、问题背景
去年接手了一个聚合查询服务,逻辑很简单:客户端传一个商品ID进来,服务端要同时去查三个下游接口——库存服务、价格服务、促销服务,把结果拼起来返回。
原来的实现是Flask + requests,典型写法:
# before: app.py
from flask import Flask, jsonify
import requests
app = Flask(__name__)
INVENTORY_URL = "http://inventory.internal/api/stock"
PRICE_URL = "http://price.internal/api/price"
PROMO_URL = "http://promo.internal/api/promo"
@app.route("/product/")
def get_product(pid):
stock = requests.get(INVENTORY_URL, params={"pid": pid}, timeout=2).json()
price = requests.get(PRICE_URL, params={"pid": pid}, timeout=2).json()
promo = requests.get(PROMO_URL, params={"pid": pid}, timeout=2).json()
return jsonify({
"pid": pid,
"stock": stock["count"],
"price": price["amount"],
"promo": promo.get("tag"),
})
这段代码的问题非常明显:三个下游请求是串行的。每个接口平均响应60~80ms,三个加起来就是200ms打底。下游稍微抖动一下,单个请求轻松上500ms。
线上表现:4核8G的机器,gunicorn开4个worker,每个worker 8个线程,单机QPS稳定在120左右,P99延迟1.2秒。大促前压测,8台机器才勉强扛住1000 QPS,成本高得离谱。
其实这三个请求彼此完全独立,没有依赖关系。串行执行纯属浪费。哪怕用线程池并发,也能把200ms压到80ms。但线程模型有GIL和上下文切换的开销,高并发下线程数会爆炸。真正合适的方案是asyncio——I/O密集型场景,协程才是正解。
二、环境与版本
- Python 3.11.6(3.11的asyncio在异常处理和task创建上比3.8快不少,官方benchmark显示TaskGroup场景有约20%提升)
- aiohttp 3.9.3
- FastAPI 0.109.2 + uvicorn 0.27.0(既然都上asyncio了,框架也一起换掉)
- gunicorn 21.2.0(只用它做进程管理,worker_class用uvicorn.workers.UvicornWorker)
- 压测工具:wrk 4.2.0,脚本模式,8线程200连接,持续60秒
如果你还在Python 3.8,建议至少升到3.10,asyncio.timeout()这个上下文管理器是3.11才有的,3.10之前只能手动wait_for,写起来啰嗦。
三、方案设计
核心思路其实就一句话:把串行的三次HTTP调用改成并发。
具体设计上考虑了几点:
- 框架选型:Flask是WSGI,天生同步,硬要在里面跑asyncio得用
asyncio.run()包一层,但每个请求都创建一个新event loop,开销巨大。直接换FastAPI,原生ASGI,event loop全局复用。 - HTTP客户端:requests是同步阻塞的,换成aiohttp。关键是ClientSession要复用,不能每个请求都新建,否则连接池形同虚设。
- 并发控制:三个请求用
asyncio.gather并发,但要加return_exceptions=True,避免一个下游挂了把整个请求带崩。 - 超时:整体超时2秒,单个下游超时1.5秒,超时的那个降级返回默认值,不影响其他两个。
- 连接池:ClientSession的
connector限制总连接数100,每个host限制20,避免打爆下游。
四、核心实现
先看重构后的核心代码:
# after: main.py
import asyncio
import aiohttp
from fastapi import FastAPI
from contextlib import asynccontextmanager
INVENTORY_URL = "http://inventory.internal/api/stock"
PRICE_URL = "http://price.internal/api/price"
PROMO_URL = "http://promo.internal/api/promo"
session: aiohttp.ClientSession | None = None
@asynccontextmanager
async def lifespan(app: FastAPI):
global session
connector = aiohttp.TCPConnector(
limit=100, # 总连接池上限
limit_per_host=20, # 单host上限
ttl_dns_cache=300, # DNS缓存5分钟
enable_cleanup_closed=True,
)
timeout = aiohttp.ClientTimeout(total=1.5, connect=0.5)
session = aiohttp.ClientSession(connector=connector, timeout=timeout)
yield
await session.close()
app = FastAPI(lifespan=lifespan)
async def fetch_json(url: str, pid: str) -> dict:
try:
async with session.get(url, params={"pid": pid}) as resp:
resp.raise_for_status()
return await resp.json()
except (aiohttp.ClientError, asyncio.TimeoutError) as e:
# 记录日志,返回降级值
return {"_error": str(e)}
@app.get("/product/{pid}")
async def get_product(pid: str):
async with asyncio.timeout(2.0): # 整体兜底超时,Python 3.11+
stock_task = asyncio.create_task(fetch_json(INVENTORY_URL, pid))
price_task = asyncio.create_task(fetch_json(PRICE_URL, pid))
promo_task = asyncio.create_task(fetch_json(PROMO_URL, pid))
stock, price, promo = await asyncio.gather(
stock_task, price_task, promo_task
)
return {
"pid": pid,
"stock": stock.get("count", 0) if "_error" not in stock else -1,
"price": price.get("amount", 0) if "_error" not in price else -1,
"promo": promo.get("tag") if "_error" not in promo else None,
}
启动命令:
gunicorn main:app \
-w 4 \
-k uvicorn.workers.UvicornWorker \
--bind 0.0.0.0:8000 \
--timeout 30 \
--max-requests 10000 \
--max-requests-jitter 1000
这里几个点值得展开说:
为什么用create_task而不是直接把协程扔给gather? 其实gather可以直接接收协程,效果一样。但显式create_task能让三个请求在进入gather之前就已经开始调度,语义更清楚,调试时也方便单独cancel某个task。
为什么timeout分两层? 单请求1.5秒是给下游的,整体2秒是给自己的。如果三个下游都超时,1.5秒就能返回;如果某个下游响应慢但没超时,整体2秒兜底,防止请求堆积。
ClientSession为什么放lifespan? 每个请求新建session会导致TCP连接无法复用,三次握手开销直接吃掉并发收益。实测新建session的版本QPS只有复用版本的三分之一。
五、踩坑与优化
坑1:event loop阻塞。 一开始我在异步函数里用了json.loads处理大响应,还有一处用了requests做埋点上报,结果QPS上不去。asyncio是单线程event loop,任何同步阻塞调用都会卡住整个loop。后来全部换成orjson和aiohttp,QPS直接翻倍。
坑2:连接池太小。 最初limit_per_host设的10,压测到800 QPS就上不去了,日志里全是Connection pool is full, discarding connection。改成20后,单机跑到1800 QPS稳定。
坑3:DNS解析。 内网服务用的是域名,aiohttp默认每次请求都走一次DNS(虽然是getaddrinfo,但在高并发下也是瓶颈)。加上ttl_dns_cache=300后,P99降了约30ms。
坑4:gunicorn worker数。 一开始按CPU核数配了4个worker,后来发现asyncio场景下worker不是越多越好。4核机器上,4个worker跑满CPU,但每个worker的event loop都有大量空闲。实测3个worker + 每个worker更高并发,整体吞吐更高。最终用4 worker是考虑到故障隔离。
坑5:asyncio.gather的异常传播。 默认return_exceptions=False,任何一个task抛异常,gather会立即抛出,但其他task不会被取消,会继续跑完。这会导致"请求已经返回了,后台还在发请求"的诡异现象。要么用return_exceptions=True,要么在except里手动cancel。我选了前者,配合fetch_json内部的try/except,双保险。
六、效果数据
压测条件:4核8G,8台机器压到2台,wrk 8线程200连接,60秒。
| 指标 | Before (Flask+requests) | After (FastAPI+aiohttp) | 提升 |
|---|---|---|---|
| 单机QPS | 120 | 1800 | 15x |
| P50延迟 | 210ms | 65ms | 3.2x |
| P99延迟 | 1200ms | 180ms | 6.7x |
| 单机CPU | 85% | 62% | - |
| 部署机器数 | 8 | 2 | -75% |
| 大促峰值承载 | 1000 QPS | 3600 QPS | 3.6x |
几个细节:
- P99从1.2秒到180ms,主要是消除了串行等待。理论上三个60ms的请求并发后应该是60ms,实测65ms,多出来的5ms是协程调度和结果聚合的开销,可以接受。
- CPU从85%降到62%,因为线程上下文切换没了。原来8个线程抢GIL,现在单event loop跑协程,切换成本极低。
- 机器从8台减到2台,一年省下来的云成本够发好几个月的年终奖(笑)。
一个反直觉的点:单机QPS到1800之后,继续加连接数收益递减。200连接时QPS 1800,400连接时QPS只到1950,P99反而涨到250ms。原因是下游三个服务本身有承载上限,我们的并发把压力传导过去了。所以最终生产环境我限了limit_per_host=15,宁可自己慢一点,也不能把下游打挂。
七、总结
这次重构的核心其实就一句话:I/O密集型场景,串行改并发,同步改异步。技术上没什么黑魔法,但落地时细节很多——event loop不能阻塞、连接池要复用、超时要分层、gather的异常语义要搞清楚,每一个踩错都会让性能打对折。
几点经验:
- 别在WSGI里硬塞asyncio,该换框架就换。FastAPI的学习成本一天就能补上,收益是十几倍QPS。
- ClientSession、连接池、DNS缓存是aiohttp性能的三驾马车,缺一不可。
- 超时一定要分层,单请求超时 + 整体超时,配合降级逻辑,下游抖动时你的服务才不会雪崩。
- 压测要压到瓶颈,找到拐点再回退10%,留出安全边际。QPS不是越高越好,稳定才是。
最后提醒一句:asyncio不是银弹。如果你的服务是CPU密集型(比如图像处理、大量计算),上asyncio只会更慢,该用多进程就用多进程。选对场景,才是真正的优化。