1. 问题背景:一个慢到被投诉的报表接口
上个月运维同事甩给我一个监控截图:“这个/api/v1/report/summary接口,P99延迟2.3秒,上游服务都超时了”。我查了下代码,这是个典型的同步串行调用:
# 改造前:同步串行调用,3个HTTP请求顺序执行
def get_summary_report(self, user_id):
# 调用订单服务
order_data = requests.get(f"http://order-service/api/orders?uid={user_id}", timeout=3).json()
# 调用用户服务
user_data = requests.get(f"http://user-service/api/users?uid={user_id}", timeout=3).json()
# 调用推荐服务
rec_data = requests.get(f"http://rec-service/api/recommend?uid={user_id}", timeout=3).json()
# 业务聚合...
return self._merge(order_data, user_data, rec_data)
三个服务调用是独立的,但这里串行执行,每次请求固定耗时 = 三个服务延迟之和。当时平均每个服务响应200-300ms,总耗时850ms左右。
2. 环境与版本:明确基线数据
先说下环境,方便你复现对比:
- Python 3.8.10(生产环境,后面会提到这个版本的坑)
- Flask 2.0.1 + gunicorn 20.1.0(4 workers, sync worker class)
- 上游服务:3个内部FastAPI服务,部署在K8s内网
- 压测工具:wrk 4.1.0,单线程100连接压测10秒
压测基线(改造前):
Running 10s test @ http://api-gateway:8080/api/v1/report/summary
100 threads and 100 connections
Thread Stats Avg Stdev Max +/- Stdev
Latency 1.21s 320.5ms 2.30s 85%
Req/Sec 120.5 15.3 180.0 75%
12208 requests in 10.09s, 3.12MB read
Requests/sec: 1210.0
这里注意:单worker QPS约120,多worker后整体吞吐上不去,瓶颈在CPU等待I/O(线程阻塞在requests上)。
3. 方案设计:异步化不是终极方案,要配合连接池
我的思路很直接:用asyncio把3个串行HTTP调用改成并发,同时用aiohttp替代requests。但有几个关键决策点:
- 框架选择:Flask是同步WSGI框架,不能直接用
async def路由。所以采用双层架构:Flask路由保持同步,内部调用asyncio.run()启动事件循环,或者用loop.run_in_executor把异步逻辑包装成线程执行。 - 连接池复用:
aiohttp.ClientSession必须全局复用,否则每次请求创建新连接池,性能反而更差。 - 超时控制:用
asyncio.wait_for给每个子任务加超时,避免单个服务慢拖垮整体。 - 限流:
asyncio.Semaphore控制并发连接数,防止打爆上游服务。
架构图(文字版):
Flask (sync) → run_in_executor → asyncio.run()
→ 创建aiohttp.ClientSession(connector=TCPConnector(limit=100))
→ asyncio.gather(3个coroutine, return_exceptions=True)
→ 聚合结果返回
4. 核心实现:asyncio重构代码
先说结论:不要直接用asyncio.run()在Flask路由里,因为每次调用会创建新事件循环,开销巨大。正确做法是全局事件循环 + run_coroutine_threadsafe,或者更简单点:用loop.run_in_executor把整个异步逻辑丢到默认线程池。
4.1 全局Session与辅助函数
# async_client.py
import asyncio
import aiohttp
from aiohttp import TCPConnector
# 全局连接池,limit=100表示最多100个并发连接
_connector = TCPConnector(limit=100, limit_per_host=20, ttl_dns_cache=300)
_session = aiohttp.ClientSession(connector=_connector)
async def fetch_json(session, url, semaphore, timeout=2.0):
"""带限流和超时的异步请求"""
async with semaphore:
try:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=timeout)) as resp:
if resp.status != 200:
return None
return await resp.json()
except asyncio.TimeoutError:
# 记录日志,返回降级数据
return None
except Exception as e:
return None
async def fetch_all(user_id):
sem = asyncio.Semaphore(50) # 限制并发50
urls = [
f"http://order-service/api/orders?uid={user_id}",
f"http://user-service/api/users?uid={user_id}",
f"http://rec-service/api/recommend?uid={user_id}",
]
results = await asyncio.gather(
*(fetch_json(_session, url, sem) for url in urls),
return_exceptions=True # 防止单个异常导致全部失败
)
return results # [order_data, user_data, rec_data] 或 None
4.2 Flask路由集成
# app.py
import asyncio
from flask import Flask, jsonify
from async_client import fetch_all
app = Flask(__name__)
# 全局事件循环,在worker启动时创建
loop = asyncio.new_event_loop()
@app.route('/api/v1/report/summary')
def get_summary():
user_id = request.args.get('uid')
if not user_id:
return jsonify({'error': 'missing uid'}), 400
# run_coroutine_threadsafe 是线程安全的
future = asyncio.run_coroutine_threadsafe(fetch_all(user_id), loop)
try:
results = future.result(timeout=5) # 整体超时5s
except asyncio.TimeoutError:
return jsonify({'error': 'upstream timeout'}), 504
# 业务聚合(略)
return jsonify(merge_data(results))
关键点:loop必须在gunicorn worker启动时创建,不能每次请求都建。用run_coroutine_threadsafe把协程提交到事件循环,主线程阻塞等待结果(future.result())。
4.3 gunicorn配置调整
因为异步逻辑跑在事件循环里,worker不需要多线程了,把worker数调成CPU核数×2,worker_class保持同步:
gunicorn -w 8 -k sync -b 0.0.0.0:8080 app:app
5. 踩坑与优化:三个大坑,每个都让我调了半天
5.1 Python 3.8的asyncio.run()坑
一开始我图省事,直接在Flask路由里写asyncio.run(fetch_all(user_id))。结果压测时发现每隔几分钟就报“Event loop is closed”。查了文档:Python 3.8的asyncio.run()每次创建新循环,且ClientSession绑定的是旧循环,下次调用时session关联的传输已经关闭。
解决方案:改用全局事件循环 + run_coroutine_threadsafe。升级Python 3.10后asyncio.run()支持复用,但生产环境一时半会儿升不了,就用全局循环方案。
5.2 aiohttp连接池耗尽
压测到500 QPS时,上游服务开始报“Connection reset by peer”。原因是TCPConnector(limit=100)默认是对所有域名共享100连接,但我们有3个上游服务,每个服务被限制到33个连接,不够用。
优化:改成limit=100, limit_per_host=50(每域名50),同时加上enable_cleanup_closed=True自动清理半开连接。实测后上游重试率从5%降到0.3%。
5.3 Semaphore的坑:放错位置
一开始Semaphore放在fetch_all()里,每次请求都创建新的信号量,等于没限流。应该全局创建,在fetch_json内部使用。上面代码里sem = asyncio.Semaphore(50)在fetch_all里创建,其实不对,应该提到模块级别:
_sem = asyncio.Semaphore(50) # 全局信号量
async def fetch_all(user_id):
urls = [...]
results = await asyncio.gather(
*(fetch_json(_session, url, _sem) for url in urls), # 用全局_sem
return_exceptions=True
)
别小看这个,一开始每个请求新建Semaphore,并发数完全不可控,直接打崩了上游。
6. 效果数据:改造前后对比
压测命令和参数完全一致(wrk -t100 -c100 -d10s),改造后结果:
Running 10s test @ http://api-gateway:8080/api/v1/report/summary
100 threads and 100 connections
Thread Stats Avg Stdev Max +/- Stdev
Latency 180.2ms 45.1ms 500.0ms 82%
Req/Sec 1800.5 210.3 2400.0 73%
18205 requests in 10.09s, 4.65MB read
Requests/sec: 1804.0
| 指标 | 改造前 | 改造后 | 提升 |
|---|---|---|---|
| 平均延迟 | 1.21s | 180ms | 6.7倍 |
| P99延迟 | 2.3s | 450ms | 5.1倍 |
| QPS | 120 | 1800 | 15倍 |
| 错误率 | 1.2%(超时) | 0.1%(降级) | 12倍 |
注意:QPS提升15倍不全是asyncio的功劳,还因为我把worker数从4调到8(用满CPU核)。但老代码即使8个worker,QPS也就200左右,因为每个worker只能处理1个请求,阻塞在I/O上。异步化后单worker就能同时处理几十个请求。
7. 总结与建议
这次改造的核心教训:
- 同步代码的并发瓶颈在I/O等待,asyncio把等待时间让给其他协程,吞吐量提升是数量级的。
- 别用
asyncio.run()在同步框架里,全局事件循环才是正道。 - 连接池和信号量是异步编程的命门,不控制迟早出事。
return_exceptions=True必须加,否则一个服务抖动就全挂。
如果你的接口是多个独立HTTP调用聚合,强烈建议试试这个方案。成本不高,收益巨大。但注意:如果业务逻辑里有CPU密集计算(比如大列表排序、正则匹配),别直接用asyncio,配合run_in_executor丢给进程池,否则事件循环会被阻塞。
最后的建议:生产环境升级Python到3.10+,asyncio.run()复用循环的体验好很多,还有官方asyncio.timeout()上下文管理器,比wait_for更优雅。