1. 问题背景:一个串行调用3个下游接口的报表服务

我们的订单报表接口 /api/orders/summary 需要聚合三个数据源:
- 用户服务user-service):返回用户等级、地域(平均耗时800ms)
- 订单服务order-service):返回订单金额、数量(平均耗时700ms)
- 库存服务inventory-service):返回库存预警信息(平均耗时600ms)

原实现用 requests 库同步调用,代码直白,但生产环境一压测就露馅:

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

app = Flask(__name__)

@app.route('/api/orders/summary')
def summary():
    user = requests.get('http://user-service/api/user', timeout=3).json()
    order = requests.get('http://order-service/api/order', timeout=3).json()
    stock = requests.get('http://inventory-service/api/stock', timeout=3).json()

    return jsonify({
        'user_level': user['level'],
        'order_amount': order['amount'],
        'stock_warning': stock['warning']
    })

性能实测(wrk -t4 -c200 -d30s):
- 平均延迟:2150ms
- P99:3100ms
- 吞吐率:180 req/s
- CPU:85%(GIL锁竞争 + 线程池切换开销)

三个请求完全独立,却串行等待,纯粹是代码习惯问题。

2. 环境与版本:明确依赖,避免玄学

  • Python 3.10.12(原生支持 asyncio.run(),不需要 nest_asyncio
  • Flask 2.3.3(注意:Flask 是 WSGI 同步框架,不能直接在视图函数里 await
  • aiohttp 3.9.1(替代 requests 的异步客户端)
  • gunicorn 21.2.0(生产部署,worker 模式改为 geventgthread,见下文坑)

3. 方案设计:用 asyncio 并发等待,而不是并行计算

核心思路:把三个 HTTP 调用变成协程,用 asyncio.gather 并发执行。但有两个关键约束:

  1. Flask 视图函数是同步的,不能直接 await。解决方案:用 asyncio.run() 包装一个异步函数,但注意不能在每个请求里都 asyncio.run()(会重复创建事件循环,性能差)。
  2. 必须限流:如果下游服务扛不住并发,直接把对方打挂。用 asyncio.Semaphore(10) 限制最大并发数为10。

架构图(文字版):

Flask 同步视图
    ↓ asyncio.run()
异步主函数
    ↓ 创建 Semaphore(10)
    ↓ asyncio.gather(协程1, 协程2, 协程3)
    ↓ 每个协程内部: async with semaphore: await aiohttp.get()

4. 核心实现:异步重构代码

第一步:初始化 aiohttp ClientSession(全局复用,避免每次请求创建连接池)

# async_client.py
import asyncio
import aiohttp

# 全局复用 session,连接池默认 100 个连接
session = None

async def init_session():
    global session
    if session is None:
        # connector 调大连接池,TCPConnector 默认 limit=100
        connector = aiohttp.TCPConnector(limit=100, ttl_dns_cache=300)
        session = aiohttp.ClientSession(connector=connector)

async def get_json(url, semaphore):
    async with semaphore:
        async with session.get(url, timeout=aiohttp.ClientTimeout(total=3)) as resp:
            return await resp.json()

第二步:重写 Flask 视图函数

# after.py - 异步版本
import asyncio
from flask import Flask, jsonify
from async_client import init_session, get_json, session

app = Flask(__name__)

# 全局信号量:限制最大10个并发下游请求
semaphore = asyncio.Semaphore(10)

def async_summary():
    """异步聚合三个服务数据"""
    async def _fetch_all():
        await init_session()  # 确保 session 已创建
        urls = [
            'http://user-service/api/user',
            'http://order-service/api/order',
            'http://inventory-service/api/stock'
        ]
        # gather 并发执行,return_exceptions=True 避免单个失败拖垮全部
        results = await asyncio.gather(
            *(get_json(url, semaphore) for url in urls),
            return_exceptions=True
        )
        # 解析结果,处理异常情况
        user, order, stock = results
        if isinstance(user, Exception):
            user = {'level': 'N/A'}
        if isinstance(order, Exception):
            order = {'amount': 0}
        if isinstance(stock, Exception):
            stock = {'warning': 'unknown'}
        return {
            'user_level': user['level'],
            'order_amount': order['amount'],
            'stock_warning': stock['warning']
        }

    # 每个请求创建新的事件循环(注意:Flask 多线程下必须这样)
    return asyncio.run(_fetch_all())

@app.route('/api/orders/summary')
def summary():
    return jsonify(async_summary())

生产部署注意asyncio.run() 每次创建新事件循环,如果 init_session() 反复执行,会警告 unclosed session。解决:把 init_session() 放在模块导入时做一次,但 Flask 的 debug 模式会重载模块,所以用全局 if session is None 判断即可。

5. 踩坑与优化:三个真实遇到的坑

坑1:asyncio.run() 在 gunicorn 多 worker 下的坑
- 现象:部署后报 RuntimeError: asyncio.run() cannot be called from a running event loop
- 原因:gunicorn 默认 sync worker 是单线程的,没问题。但如果你用了 gevent worker,gevent 会 patch 标准库,导致 asyncio 检测到已存在事件循环。
- 解决:gunicorn worker 类改用 gthread(线程池模式),配置 --worker-class gthread --threads 4。或者干脆不用 gunicorn,用 uvicorn 跑 ASGI 框架。但我为了兼容现有 Flask 代码,选择了 gthread。

坑2:aiohttp 连接池耗尽
- 现象:压测时大量 TimeoutError: Connect timeout,查看下游服务 CPU 正常。
- 原因:aiohttp 默认 TCPConnector(limit=100),但我们的信号量是10,按理说不会超过10个并发连接。排查发现是 DNS 解析阻塞——aiohttp 默认用线程池解析 DNS,高并发下线程池满了。
- 解决:TCPConnector(ttl_dns_cache=300),缓存DNS 5分钟。

坑3:asyncio.Semaphore 作用域
- 如果信号量定义在函数内部,每个请求会创建新的信号量,限流形同虚设。必须定义为模块级全局变量

6. 效果数据:压测对比

wrk -t4 -c200 -d60s http://localhost:5000/api/orders/summary 压测,结果:

指标 同步版本 异步版本 提升
平均延迟 2150ms 420ms 5.1x
P99延迟 3100ms 780ms 3.9x
吞吐率 180 req/s 920 req/s 5.1x
CPU占用 85% 55% -35%
错误率 0.2% 0.1% -

额外验证:把下游服务模拟延迟提高到2秒(用 time.sleep(2)),同步版本延迟直接飙到6.1秒(串行),异步版本仍然稳定在2.1秒(并行),说明并发红利完全取决于下游延迟。

7. 总结与进一步优化建议

  • 适用场景:I/O密集型(HTTP调用、数据库查询、文件读取),不适用于CPU密集型任务(asyncio不解决GIL问题)。
  • 不要过度设计:如果下游只有1个接口,同步就够。3个以上独立调用才值得改造。
  • 后续可做
  • asyncio.Queue 做请求合并(类似 Hystrix 的 request collapsing),批量调用下游。
  • 超时控制用 asyncio.wait_for,替代 aiohttp 的 timeout 参数,更精细。
  • 如果想完全摆脱 Flask 同步限制,可以直接迁移到 FastAPI(原生 async 支持),但迁移成本需评估。

最后说一句:异步编程不是银弹,但它能让你的服务在高并发下从「勉强能跑」变成「游刃有余」。关键是把业务拆分成可并发的子任务,然后用最简单的 gather + Semaphore 组合拳,别一上来就上消息队列。