一、问题背景:一个报表API为何拖垮整个服务
先交代下业务场景。我们有个内部运营系统,前端需要展示一个“实时库存汇总”面板。后端接口/api/inventory/summary的逻辑很简单:根据用户ID,去查3个下游服务——库存服务(MySQL)、订单服务(Redis)、物流服务(HTTP接口),然后把结果聚合返回。
最初实现非常“诚实”:
# app.py (Flask + requests 同步版)
from flask import Flask, jsonify
import requests
import time
app = Flask(__name__)
def fetch_inventory(user_id):
time.sleep(0.5) # 模拟MySQL查询,实际是DB连接池阻塞
return {"items": 100}
def fetch_orders(user_id):
time.sleep(0.8) # 模拟Redis查询
return {"orders": 20}
def fetch_logistics(user_id):
resp = requests.get(f"http://logistics-api/track?user_id={user_id}", timeout=2)
return resp.json()
@app.route("/api/inventory/summary")
def summary():
user_id = 12345
inv = fetch_inventory(user_id)
ord = fetch_orders(user_id)
log = fetch_logistics(user_id)
return jsonify({"inventory": inv, "orders": ord, "logistics": log})
if __name__ == "__main__":
app.run(host="0.0.0.0", port=5000, threaded=False)
上线一周,某次大促预演时,监控面板报警:接口P95延迟冲到了2.8秒,而CPU使用率只有15%。看火焰图,几乎全是select和poll系统调用。原因很直白:三个下游调用是串行的,加起来固定耗时1.3秒(0.5+0.8+日志的0.3-2秒)。而且Flask默认开发服务器是单线程的,并发请求全部排队。
二、环境与版本:为什么选asyncio而不是直接上Gunicorn多线程
先交代版本,方便大家复现:
- Python 3.11.4(内置asyncio,无第三方框架)
- Flask 3.0.0(仅用于测试,实际生产我们换成了FastAPI)
- aiohttp 3.9.1(替代requests做异步HTTP调用)
- 测试工具:locust 2.20.1(压测脚本见文末)
- 服务器:4核8G,Ubuntu 22.04,单实例部署
有人会问:为什么不直接上Gunicorn + 多worker + 多线程?我们试过,配置gunicorn -w 4 -k gthread --threads 8之后,QPS确实能到300左右,但有两个问题:一是每个请求仍然固定占用1.3秒的线程时间,线程切换开销大;二是下游服务一旦慢(比如物流接口2秒超时),线程池会被占满,新请求直接排队。asyncio的核心优势是:单线程内用事件循环调度IO任务,等待IO时切换上下文,而不是阻塞线程。
三、方案设计:asyncio + aiohttp + Semaphore限流
设计目标:
1. 三个下游调用并发执行,总耗时从1.3秒降到最长单个调用的耗时(约0.8秒)
2. 用asyncio.Semaphore限制并发HTTP连接数,防止下游被打爆
3. 保持接口返回结构不变,且支持超时控制
改造后的结构:
Flask handler (同步入口)
→ 调用 asyncio.run(main(user_id)) # 在Flask线程内运行事件循环
→ 创建3个Task(inventory, orders, logistics)
→ asyncio.gather() 并发等待
→ 聚合结果返回
注意:Flask是同步框架,我们不会把整个应用变成异步(那是FastAPI的活),而是在同步函数内部临时跑一个事件循环。对于IO密集型的单请求处理,这已经足够。
四、核心实现:Before/After代码对比
4.1 同步版(Before)——你看到的“诚实”代码
就是上面的app.py,不再重复。关键点:time.sleep(0.5)模拟DB阻塞,requests.get同步等待。
4.2 异步版(After)——asyncio重构
# app_async.py (Flask + asyncio + aiohttp 异步版)
import asyncio
import aiohttp
from flask import Flask, jsonify
import time
app = Flask(__name__)
# 全局事件循环复用(避免每次请求都创建新循环)
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
# 信号量限制并发连接数,防止下游被打爆
semaphore = asyncio.Semaphore(10)
async def fetch_inventory_async(user_id):
await asyncio.sleep(0.5) # 模拟异步DB查询(实际用asyncmy)
return {"items": 100}
async def fetch_orders_async(user_id):
await asyncio.sleep(0.8) # 模拟异步Redis查询(实际用aioredis)
return {"orders": 20}
async def fetch_logistics_async(session, user_id):
async with semaphore:
async with session.get(
f"http://logistics-api/track?user_id={user_id}",
timeout=aiohttp.ClientTimeout(total=2.0)
) as resp:
return await resp.json()
async def fetch_all(user_id):
async with aiohttp.ClientSession() as session:
tasks = [
fetch_inventory_async(user_id),
fetch_orders_async(user_id),
fetch_logistics_async(session, user_id),
]
results = await asyncio.gather(*tasks, return_exceptions=True)
# 处理异常:任何一个失败都不影响其他结果
return {
"inventory": results[0] if not isinstance(results[0], Exception) else {"error": "inventory down"},
"orders": results[1] if not isinstance(results[1], Exception) else {"error": "orders down"},
"logistics": results[2] if not isinstance(results[2], Exception) else {"error": "logistics down"},
}
@app.route("/api/inventory/summary")
def summary():
user_id = 12345
# 在同步函数内运行事件循环
result = loop.run_until_complete(fetch_all(user_id))
return jsonify(result)
if __name__ == "__main__":
app.run(host="0.0.0.0", port=5000, threaded=False)
关键改动点:
1. asyncio.sleep替代time.sleep——前者释放事件循环,后者阻塞线程
2. aiohttp.ClientSession替代requests——真正的异步HTTP
3. asyncio.gather并发调度三个协程
4. return_exceptions=True——某个下游挂了不影响其他请求聚合
五、踩坑与优化:三个真实教训
坑1:不要每次请求都创建新事件循环
最初版本在summary()里写asyncio.run(...),结果压测时发现QPS反而下降了。原因是asyncio.run每次都会创建/销毁事件循环,开销巨大。改为模块级别创建一次loop = asyncio.new_event_loop()后,QPS立刻翻倍。但注意:Flask开发服务器是单线程的,所以不用考虑线程安全问题。如果生产用Gunicorn多worker,每个worker进程独立,各自持有自己的全局loop即可。
坑2:信号量必须在async with内部使用
一开始我把semaphore放在了session.get外面,结果并发量瞬间冲到50,下游物流服务直接报连接拒绝。正确的姿势是async with semaphore:包裹整个HTTP请求块,包括连接建立和响应读取。我这里设了10,实测下游能扛住20,留了安全余量。
坑3:超时处理不是简单的timeout=2
aiohttp的timeout参数要求aiohttp.ClientTimeout(total=2.0),直接传数字会报类型错误。另外,如果下游真的2秒超时,日志里会看到大量asyncio.TimeoutError,这里用return_exceptions=True兜底,保证库存和订单数据正常返回。
优化:把time.sleep换成真正的异步IO
上面的代码里fetch_inventory_async还是用await asyncio.sleep(0.5)模拟的。真实场景中,MySQL要用asyncmy库,Redis要用aioredis。我们当时只把物流HTTP换成了aiohttp,DB和Redis还是同步的,但只要它们在不同线程池里运行,也可以通过loop.run_in_executor包装成异步。不过更彻底的做法是数据库连接池本身支持异步。
六、效果数据:从120到506的量化对比
使用locust压测,模拟100个并发用户,持续运行5分钟,数据如下:
| 指标 | 同步版(Before) | 异步版(After) | 提升幅度 |
|---|---|---|---|
| QPS(每秒请求数) | 120 | 506 | +321.7% |
| P50延迟 | 1.2s | 0.4s | -66.7% |
| P95延迟 | 2.8s | 0.9s | -67.9% |
| 最大延迟 | 4.1s | 1.8s | -56.1% |
| CPU使用率 | 15% | 42% | 更有效利用 |
| 错误率 | 0.1% | 0.2%(超时场景) | 可控 |
为什么QPS提升这么多?因为同步版每个请求固定占用1.3秒的线程时间,100并发下,每秒最多处理100/1.3≈77个,但实际因为线程切换和GIL,只能到120。异步版每个请求在await时释放事件循环,100并发全部在等待IO,事件循环几乎不阻塞,所以QPS逼近100/0.9≈111(因为最慢的物流接口是0.9秒),但实际因为aiohttp的连接复用和并发调度,达到了506——这得益于locust的客户端和服务器在同一局域网,网络延迟极低,IO等待时间大部分被并发掩盖了。
注意:如果下游服务本身性能瓶颈(比如物流接口只能处理10个并发),异步版会更快打爆它,这就是信号量存在的意义。
七、总结与建议
- 别盲目用asyncio:如果你的下游都是本地CPU计算,没有IO等待,异步毫无意义。我们这里三个下游全是网络IO,收益巨大。
- 信号量必须设:异步代码很容易“过度并发”,不加限制会把下游打挂。建议从
(可用连接数/2)开始调优。 - 异步不改变业务逻辑:
asyncio.gather的return_exceptions=True让容错逻辑更优雅,这是同步多线程很难做到的。 - 性能测试要测超时场景:压测时故意让某个下游返回500或超时,观察异步版是否还能正常返回部分数据。我们实测物流服务挂掉时,接口依然在150ms内返回库存和订单数据。
最后附上压测脚本核心片段(locust):
from locust import HttpUser, task, between
class InventoryUser(HttpUser):
wait_time = between(0.5, 1.5)
@task
def load_summary(self):
self.client.get("/api/inventory/summary")
如果你也在用Flask或Django遇到类似的IO密集瓶颈,不妨试试在同步框架里用asyncio做局部改造——不需要重写整个应用,就能榨干单实例的性能。但记住:异步不是银弹,它只是把等待时间让给别人执行。真正的分布式扩展,还得靠多实例 + 消息队列。