一、问题背景:同步阻塞让8核CPU集体“摸鱼”

事情要从一个订单导出接口说起。业务逻辑很简单:前端请求导出过去30天的订单,后端需要调用三个外部服务:用户服务(获取买家信息)、库存服务(获取SKU详情)、物流服务(获取轨迹),然后聚合数据。单个请求串行耗时约600ms,但线上QPS过200就开始超时。

py-spy dump查看线程栈,发现所有worker线程全部阻塞在requests.get()上,CPU利用率不足5%。典型的IO密集型瓶颈——8核CPU都在等网络IO,资源浪费严重。

二、环境与版本:Python 3.11 + 混合异步方案

改造前环境:
- Python 3.8.10(无法用asyncio.TaskGroup
- Flask 2.2.2 (同步WSGI)
- requests 2.28.1

改造后环境:
- Python 3.11.4(享受asyncio.timeout和异常组)
- Flask 2.3.2 + gunicorn 20.1.0(同步worker)
- httpx 0.24.1(支持HTTP/2和异步)
- asyncio 内置库

关键决策:没有用aiohttp替换Flask,因为Flask生态成熟,且我们团队不打算全量迁到FastAPI。方案是同步入口 + 异步核心——在Flask视图函数里用asyncio.run()驱动协程。

三、方案设计:两层并发模型

核心思路分三层:

  1. 协程级并发:三个外部API调用用asyncio.gather()并发执行,将串行600ms压到并行180ms
  2. 信号量限流:为每个下游服务创建独立的asyncio.Semaphore(50),防止突发流量打垮第三方
  3. 线程池兜底:Flask的每个请求是独立线程,我们只把IO密集部分丢给事件循环,CPU密集(如JSON序列化)保留在worker线程

架构图:

graph LR
A[客户端] --> B[gunicorn worker线程]
B -->|asyncio.run| C[事件循环]
C --> D[Semaphore(50)] --> E[httpx调用服务A]
C --> F[Semaphore(30)] --> G[httpx调用服务B]
C --> H[Semaphore(80)] --> I[httpx调用服务C]
E & G & I --> J[gather聚合]
J --> B

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

4.1 改造前:纯同步串行代码

# before.py
import requests
from flask import Flask, jsonify

app = Flask(__name__)

def fetch_user(order):
    # 模拟外部API,实际耗时200ms
    resp = requests.get(f"http://user-svc/{order['uid']}", timeout=2)
    return resp.json()

def fetch_stock(order):
    resp = requests.get(f"http://stock-svc/{order['sku']}", timeout=2)
    return resp.json()

def fetch_logistics(order):
    resp = requests.get(f"http://logi-svc/{order['oid']}", timeout=2)
    return resp.json()

@app.route('/export')
def export():
    orders = get_orders_from_db()  # 假设已实现
    result = []
    for o in orders[:10]:  # 每页10条
        user = fetch_user(o)      # 200ms
        stock = fetch_stock(o)    # 200ms
        logi = fetch_logistics(o) # 200ms
        result.append({...})
    return jsonify(result)

4.2 改造后:asyncio + 信号量 + 超时控制

# after.py (Python 3.11)
import asyncio
import httpx
from flask import Flask, jsonify

app = Flask(__name__)

# 每个下游独立的信号量,防止单点打爆
USER_SEM = asyncio.Semaphore(50)
STOCK_SEM = asyncio.Semaphore(30)
LOGI_SEM = asyncio.Semaphore(80)

async def fetch_user(client, order):
    async with USER_SEM:
        resp = await client.get(f"http://user-svc/{order['uid']}", timeout=2.0)
        resp.raise_for_status()
        return resp.json()

async def fetch_stock(client, order):
    async with STOCK_SEM:
        resp = await client.get(f"http://stock-svc/{order['sku']}", timeout=2.0)
        return resp.json()

async def fetch_logistics(client, order):
    async with LOGI_SEM:
        resp = await client.get(f"http://logi-svc/{order['oid']}", timeout=2.0)
        return resp.json()

async def fetch_order_detail(client, order):
    # 关键点:三个请求并发执行,总耗时约等于最慢的一个
    user_task = asyncio.create_task(fetch_user(client, order))
    stock_task = asyncio.create_task(fetch_stock(client, order))
    logi_task = asyncio.create_task(fetch_logistics(client, order))

    # Python 3.11可以用asyncio.timeout,老版本用wait_for
    try:
        async with asyncio.timeout(3.0):  # 整体超时3秒
            user, stock, logi = await asyncio.gather(
                user_task, stock_task, logi_task
            )
            return merge_data(order, user, stock, logi)
    except asyncio.TimeoutError:
        # 超时任务需要取消,否则会泄漏
        for t in (user_task, stock_task, logi_task):
            t.cancel()
        raise

async def process_batch(orders):
    async with httpx.AsyncClient(http2=True, limits=httpx.Limits(max_connections=200)) as client:
        # 限制并发批次,控内存
        sem = asyncio.Semaphore(20)
        async def wrapped(o):
            async with sem:
                return await fetch_order_detail(client, o)
        return await asyncio.gather(*[wrapped(o) for o in orders])

@app.route('/export')
def export():
    orders = get_orders_from_db()
    # 关键:asyncio.run创建独立事件循环,避免线程间共享循环
    result = asyncio.run(process_batch(orders[:10]))
    return jsonify(result)

4.3 部署配置:gunicorn的坑

这里有个重要坑:千万不要用gunicorn的geventeventlet worker配合asyncio,事件循环会被Greenlet切换搞乱。我用的是同步worker,每个worker内部自己跑asyncio:

# gunicorn.conf.py
workers = 4          # 4个进程,每个进程8线程
threads = 8          # 每个进程8个同步线程
worker_class = 'sync' # 必须用sync,不能用gevent
timeout = 30
keepalive = 5

这样配置下,单进程最多8个并发请求同时进入asyncio.run(),每个事件循环内可处理50+并发IO。总并发能力 = 4进程 × 8线程 × 50信号量 = 1600个并发IO等待。

五、踩坑与优化:三个血泪教训

5.1 坑1:asyncio.run()在线程中的重复创建

Flask的每个worker线程会多次调用asyncio.run()。每次调用都会创建新的事件循环,而asyncio.Semaphorehttpx.AsyncClient如果定义在函数外部,会绑定到第一个事件循环,后续线程调用直接报错。

解法:把Semaphore和AsyncClient全部移到process_batch内部创建,用完后自动销毁。虽然有点浪费,但安全。

5.2 坑2:信号量放错位置导致死锁

我最初把Semaphore放在fetch_order_detail内部,结果当batch大小超过Semaphore容量时,外层gather等待内层信号量释放,而内层在等待外层gather完成——死锁。信号量必须放在每个外部调用的最外层

5.3 坑3:asyncio.TimeoutError不会自动取消子任务

asyncio.wait_forasyncio.timeout超时后,子任务不会自动取消,它们会继续在后台运行直到完成。这在高并发下会造成任务泄漏。必须在except块中显式cancel。

5.4 性能优化:HTTP/2多路复用

启用httpx.AsyncClient(http2=True)后,同一个TCP连接可以并发传输多个请求,减少了TLS握手开销。实测在100并发下,连接复用率提升了40%。

六、效果数据:性能对比

使用locust压测,10分钟持续负载,数据如下:

指标 改造前(同步) 改造后(asyncio) 提升倍数
QPS (50并发) 187 1520 8.1x
P99延迟 2.31s 340ms 6.8x
平均延迟 890ms 190ms 4.7x
错误率 4.2% 0.1% -
CPU利用率 5% 35% 7x

压测环境:8核C5.xlarge,100Mbps内网,下游API模拟延迟200ms±30ms。

关键发现:QPS提升并非线性,因为asyncio将3次串行IO变成并行,理论最大提升3倍;额外的5倍提升来自HTTP/2连接复用和更高效的调度。当并发从50升到500时,同步版本直接雪崩(错误率>30%),异步版本只是延迟线性增加。

七、总结与后续改进

这次改造的核心收益不是“异步魔法”,而是把等待时间还给CPU。但要注意两点:

  1. 不是所有场景都适合:如果下游服务本身是CPU密集或数据库慢查询,异步无济于事。必须先用cProfile确认瓶颈在IO等待。
  2. 监控更复杂:asyncio的调用栈在py-spy里看是“扁平”的,需要自己打点追踪。

后续计划:
- 将Flask迁移到FastAPI,直接用原生异步路由,去掉asyncio.run()的桥接开销(预计再提升15%)
- 引入asyncio.TaskGroup(Python 3.11)替代gather,自动管理异常传播
- 为每个下游服务配置独立的超时和重试策略,避免整体超时误伤

如果你也遇到“CPU空闲但接口慢”的问题,先别急着加机器,试试这个改造思路。有问题欢迎评论区交流。


参考配置:Python 3.11.4, httpx 0.24.1, gunicorn 20.1.0, Linux 5.15 kernel