1. 问题背景:一个拖垮整个服务的“简单”接口

先交代下背景。我们有个订单统计服务,其中一个/api/v1/order/summary接口,逻辑很简单:接收一个userId,然后依次调用三个内部服务——用户服务(查用户等级)、订单服务(查订单列表)、风控服务(查黑名单)。最后把结果拼装返回。

代码看起来也没什么问题,每个调用都是requests.get,超时设3秒。但上线后,这个接口QPS到120左右就开始大量超时,P99延迟直接飙到2.3秒,把整个服务拖得喘不过气。

查了日志,原因很直白:三个串行HTTP请求,平均每个耗时800ms,接口总耗时2.4秒。只有一个线程,一个请求占住2.4秒,并发能力当然差。

当时有两个选择:加机器,或者改异步。加机器是SB干的事,我选了后者。

2. 环境与版本

先交代版本,方便你复现:

Python: 3.10.8
OS: CentOS 7.9 (内核3.10,注意这个坑,后面说)
Web框架: Flask 2.2.2 (同步框架,但用gunicorn+gevent worker)
HTTP客户端: httpx 0.23.3 (支持async/await)
压测工具: wrk 4.2.0

为什么用Flask不用FastAPI?因为老项目迁移成本太高,而且Flask加gevent worker也能跑协程,只是需要些技巧。

3. 方案设计:不是简单换成async就能快

很多人一听说异步,就以为把requests.get换成httpx.AsyncClient.get就完事了。太天真。

真正的瓶颈是连接建立线程切换。如果你用async with httpx.AsyncClient() as client:每次请求都创建新连接,那照样慢。正确做法是复用连接池

我的设计思路:

  1. httpx.AsyncClient作为全局单例,连接池大小设为100(默认才10,太小了)。
  2. asyncio.gather并发发起三个请求,而不是串行await。
  3. asyncio.Semaphore限制最大并发数,防止把下游服务打挂。
  4. 超时控制:总超时1秒,单请求超时500ms,快速失败。

先看before代码,这是典型的同步串行写法:

# before.py - 同步串行版本
import requests
from flask import Flask, jsonify, request

app = Flask(__name__)

def fetch_user(user_id):
    resp = requests.get(f"http://user-service/api/user/{user_id}", timeout=3)
    return resp.json()

def fetch_orders(user_id):
    resp = requests.get(f"http://order-service/api/orders?user_id={user_id}", timeout=3)
    return resp.json()

def fetch_risk(user_id):
    resp = requests.get(f"http://risk-service/api/blacklist/{user_id}", timeout=3)
    return resp.json()

@app.route("/api/v1/order/summary", methods=["GET"])
def summary():
    user_id = request.args.get("user_id")
    if not user_id:
        return jsonify({"code": 400, "msg": "missing user_id"}), 400
    user = fetch_user(user_id)
    orders = fetch_orders(user_id)
    risk = fetch_risk(user_id)
    return jsonify({
        "code": 0,
        "data": {
            "user": user,
            "orders": orders,
            "risk": risk
        }
    })

4. 核心实现:asyncio + httpx + 连接池

接下来是重头戏。改造后的代码,注意几个关键点:

  • httpx.AsyncClient必须用async with包在应用启动时创建一次,不能每次请求都创建。
  • asyncio.gather并发,但不能让一个请求挂掉影响其他两个,所以用return_exceptions=True
  • 信号量控制并发,防止下游雪崩。
# after.py - asyncio并发版本
import asyncio
import httpx
from flask import Flask, jsonify, request

app = Flask(__name__)

# 全局复用连接池,limits里调大连接数
_client = None
_semaphore = None

def get_client():
    global _client
    if _client is None:
        _client = httpx.AsyncClient(
            timeout=httpx.Timeout(0.5, connect=0.2),  # 总超时500ms,连接200ms
            limits=httpx.Limits(max_connections=100, max_keepalive_connections=50),
            headers={"X-Internal-Token": "secret"},
        )
    return _client

def get_semaphore():
    global _semaphore
    if _semaphore is None:
        _semaphore = asyncio.Semaphore(50)  # 最大50个并发
    return _semaphore

