1. 问题背景:一次线上事故引发的重构

周三下午,运营反馈“订单导出一直转圈”。我打开监控面板:API平均延迟2.8s,P99直接飙到5s+,但CPU利用率只有35%,内存正常。直觉告诉我——线程被IO卡住了。

看代码,逻辑很简单:接收订单ID列表 → 逐个请求第三方物流接口获取轨迹 → 聚合后返回。同步代码大概长这样:

def export_orders(order_ids):
    results = []
    for oid in order_ids:
        # 每个请求等待2秒(第三方接口P99)
        trace = requests.post(LOGISTICS_URL, json={"oid": oid}, timeout=3)
        results.append(parse(trace))
    return results

假设100个订单,最坏情况需要200秒。虽然用了线程池(max_workers=8),但每个线程依然阻塞等待IO。线程切换开销 + GIL限制,导致吞吐量上不去。这是典型的“IO密集型任务用错了并发模型”。

2. 环境与版本:选型说明

  • Python 3.10.9(asyncio在3.10后稳定,没有3.11的TaskGroup但够用)
  • Flask 2.2.3(WSGI框架,后续会说明为什么不能直接跑async)
  • httpx 0.24.1(支持async/await的HTTP客户端,替代requests)
  • gunicorn 20.1.0(uvicorn等ASGI服务器暂时不用,团队运维栈是gunicorn)
  • 部署:4核8G云主机,单实例

关键决策:不用FastAPI(团队没时间迁移路由层),保留Flask,只在内部逻辑用asyncio。

3. 方案设计:三层异步化

第一层:IO并发化。把所有requests.post换成httpx.AsyncClient,用asyncio.gather并发发起请求。

第二层:信号量限流。第三方API有并发限制(我们测过超过50并发会触发限流),用asyncio.Semaphore(30)控制最大并发数。

第三层:事件循环生命周期管理。Flask是同步框架,每个请求在独立线程中执行。我们需要在每个请求线程内创建独立的事件循环,处理完关闭,避免跨线程共享loop。

架构图(文字版):
Flask View → 同步函数 → 创建新事件循环 → 异步子协程(信号量控制)→ httpx异步请求 → 聚合结果 → 返回

4. 核心实现:改造前后的完整代码

Before(同步版,部分代码)

from flask import Flask, request, jsonify
import requests

app = Flask(__name__)

@app.route('/api/export', methods=['POST'])
def export():
    order_ids = request.json['order_ids']
    results = []
    for oid in order_ids:
        # 同步阻塞,每次耗时约2秒
        resp = requests.post(
            "https://logistics.example.com/trace",
            json={"order_id": oid},
            timeout=3
        )
        results.append(resp.json()["data"])
    return jsonify({"data": results})

After(异步优化版,完整可运行)

import asyncio
import httpx
from flask import Flask, request, jsonify

app = Flask(__name__)

# 全局httpx异步客户端,连接池复用
client = httpx.AsyncClient(
    base_url="https://logistics.example.com",
    timeout=httpx.Timeout(3.0, connect=1.0),
    limits=httpx.Limits(max_connections=50, max_keepalive_connections=20)
)

# 信号量:控制第三方API并发数,避免触发限流
semaphore = asyncio.Semaphore(30)

async def fetch_trace(order_id: str) -> dict:
    """异步获取单个订单轨迹,带信号量保护"""
    async with semaphore:
        payload = {"order_id": order_id}
        resp = await client.post("/trace", json=payload)
        resp.raise_for_status()
        data = resp.json()
        return {"order_id": order_id, "trace": data["data"]}

async def async_main(order_ids: list) -> list:
    """并发聚合所有订单数据"""
    tasks = [fetch_trace(oid) for oid in order_ids]
    # gather默认返回顺序与输入一致,ensure_async包装
    results = await asyncio.gather(*tasks, return_exceptions=True)
    # 过滤异常,记录错误日志
    valid_results = []
    for r in results:
        if isinstance(r, Exception):
            # 生产环境应记录到日志中心,这里简略处理
            valid_results.append({"error": str(r)})
        else:
            valid_results.append(r)
    return valid_results

