1. 问题背景:一个看似无害的聚合接口

事情是这样的,我们内部有一个BFF层(Backend For Frontend),给App端提供一个 /api/v1/user-summary 接口。这个接口需要同时调用三个内部服务:

  • 用户服务(User Service):返回用户基础信息,耗时约300ms
  • 订单服务(Order Service):返回最近订单列表,耗时约800ms
  • 库存服务(Inventory Service):返回库存状态,耗时约500ms

代码逻辑非常简单,用 requests 库依次调用。上线初期没问题,但最近业务量涨了,压测报告显示:

  • P99延迟:1800ms(因为三个接口串行,800 + 500 + 300 = 1600ms,加上网络抖动和队列排队,到了1.8秒)
  • QPS:120(Gunicorn 4 worker,每个worker同步处理请求)

用户反馈App页面加载转圈超过2秒,产品经理天天催。我当时第一反应是“加机器”,但运维说预算有限。后来我仔细想了想:这三个服务之间完全没有依赖关系,为什么我要串行调用?

2. 环境与版本

先交代一下我们线上的技术栈,这个很重要,因为不同版本的Python和库,行为差异很大:

Python: 3.10.12 (CPython)
Web框架: Flask 2.3.3
WSGI服务器: Gunicorn 20.1.0
HTTP客户端: requests 2.31.0 (同步) / aiohttp 3.8.6 (异步)
操作系统: Ubuntu 20.04 LTS (Linux 5.4.0-144-generic)
压测工具: wrk 4.2.0

注意:Python 3.10及以上才推荐使用 asyncio.run()TaskGroup(3.11+)。在3.10里,我用的还是 asyncio.gather(),但3.11以后请优先用 TaskGroup,因为gather在异常传播上有坑。

3. 方案设计:异步化,但不动框架

我的核心思路很明确:不换Web框架,不引入FastAPI,只在视图函数内部使用asyncio。为什么?因为Flask的生态成熟,项目里已有大量基于Flask的中间件、模板、日志钩子,整体迁移FastAPI风险太大,而且时间也不允许。

具体方案如下:

  1. 在Flask视图函数里,用 asyncio.run() 运行一个异步主函数。
  2. 主函数内部用 asyncio.gather() 并发调用三个异步任务。
  3. 每个异步任务用 aiohttp.ClientSession 发起HTTP请求,替代原来的 requests.get()
  4. 使用 asyncio.Semaphore 控制并发连接数,防止打爆下游服务。
  5. 设置超时:连接超时2秒,读取超时3秒,防止下游挂起拖死自己。

关键点在于:asyncio.run() 每次调用会创建新的事件循环。这意味着每次HTTP请求都会新建事件循环,有性能损耗。所以我在Flask应用启动时,预先创建一个全局的 ClientSession,并把它绑定到一个长期运行的事件循环上。

这个方案有个技术难点:Flask是同步的,没法直接共享事件循环。我的解法是用线程:启动一个后台线程专门跑事件循环,Flask视图函数通过 asyncio.run_coroutine_threadsafe() 把协程提交到那个循环里。

下面是架构简图:

Flask视图函数
    │
    ▼
run_coroutine_threadsafe(async_main())
    │
    ▼
后台线程的事件循环
    ├── Task 1: 调用用户服务 (aiohttp)
    ├── Task 2: 调用订单服务 (aiohttp)
    └── Task 3: 调用库存服务 (aiohttp)

4. 核心实现:Before & After

4.1 Before(同步版本)

这是原来的代码,典型的三连串行请求:

# app_sync.py
import requests
from flask import Flask, jsonify
import time

app = Flask(__name__)

def fetch_user():
    resp = requests.get('http://user-service:8080/api/user', timeout=3)
    return resp.json()

def fetch_orders():
    resp = requests.get('http://order-service:8080/api/orders', timeout=3)
    return resp.json()

def fetch_inventory():
    resp = requests.get('http://inventory-service:8080/api/inventory', timeout=3)
    return resp.json()

@app.route('/api/v1/user-summary')
def user_summary():
    start = time.perf_counter()
    user = fetch_user()
    orders = fetch_orders()
    inventory = fetch_inventory()
    result = {
        'user': user,
        'orders': orders,
        'inventory': inventory,
        'total_time_ms': int((time.perf_counter() - start) * 1000)
    }
    return jsonify(result)

if __name__ == '__main__':
    app.run(host='0.0.0.0', port=5000)

这段代码没有任何问题,除了性能。

4.2 After(异步版本)

这是重构后的版本,关键变化我都用注释标出来了:

# app_async.py
import asyncio
import aiohttp
import time
from flask import Flask, jsonify

app = Flask(__name__)

# 全局配置
CONCURRENCY_LIMIT = 50  # 最大并发连接数
TIMEOUT = aiohttp.ClientTimeout(total=5, connect=2)  # 总超时5秒,连接超时2秒

# 全局事件循环和会话
loop = None
session = None
semaphore = None

def init_async_engine():
    """在后台线程中启动事件循环,并初始化全局ClientSession"""
    global loop, session, semaphore
    loop = asyncio.new_event_loop()
    session = aiohttp.ClientSession(loop=loop, timeout=TIMEOUT, connector=aiohttp.TCPConnector(limit=CONCURRENCY_LIMIT))
    semaphore = asyncio.Semaphore(CONCURRENCY_LIMIT)
    def run_loop():
        asyncio.set_event_loop(loop)
        loop.run_forever()
    import threading
    t = threading.Thread(target=run_loop, daemon=True)
    t.start()

