一、问题背景:一个看似“没问题”的聚合API

我们有个订单服务,前端需要一次性展示订单详情、用户信息、物流轨迹、优惠券状态。最初实现是同步串行调用四个下游服务:

def get_order_detail(order_id):
    user = requests.get(f"http://user-service/user/{order_id}").json()
    logistics = requests.get(f"http://logistics-service/track/{order_id}").json()
    coupon = requests.get(f"http://coupon-service/coupon/{order_id}").json()
    ...

线上运行半年,日均调用量不大(约200万),一直没出问题。直到双十一前压测,50并发下QPS 37,P99 4.8s——这数据根本没法上线。四个下游服务响应时间都在200-500ms,串行累加后单请求耗时轻松超过1.5s。

二、环境与版本:先交代清楚

  • Python 3.10.11(3.10+的asyncio才比较稳定,3.8的坑太多)
  • Flask 2.1.2(注意:Flask本身是WSGI同步框架,asyncio需要自己桥接)
  • aiohttp 3.8.4(替代requests的关键)
  • wrk 4.2.0(压测工具)
  • 服务器:8核8G CentOS 7.9,下游服务在同一内网,延迟约50ms

三、方案设计:不是简单替换,而是分层改造

直接上asyncio不能解决所有问题。我拆了三步:

  1. 把四个串行IO改为并发协程:用asyncio.gather同时发起四个请求
  2. 限制并发度:防止突发流量把下游打爆,用Semaphore(100)做信号量
  3. 超时熔断:每个请求用asyncio.wait_for设置500ms超时,超时直接返回降级数据

核心设计原则:异步化不改变业务逻辑,只改变IO等待方式。下游接口不兼容协程没关系,aiohttp直接替代requests,返回体结构相同。

四、核心实现:before/after代码对比

Before:同步串行版(50行)

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

app = Flask(__name__)

def fetch_user(order_id):
    return requests.get(f"http://user-service/user/{order_id}", timeout=1).json()

def fetch_logistics(order_id):
    return requests.get(f"http://logistics-service/track/{order_id}", timeout=1).json()

def fetch_coupon(order_id):
    return requests.get(f"http://coupon-service/coupon/{order_id}", timeout=1).json()

@app.route("/order/")
def get_order(order_id):
    user_data = fetch_user(order_id)
    logistics_data = fetch_logistics(order_id)
    coupon_data = fetch_coupon(order_id)
    return jsonify({
        "user": user_data,
        "logistics": logistics_data,
        "coupon": coupon_data
    })

After:异步并发版(80行)

# app.py - 异步版本
import asyncio
import aiohttp
from flask import Flask, jsonify

app = Flask(__name__)
semaphore = asyncio.Semaphore(100)  # 控制最大并发100

async def fetch_json(session, url, timeout=0.5):
    """带超时和信号量的异步GET"""
    async with semaphore:
        try:
            async with session.get(url, timeout=aiohttp.ClientTimeout(total=timeout)) as resp:
                return await resp.json()
        except (asyncio.TimeoutError, aiohttp.ClientError):
            return {"error": "service_unavailable"}  # 降级返回

async def fetch_all(order_id):
    """并发获取所有下游数据"""
    timeout = aiohttp.ClientTimeout(total=1.0)
    async with aiohttp.ClientSession(timeout=timeout) as session:
        user_task = fetch_json(session, f"http://user-service/user/{order_id}")
        logistics_task = fetch_json(session, f"http://logistics-service/track/{order_id}")
        coupon_task = fetch_json(session, f"http://coupon-service/coupon/{order_id}")
        results = await asyncio.gather(user_task, logistics_task, coupon_task)
        return results

@app.route("/order/")
def get_order(order_id):
    # 关键:在同步WSGI中运行协程
    user_data, logistics_data, coupon_data = asyncio.run(fetch_all(order_id))
    return jsonify({
        "user": user_data,
        "logistics": logistics_data,
        "coupon": coupon_data
    })

五、踩坑与优化:三个必须说的坑

坑1:Flask + asyncio.run 的事件循环冲突
第一次测试时,我用loop = asyncio.new_event_loop()手动管理,结果每次请求都报“Event loop is closed”。原因是Flask多线程处理请求时,每个线程都尝试创建自己的事件循环。解决方案很简单:直接asyncio.run(),它会自动创建新循环并关闭旧循环。但注意:如果后续要升级到FastAPI或Starlette,就不要这么写了,直接用asyncio.get_running_loop()

坑2:aiohttp连接池泄漏
如果每次请求都新建ClientSession,压测时会出现Connection pool is full。正确做法:把ClientSession定义为全局单例,用async with管理生命周期。但全局session需要在应用启动时创建,在请求中复用。我用了模块级变量:

session = None

def get_session():
    global session
    if session is None:
        session = aiohttp.ClientSession()
    return session

坑3:超时设置要双重保险
aiohttp.ClientTimeout(total=1.0)是总超时,但TCP连接建立超时默认是5秒。如果下游服务宕机,总超时可能不起作用。必须显式设置connect=0.5

timeout = aiohttp.ClientTimeout(total=1.0, connect=0.5)

六、效果数据:压测对比

用wrk压测,50并发跑60秒,结果如下:

指标 同步版本 异步版本 提升倍数
QPS 37 312 8.4x
P99延迟 4.8s 410ms 11.7x
平均延迟 1.2s 320ms 3.75x
错误率 0% 0.3%(超时降级) -

注意错误率:异步版本有0.3%的请求因为超时降级返回了错误数据。这是预期的——我把超时设为500ms,但下游偶尔有1s以上的毛刺。降级方案是:返回空对象并记录日志,前端根据error字段提示“部分信息加载失败”

七、总结:什么时候该用asyncio?

不是所有场景都适合asyncio。我的经验是:

  • IO密集型(大量网络请求、数据库查询)→ 必须用asyncio,效果立竿见影
  • CPU密集型(计算、加密、序列化)→ 用多进程,asyncio反而拖慢速度
  • 混合负载 → 用asyncio + concurrent.futures.ThreadPoolExecutor组合

这次改造最值的点在于:只改了IO层,业务逻辑零变化。如果你也在用Flask + Requests做聚合API,建议先试试asyncio + aiohttp。如果追求更高性能,可以直接迁移到FastAPI(原生支持async def),但那是另一个故事了。

最后提醒:压测数据仅供参考,真实环境的网络波动、下游容量都会影响结果。改造前一定要先压测,改造后一定要灰度发布。