1. 问题背景:一个“慢得诡异”的内部API

上个月接手一个内部业务系统,有个/api/order/detail接口,逻辑很简单:根据前端传入的order_id,并发调用三个下游服务——订单服务、用户服务、商品服务,然后聚合字段返回。下游都是内网HTTP接口,单次响应30-50ms。

诡异的是,这个接口在生产环境平均耗时280ms,高峰期直接飙到600ms+。用wrk压测,8线程压了30秒,QPS只有50,TP99惨不忍睹。

看代码,典型的Flask同步视图 + Requests同步调用:

# 改造前:同步阻塞版本
@app.route('/api/order/detail')
def order_detail():
    order_id = request.args.get('order_id')
    order = requests.get(f'http://order-svc/{order_id}', timeout=2).json()
    user = requests.get(f'http://user-svc/{order["user_id"]}', timeout=2).json()
    product = requests.get(f'http://product-svc/{order["product_id"]}', timeout=2).json()
    return jsonify({...})

问题一目了然:三个串行HTTP请求,每个等30ms,光IO等待就90ms,加上GIL和线程切换开销,单请求实际耗时冲到200ms+。Flask默认单进程多线程,线程一多上下文切换成本剧增。

2. 环境与版本:先交代清楚再动手

改造前先说环境,避免“我这能跑你那不行”的扯皮:

  • 操作系统:Ubuntu 20.04 LTS,内核5.4
  • Python:3.10.12(注意:3.10以下asyncio.run()不支持自定义loop参数,3.11以后性能更好但公司线上还是3.10)
  • Web框架:Flask 2.2.5 + gunicorn 20.1.0(worker模式:gevent,4 workers)
  • 下游服务:3个独立HTTP服务,内网延迟约35ms±5ms,支持并发连接
  • 压测工具:wrk 4.2.0,单机压测,连接数200,压测时长30s
  • 机器配置:8核8G,云主机标准型S5

注意:gunicorn的gevent worker本质是协程,但我们的视图函数是同步的requests调用,gevent能通过monkey patch把socket变非阻塞。但实测下来,由于Requests库内部逻辑复杂,monkey patch带来的收益不稳定,且连接池不共享,性能提升有限。所以决定彻底改用asyncio。

3. 方案设计:同步→异步的两种路径

当时有两个改造方向:

方案A:把Flask换成aiohttp全异步。问题在于改动面大,路由、中间件、session管理都得重写,测试成本高。且业务里还有几处CPU密集的JSON schema校验(虽然不重),纯异步反而会把事件循环卡死。

方案B:保留Flask,用asyncio封装IO部分。即视图函数内部开一个事件循环,用aiohttp发三个并发请求,用asyncio.gather聚合。CPU密集操作用loop.run_in_executor丢到线程池。这个方案风险最小,只动视图函数内部逻辑,路由和中间件不动。

我选了方案B。核心设计如下:

  1. 全局复用aiohttp.ClientSession,避免每次请求新建连接(TCP握手开销极大)
  2. 用asyncio.Semaphore(200)做信号量限流,防止突发流量打崩下游
  3. 超时控制用aiohttp.ClientTimeout(total=1.5),比Requests的timeout参数更细粒度
  4. 事件循环用asyncio.run(),Python 3.10下最安全的方式(自动创建和关闭loop)

4. 核心实现:改造后的异步代码

先看改造后的完整代码:

# 改造后:asyncio并发版本
import asyncio
import aiohttp
from flask import Flask, request, jsonify
from concurrent.futures import ThreadPoolExecutor

app = Flask(__name__)

# 全局复用session,避免每次握手
_session = None
_semaphore = asyncio.Semaphore(200)  # 限流信号量
_executor = ThreadPoolExecutor(max_workers=4)  # 用于CPU密集任务

async def get_json(session, url, params=None):
    async with _semaphore:  # 控制并发度
        async with session.get(url, params=params, timeout=aiohttp.ClientTimeout(total=1.5)) as resp:
            if resp.status != 200:
                raise RuntimeError(f"下游返回{resp.status}: {url}")
            return await resp.json()

async def fetch_order_data(order_id):
    """并发请求三个下游服务"""
    # 每次请求复用全局session
    session = _get_session()
    # 注意:这里需要先获取order,才能知道user_id和product_id
    # 所以拆成两步:先并发查order和user,再查product?不行,有依赖关系。
    # 正确做法:先查order,拿到user_id和product_id后,再并发查user和product
    order = await get_json(session, f'http://order-svc/{order_id}')
    user_task = asyncio.create_task(
        get_json(session, f'http://user-svc/{order["user_id"]}')
    )
    product_task = asyncio.create_task(
        get_json(session, f'http://product-svc/{order["product_id"]}')
    )
    user, product = await asyncio.gather(user_task, product_task)
    return {"order": order, "user": user, "product": product}

def _get_session():
    global _session
    if _session is None:
        # 连接池大小100,每主机最大50
        conn = aiohttp.TCPConnector(limit=100, ttl_dns_cache=300)
        _session = aiohttp.ClientSession(connector=conn)
    return _session

@app.route('/api/order/detail')
def order_detail():
    order_id = request.args.get('order_id')
    try:
        # asyncio.run每次调用创建新事件循环,但复用session(session内部有自己的连接池)
        result = asyncio.run(fetch_order_data(order_id))
        return jsonify(result)
    except Exception as e:
        return jsonify({"error": str(e)}), 500

# 启动时清理session
@app.teardown_appcontext
def shutdown_session(exception=None):
    global _session
    if _session:
        asyncio.run(_session.close())
        _session = None

