一、问题背景:同步阻塞让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()驱动协程。
三、方案设计:两层并发模型
核心思路分三层:
- 协程级并发:三个外部API调用用
asyncio.gather()并发执行,将串行600ms压到并行180ms - 信号量限流:为每个下游服务创建独立的
asyncio.Semaphore(50),防止突发流量打垮第三方 - 线程池兜底: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的gevent或eventlet 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.Semaphore和httpx.AsyncClient如果定义在函数外部,会绑定到第一个事件循环,后续线程调用直接报错。
解法:把Semaphore和AsyncClient全部移到process_batch内部创建,用完后自动销毁。虽然有点浪费,但安全。
5.2 坑2:信号量放错位置导致死锁
我最初把Semaphore放在fetch_order_detail内部,结果当batch大小超过Semaphore容量时,外层gather等待内层信号量释放,而内层在等待外层gather完成——死锁。信号量必须放在每个外部调用的最外层。
5.3 坑3:asyncio.TimeoutError不会自动取消子任务
用asyncio.wait_for或asyncio.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。但要注意两点:
- 不是所有场景都适合:如果下游服务本身是CPU密集或数据库慢查询,异步无济于事。必须先用cProfile确认瓶颈在IO等待。
- 监控更复杂: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