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等待,再动手。别为了用异步而异步,否则代码复杂度带来的维护成本可能超过性能收益。
(完)