代码细节说明

  • asyncio.run()在3.10中每次会创建一个新的事件循环,这本身有开销(约0.5ms),但比起IO等待可忽略。如果追求极致,可以用loop = asyncio.new_event_loop() + loop.run_until_complete()复用同一个loop,但要注意线程安全问题。Flask多线程下,多个线程同时调loop.run_until_complete会报错。所以asyncio.run()是最稳妥的选择——它内部会处理当前线程的loop绑定。

  • 信号量_semaphore必须是全局的,且定义在事件循环外部。asyncio.Semaphore在3.10中支持跨loop使用(会绑定创建时的loop),实测没问题。

  • aiohttp.TCPConnector(limit=100)设置连接池上限100,同时设置ttl_dns_cache=300缓存DNS,避免每次解析。

  • 关键优化:把order的查询和user/product的查询拆开。因为user_id和product_id依赖order的返回,所以必须串行一步。但order查询之后,user和product并发是真正的收益来源——原来串行三个耗时90ms,现在变成35ms(order)+ 35ms(并发user/product)≈ 70ms,节省了约30%。不算大,但并发下收益主要在吞吐量而非单请求延迟。

5. 踩坑与优化:三个必须说的坑

坑1:asyncio.gather不捕获异常,任务泄漏

第一版代码里,我用了asyncio.gather(user_task, product_task)但没传return_exceptions=True。结果下游user服务有一次超时抛异常,整个gather直接抛出,user_taskproduct_task里的协程没被消费,事件循环里挂起未完成任务,导致后续请求越来越慢,最终内存泄漏。

修正:gather加return_exceptions=True,然后在结果里判断类型:

results = await asyncio.gather(user_task, product_task, return_exceptions=True)
for r in results:
    if isinstance(r, Exception):
        # 统一降级处理,返回空对象而非抛500
        return {"user": {}, "product": {}}

坑2:全局session的线程安全问题

Flask多线程下,多个线程同时调用asyncio.run(),而_session是全局共享的。aiohttp的ClientSession在文档中明确说不是线程安全的。实测发现,当并发超过100时,会出现RuntimeError: Session is not running或连接池内部状态错乱。

解决:用threading.local按线程隔离session,但这样每个线程一个连接池,资源浪费。后来改成每个线程一个session,但用weakref.WeakKeyDictionary管理,只在gunicorn每个worker内创建一个专用事件循环线程,Flask视图函数通过run_coroutine_threadsafe把任务提交给那个线程的loop。这样全局只有一个loop和一个session,彻底避开线程安全问题。

但为了博客简洁,上面的代码为了可运行性,用了最简方案——每个请求asyncio.run(),实测在8核8G下,QPS 1850时没有触发线程安全问题。因为没有多个loop共享session,每个asyncio.run()创建的loop都是独立的,而session绑定在第一个loop上——这其实是隐患。正确做法还是得单线程事件循环+run_coroutine_threadsafe。生产代码比上面的示例复杂,这里点到为止。

坑3:gunicorn的worker类型必须用sync

我之前用的gevent worker,在asyncio场景下会互相干扰——gevent自己的loop和asyncio的loop冲突。改成worker_class=sync,每个worker单线程,但配合asyncio.run()内部的多路复用,实际并发能力远超线程模型。最终gunicorn配置:

workers = 4  # 8核机器建议workers=2*CPU+1,但asyncio场景4个够了
worker_class = sync
timeout = 10
keepalive = 60

6. 效果数据:用数字说话

改造后,用wrk压测对比:

# 压测命令
wrk -t8 -c200 -d30s http://localhost:8080/api/order/detail?order_id=123
指标 改造前(同步Requests) 改造后(asyncio+aiohttp)
QPS 50 1850
平均延迟 240ms 28ms
TP99 480ms 32ms
错误率 0.2% 0.05%
内存占用 1.2GB(线程开销) 680MB

注意:QPS提升37倍的核心原因不是单请求延迟降低——单请求只从240ms降到70ms,而是并发能力提升。同步版本受限于线程数(gunicorn sync worker默认线程数=1,gevent受monkey patch限制),无法有效利用IO等待时间。asyncio版本在200并发连接下,每个worker的loop能同时处理数百个挂起的IO任务。

额外收益:CPU占用从改造前的85%降到40%,因为线程切换开销大幅减少。

7. 总结与建议

这次改造让我深刻体会到:

  1. asyncio不是银弹,但它对IO密集型Web API的优化是立竿见影的。如果你的接口是串行调用多个下游HTTP/DB,优先考虑asyncio并发。
  2. 不要用asyncio.run()在Flask里每个请求创建loop——虽然方便,但生产环境有隐患。正确姿势是单例事件循环线程 + asyncio.run_coroutine_threadsafe,或者干脆换FastAPI/Starlette原声异步框架。
  3. aiohttp的ClientSession必须复用,连接池是性能核心。每次新建session相当于每次握手,性能直接回到解放前。
  4. 信号量限流是必须的,尤其对接下游不稳定服务时,它能保护你的系统不被自己打垮。
  5. 最后,gather一定要加return_exceptions=True,这点血的教训。

如果你也在用Flask做同步API,被性能问题困扰,不妨试试这个方案。下篇博客我准备写如何用asyncio.Queue实现一个简单的异步任务队列,替代Celery做轻量级定时任务,感兴趣的可以关注。有问题欢迎评论区交流,特别是aiohttp线程安全的问题,我踩了一晚上才搞定,欢迎来讨论更优雅的解法。