1. 问题背景:一次“正常”的发布引发的雪崩
事情发生在上周三。我们有个BFF层服务,用Flask 2.2写的,部署在8核16G的容器里,Gunicorn配了4个worker。某天发布了一个新功能——首页聚合接口/api/home,需要并行调用用户服务、商品服务、库存服务、推荐服务和营销服务,取数据后拼装返回。
代码很简单,就是requests.get依次调用5个下游:
def get_home_data(user_id):
user = requests.get(f"http://user-service/users/{user_id}", timeout=3).json()
products = requests.get(f"http://product-service/products?ids={user['fav_ids']}", timeout=3).json()
stock = requests.get(f"http://inventory-service/stock?ids={[p['id'] for p in products]}", timeout=3).json()
...
return assemble(user, products, stock, rec, mkt)
发布后监控立刻报警:接口P99从300ms飙到3.8s,QPS从180掉到70,Gunicorn的worker CPU 100%,线程池队列积压。原因很明显——串行调用,每个下游耗时300-500ms,加上网络抖动,一个请求要等2s+。4个worker × 默认线程池(20) = 80并发上限,下游一抖动,线程全被占住,新请求排队。
当时有两个方案:改多线程,或者改异步。多线程能解决阻塞问题,但线程切换开销和GIL限制,提升有限。我选了asyncio——因为下游是纯IO密集型,协程切换几乎零成本。
2. 环境与版本:Python 3.10 + Flask 2.2 + httpx 0.24
先说清楚环境,版本很重要,因为asyncio的API在3.10和3.11有差异,尤其asyncio.run()和loop.run_until_complete()的坑。
- Python 3.10.12
- Flask 2.2.3
- Gunicorn 20.1.0 (4 workers, sync worker class)
- httpx 0.24.1(支持异步的HTTP客户端,比aiohttp轻量)
- asyncio 内置,版本跟随Python
重要:Flask是同步框架,不能直接在view函数里await。需要把异步代码包在asyncio.run()里,或者用asyncio.get_event_loop().run_until_complete()。但注意,asyncio.run()每次调用会创建新的事件循环,如果有全局连接池,会失效。所以更优做法是用asyncio.new_event_loop() + run_until_complete,复用同一个loop。
3. 方案设计:协程并发 + 信号量限流 + 超时控制
改造方案核心三点:
- 并发化:把5个独立的下游调用从串行改为并发协程,用
asyncio.gather()同时发起。 - 限流:防止下游被打爆,用
asyncio.Semaphore(10)控制最大并发协程数。 - 超时:每个协程用
asyncio.wait_for()包裹,设置2.5s超时,防止单个下游拖垮整个请求。
架构上,Flask view函数保持同步,内部调用一个异步入口函数:
def get_home_data(user_id):
return asyncio.get_event_loop().run_until_complete(
_async_get_home_data(user_id)
)
注意:asyncio.get_event_loop()在Python 3.10里如果在没有运行loop的线程调用,会创建一个新loop。但在Gunicorn worker里,每次请求都会调用,所以我们要在worker启动时创建全局loop,避免重复创建消耗资源。这里我用了一个模块级变量:
# async_util.py
import asyncio
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
def run_async(coro):
return loop.run_until_complete(coro)
4. 核心实现:before/after代码对照
Before:同步串行(性能瓶颈)
# before.py
import requests
from flask import Flask, jsonify, request
app = Flask(__name__)
@app.route('/api/home')
def home():
user_id = request.args.get('user_id')
# 串行调用5个下游,每个耗时300-500ms
user = requests.get(f"http://user-service/users/{user_id}", timeout=3).json()
product_ids = user['fav_ids'][:10]
products = requests.get(f"http://product-service/products?ids={product_ids}", timeout=3).json()
stock = requests.get(f"http://inventory-service/stock?ids={[p['id'] for p in products]}", timeout=3).json()
rec = requests.get(f"http://rec-service/recommend?user_id={user_id}", timeout=3).json()
mkt = requests.get(f"http://mkt-service/coupons?user_id={user_id}", timeout=3).json()
return jsonify(assemble(user, products, stock, rec, mkt))
After:异步并发改造
# after.py
import asyncio
import httpx
from flask import Flask, jsonify, request
from async_util import run_async
app = Flask(__name__)
# 全局httpx异步客户端,复用连接池
client = httpx.AsyncClient(timeout=httpx.Timeout(3.0), limits=httpx.Limits(max_connections=100))
# 限流信号量:最多10个并发下游请求
semaphore = asyncio.Semaphore(10)
async def fetch_json(url, params=None):
async with semaphore:
try:
# asyncio.wait_for 强制超时
resp = await asyncio.wait_for(client.get(url, params=params), timeout=2.5)
return resp.json()
except (httpx.TimeoutException, asyncio.TimeoutError):
return None # 超时返回None,由调用方降级
except Exception as e:
print(f"Error fetching {url}: {e}")
return None
async def _get_home_data(user_id):
# 并发发起5个请求
user_task = fetch_json(f"http://user-service/users/{user_id}")
rec_task = fetch_json(f"http://rec-service/recommend", params={"user_id": user_id})
mkt_task = fetch_json(f"http://mkt-service/coupons", params={"user_id": user_id})
user, rec, mkt = await asyncio.gather(user_task, rec_task, mkt_task)
if not user:
return {"error": "user service unavailable"}, 503
# 第二个依赖第一个的结果,但产品列表和库存也可以并行
product_ids = user['fav_ids'][:10]
products_task = fetch_json(f"http://product-service/products", params={"ids": product_ids})
# 库存依赖产品ID,但这里我们先用占位,实际可再拆
stock_task = fetch_json(f"http://inventory-service/stock", params={"ids": product_ids})
products, stock = await asyncio.gather(products_task, stock_task)
return assemble(user, products or [], stock or [], rec or [], mkt or [])
@app.route('/api/home')
def home():
user_id = request.args.get('user_id')
data, status = run_async(_get_home_data(user_id))
return jsonify(data), status
if __name__ == '__main__':
app.run(threaded=False) # 注意:不能用threaded=True,否则loop冲突
关键点:
- httpx.AsyncClient是全局单例,复用TCP连接,避免每次握手。
- asyncio.gather()并发执行,5次调用从串行2s+降到并发400ms左右。
- 信号量Semaphore(10)防止下游被突发流量打死,实测峰值时下游收到120并发请求,加限流后稳定在10。
- 超时用asyncio.wait_for包裹,且内部捕获asyncio.TimeoutError,返回None进行降级——比直接抛异常好,保证主流程可用。
5. 踩坑与优化:三个真实教训
坑1:asyncio.run() 每次创建新loop,连接池失效
一开始我用asyncio.run(_get_home_data(user_id)),结果每次请求都新建事件循环,httpx.AsyncClient的连接池每次都被丢弃,握手开销巨大,性能反而没提升多少。改成run_until_complete复用全局loop后,连接复用率上去了,性能才真正起飞。
坑2:Gunicorn sync worker + threaded=True 导致loop冲突
Flask的dev server默认threaded=True,每个请求一个线程。如果多个线程同时调用同一个loop的run_until_complete,会报RuntimeError: This event loop is already running。解决办法:生产环境用sync worker(单线程),或者用asyncio.locks.Lock保护loop调用。我直接改成threaded=False,因为Gunicorn本身是多进程。
坑3:超时设置过短导致毛刺
最初超时设1.5s,下游正常时没问题,但偶尔网络抖动,超过1.5s就直接降级返回None,导致前端看到部分数据缺失。调优后设为2.5s,且对关键数据(用户信息)做重试(最多1次),对非关键数据(营销)直接降级。这样P99和可用性都得到保证。
优化:按依赖关系分层并发
最开始的代码是5个请求全部一起gather,但库存接口需要商品ID,所以实际上库存要等商品返回。我拆成两层:第一层并发用户+推荐+营销,第二层并发商品+库存。效果:总耗时从第一版并发所有(约400ms)降到310ms,因为依赖链路上减少了一个串行等待。
6. 效果数据:压测对比
用wrk压测,环境:8核16G容器,Gunicorn 4 workers,单接口/api/home,模拟下游服务延迟300ms(用Mock服务)。
| 指标 | Before (串行) | After (异步并发) | 提升 |
|---|---|---|---|
| QPS | 120 | 410 | +241% |
| P50延迟 | 850ms | 310ms | -63.5% |
| P99延迟 | 3800ms | 680ms | -82% |
| 线程池占用 | 100% (20线程全占) | 2% (协程闲置状态) | 几乎零占用 |
| CPU使用率 | 75% | 55% | 下降(因为减少了线程切换) |
压测命令:
wrk -t8 -c200 -d60s http://localhost:8000/api/home?user_id=123
注意:QPS提升来源于并发而非CPU优化。8核CPU下,同步阻塞时线程切换开销大,异步协程切换开销仅约1µs,且Gunicorn worker数可以不变,但吞吐翻倍。
7. 总结:异步化适用的边界
这次改造收益明显,但并非所有场景都适合asyncio。总结三点经验供参考:
-
适用场景:纯IO密集型、下游API多、依赖关系可分层并发。CPU密集型(如加解密、图像处理)不适合,应改用多进程。
-
框架选择:Flask做异步需要自己套loop,比较别扭。新项目建议直接用FastAPI或Sanic,原生支持async。但存量Flask项目用上述模式改造,成本很低。
-
监控与降级:异步代码里异常容易被吞,务必在协程内捕获并记录日志。我加了
try/except和超时返回None的降级逻辑,保证主流程不挂。
最后,代码已上线一个月,线上P99稳定在650ms左右,QPS峰值500+,再也没收到过下游抖动引发的雪崩告警。如果你也遇到类似的BFF层阻塞问题,可以试试这个改造路径——成本不高,收益巨大。有问题欢迎评论区交流。