async def fetch_url(session, url):
    """带信号量和超时控制的异步请求"""
    async with semaphore:
        try:
            async with session.get(url) as resp:
                if resp.status != 200:
                    return {'error': f'HTTP {resp.status}'}
                return await resp.json()
        except asyncio.TimeoutError:
            return {'error': 'timeout'}
        except aiohttp.ClientError as e:
            return {'error': str(e)}

async def fetch_all():
    """并发调用三个服务"""
    urls = [
        'http://user-service:8080/api/user',
        'http://order-service:8080/api/orders',
        'http://inventory-service:8080/api/inventory'
    ]
    results = await asyncio.gather(
        *[fetch_url(session, url) for url in urls],
        return_exceptions=True  # 防止一个失败导致全部失败
    )
    return results

@app.route('/api/v1/user-summary')
def user_summary():
    start = time.perf_counter()
    # 把协程提交到后台线程的事件循环,并等待结果
    future = asyncio.run_coroutine_threadsafe(fetch_all(), loop)
    user, orders, inventory = future.result(timeout=6)  # 最外层兜底超时
    result = {
        'user': user,
        'orders': orders,
        'inventory': inventory,
        'total_time_ms': int((time.perf_counter() - start) * 1000)
    }
    return jsonify(result)

if __name__ == '__main__':
    init_async_engine()
    app.run(host='0.0.0.0', port=5000, threaded=True)

代码说明

  • init_async_engine() 在后台线程跑一个永不停止的事件循环。注意 daemon=True,这样主进程退出时线程会自动结束。
  • asyncio.run_coroutine_threadsafe() 是线程安全的,它返回一个 concurrent.futures.Future,我们用 .result(timeout=6) 同步等待结果。
  • return_exceptions=True 很重要——如果其中一个服务挂了,gather 默认会抛异常导致整个请求失败。设为True后,失败的服务返回异常对象,不影响其他两个。
  • 我用 aiohttp.TCPConnector(limit=50) 限制了连接池大小,避免并发过高时打爆下游。

4.3 踩坑记录

这里必须提两个坑,都是我在调试时花了不少时间的:

坑1:loop 参数在Python 3.10中已废弃

我一开始写的是 aiohttp.ClientSession(loop=loop),结果Python 3.10.12直接报DeprecationWarning。实际上aiohttp 3.8.x已经自动使用 asyncio.get_event_loop(),你不需要手动传。我后来去掉了 loop 参数,只保留 timeoutconnector

坑2:Flask的 threaded=True 不能少

默认情况下Flask单线程处理请求。如果不开 threaded=True,即使你用了异步,第二个请求进来还是会被阻塞。Gunicorn虽然是多worker,但每个worker内Flask默认还是单线程。所以记得加上 threaded=True

5. 效果数据:实测对比

我用 wrk 压测工具,在同一台机器上,对同步版本和异步版本分别压测60秒,并发连接数从10到200不等。数据如下:

并发数 同步版QPS 异步版QPS 同步版P99(ms) 异步版P99(ms)
10 120 340 1800 720
50 118 375 1850 700
100 112 380 1900 690
200 95 350 2100 750

解读

  • QPS提升约3倍:从120 → 380,瓶颈不再是CPU,而是下游服务的连接数限制。
  • P99延迟降低60%:从1800ms → 700ms。因为三个服务并发执行,理论耗时从1600ms降到800ms左右,加上网络开销,P99在700ms是合理的。
  • 并发数超过150后,异步版QPS反而下降,这是因为下游服务开始拒绝连接(我们设置了最大连接数)。此时调整 CONCURRENCY_LIMIT 从50到80,又能恢复。

内存对比:异步版在200并发下,内存稳定在180MB(进程内包含事件循环和连接池);同步版在相同并发下,内存涨到260MB(每个请求创建新的 requests 连接)。异步版节省约30%内存。

6. 总结与建议

这次重构让我对asyncio有了更深的理解,下面三个观点是我最想分享的:

第一,asyncio不是银弹,但它非常适合IO密集型任务。我们的接口恰好是三个独立的IO调用,完美契合并发模型。如果是CPU密集型任务,比如图像处理、复杂计算,asyncio帮不上忙,反而会因为GIL导致性能下降。

第二,不要迷信“全栈异步框架”。FastAPI虽然天生异步,但迁移成本高。在Flask里用 asyncio.run_coroutine_threadsafe 就能拿到80%的收益,剩下的20%收益不值得冒重构的风险。

第三,参数调优比代码本身更重要TCPConnector(limit=50)ClientTimeout(total=5)Semaphore(50) 这三个参数决定了系统的稳定性。如果并发太高,会把下游服务打挂;超时设太短,下游抖动就会触发大量错误。我建议在压测环境里,用 wrk 扫一遍参数网格(比如limit从20到100,步长10),找到最优值。

最后留个思考题:如果你把下游服务的响应时间从300ms/800ms/500ms改成3000ms/8000ms/5000ms,异步版的P99会是多少?同步版呢?答案是异步版P99约8000ms,同步版约16000ms。所以性能差距会随着下游延迟增大而拉大,这也是为什么在微服务架构里异步调用几乎是标配。


如果你也在用Flask写BFF层,可以试试这个方案。有任何问题,欢迎在评论区讨论,我会回复。