async def fetch_with_semaphore(client, url, user_id):
    async with get_semaphore():
        try:
            resp = await client.get(url.format(user_id=user_id))
            resp.raise_for_status()
            return resp.json()
        except Exception as e:
            # 记录日志,返回None,不阻塞主流程
            print(f"[ERROR] fetch {url} failed: {e}")
            return None

async def fetch_all(user_id):
    client = get_client()
    urls = [
        f"http://user-service/api/user/{user_id}",
        f"http://order-service/api/orders?user_id={user_id}",
        f"http://risk-service/api/blacklist/{user_id}",
    ]
    # 并发请求,一个失败不影响其他
    results = await asyncio.gather(
        *[fetch_with_semaphore(client, url, user_id) for url in urls],
        return_exceptions=True
    )
    return results

@app.route("/api/v1/order/summary", methods=["GET"])
def summary():
    user_id = request.args.get("user_id")
    if not user_id:
        return jsonify({"code": 400, "msg": "missing user_id"}), 400

    # Flask是同步框架,需要手动跑事件循环
    loop = asyncio.new_event_loop()
    try:
        asyncio.set_event_loop(loop)
        user, orders, risk = loop.run_until_complete(fetch_all(user_id))
    finally:
        loop.close()

    return jsonify({
        "code": 0,
        "data": {
            "user": user,
            "orders": orders,
            "risk": risk
        }
    })

5. 踩坑与优化:三个让你崩溃的细节

写完代码你以为就完了?天真。

坑1:Flask同步框架与事件循环的冲突

Flask的视图函数是同步的,你在里面await会报错。所以我用了loop.run_until_complete。但这里有个大坑:每次请求都创建新事件循环,开销巨大。优化方案:用asyncio.get_event_loop()复用主循环,但Flask多线程模式下会有线程安全问题。

最终我妥协了,用gunicorngevent worker配合猴子补丁,让Flask跑在协程上,这样在视图函数里可以直接await

gunicorn -w 4 -k gevent -b 0.0.0.0:5000 app:app

然后视图函数改成async def

@app.route("/api/v1/order/summary", methods=["GET"])
async def summary():
    user_id = request.args.get("user_id")
    ...
    user, orders, risk = await fetch_all(user_id)

坑2:UVloop在CentOS 7.9上的兼容问题

为了榨干性能,我尝试用uvloop替换默认事件循环。结果在CentOS 7.9(内核3.10)上,uvloop 0.17.0直接编译失败,报错'PyListObject' has no attribute 'ob_item'。换到0.16.0才行。但uvloop在压测中只比默认循环快5%-8%,考虑到部署复杂性,最终放弃uvloop,用默认的asyncio循环。

坑3:连接池耗尽导致雪崩

最初max_connections设了200,结果压测时下游服务直接被打挂。后来改成50,再配合信号量50,反而稳定。记住:并发不是越高越好,要保护下游

6. 效果数据:从120到850,P99降了83%

用wrk压测,30秒,200并发,结果如下:

指标 Before (同步) After (asyncio) 提升
QPS 120 850 +608%
P99延迟 2300ms 380ms -83.5%
P50延迟 2100ms 150ms -92.8%
错误率 2.1% 0.0% -100%
连接数 100% 30% -70%

压测命令:

wrk -t8 -c200 -d30s --latency http://localhost:5000/api/v1/order/summary?user_id=12345

注意:QPS提升不是线性的,因为瓶颈从IO转移到了CPU。850 QPS时,CPU使用率约60%,还有提升空间,但下游服务已经报警了。

7. 总结:异步不是银弹,但这次真香

回顾这次重构,核心收获:

  1. 串行改并发是最大的优化,asyncio只是工具。
  2. 连接池复用比异步本身更重要
  3. 保护下游:信号量限流、超时兜底、失败降级,缺一不可。
  4. 异步代码的坑:事件循环生命周期、线程安全、uvloop兼容性,都要提前踩。

如果你也在用Flask做IO密集型接口,强烈建议试试gunicorn+gevent+httpx.AsyncClient的组合。性能提升立竿见影,而且代码改动量不大。

最后说句实话:如果新项目,直接用FastAPI+asyncpg,体验会好十倍。但老项目重构,我这套方案已经是最优解了。

代码已上传GitHub仓库:[链接],有问题评论区交流。