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. 方案设计:协程并发 + 信号量限流 + 超时控制

改造方案核心三点:

  1. 并发化:把5个独立的下游调用从串行改为并发协程,用asyncio.gather()同时发起。
  2. 限流:防止下游被打爆,用asyncio.Semaphore(10)控制最大并发协程数。
  3. 超时:每个协程用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。总结三点经验供参考:

  1. 适用场景:纯IO密集型、下游API多、依赖关系可分层并发。CPU密集型(如加解密、图像处理)不适合,应改用多进程。

  2. 框架选择:Flask做异步需要自己套loop,比较别扭。新项目建议直接用FastAPI或Sanic,原生支持async。但存量Flask项目用上述模式改造,成本很低。

  3. 监控与降级:异步代码里异常容易被吞,务必在协程内捕获并记录日志。我加了try/except和超时返回None的降级逻辑,保证主流程不挂。

最后,代码已上线一个月,线上P99稳定在650ms左右,QPS峰值500+,再也没收到过下游抖动引发的雪崩告警。如果你也遇到类似的BFF层阻塞问题,可以试试这个改造路径——成本不高,收益巨大。有问题欢迎评论区交流。