@app.route('/api/export', methods=['POST'])
def export_async():
    order_ids = request.json['order_ids']
    if not order_ids:
        return jsonify({"error": "empty list"}), 400

    # 关键:在同步Flask线程中创建独立事件循环
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    try:
        results = loop.run_until_complete(async_main(order_ids))
    finally:
        loop.close()
        # 注意:不能set_event_loop(None),否则同线程后续请求会报错

    return jsonify({"data": results, "total": len(results)})

if __name__ == "__main__":
    app.run(threaded=True, processes=1)

5. 踩坑与优化:三个真实教训

坑1:Flask + asyncio的兼容性陷阱
不要直接在Flask view函数前加async def。Flask 2.2是WSGI框架,不支持ASGI异步请求处理。如果你写async def export(),Flask会把它当普通函数调用,返回一个coroutine对象,导致客户端收到TypeError: Object of type coroutine is not JSON serializable。解决方案就是我在代码里做的:同步函数内部用run_until_complete跑事件循环。

坑2:连接池耗尽问题
最初我没设置max_connections,默认httpcore连接池是100。压测时发现偶尔有ConnectionPoolTimeout。后来调整limits=httpx.Limits(max_connections=50),结合信号量30,确保并发请求数永远小于连接池上限。

坑3:循环事件循环的坑
第一次写的时候,我在async_main里直接用了asyncio.get_event_loop(),结果在gunicorn多worker模式下,同一线程可能复用一个已关闭的loop,报RuntimeError: Event loop is closed。修复:强制每次new_event_loop(),用完close()。如果追求极致性能,可以考虑每个worker只创建一次loop,但需要线程局部存储,代码复杂度增加,我选择了简单方案。

优化记录
- 第一版直接改async,QPS从120→700,但P99不稳定,因为没限流。
- 加信号量后,QPS稳定在1800,P99降到180ms。
- 加了连接池复用后,CPU占用从40%降到25%(减少TCP握手开销)。

6. 效果数据:改造前后的硬指标对比

测试环境:4核8G,gunicorn 4 workers,压测工具wrk(100并发,60秒)。

指标 同步版(Before) 异步版(After) 提升幅度
QPS 120 1850 15.4倍
平均延迟 820ms 54ms 93.4%下降
P99延迟 2300ms 180ms 92.2%下降
错误率(超时+5xx) 15% 0.3% 98%下降
CPU利用率 35% 25% 反而下降

数据解读:
- 同步版瓶颈在线程阻塞,4核8G开8线程,每线程串行等待2秒,QPS天花板=8/2=4,实测120是因为有keepalive和部分请求更快。
- 异步版用事件循环+httpx,IO等待期间CPU可以处理其他请求,整体吞吐量接近IO并发上限(受限于第三方API限流,30并发)。
- 延迟下降是因为并发请求同时发出,单个请求的等待时间从“排队等前面99个”缩短为“信号量队列平均等待时间+单次IO时间”。

7. 总结与适用边界

什么情况适合这个方案
- 内部逻辑有大量外部HTTP/RPC调用,且互相无依赖。
- 第三方API有并发限制,需要限流。
- 团队暂时不能迁移到FastAPI/Starlette,但想享受异步红利。

什么情况不适合
- 你的代码是CPU密集型(像图片处理、加解密),asyncio没用,得用多进程。
- 上游API响应极快(<5ms),并发提升有限,反而增加代码复杂度。
- 你需要WebSocket长连接,还是老老实实换ASGI服务器。

最后说一句:很多同学觉得asyncio难,其实核心就是事件循环 + 协程 + await。我这次改造最值钱的部分不是代码,而是理解了“异步不是魔法,是把等待时间让出去”。如果你也遇到类似问题,先做profiling,确认瓶颈确实是IO等待,再动手。别为了用异步而异步,否则代码复杂度带来的维护成本可能超过性能收益。

(完)