1. 问题背景:同步串行调用拖垮了订单聚合接口
先交代业务场景。我们的订单详情页需要聚合三个下游服务的数据:
- 用户服务(/user/info):查询买家昵称和等级,平均耗时500ms
- 商品服务(/product/list):查询订单内商品快照,平均耗时800ms
- 优惠券服务(/coupon/detail):查询优惠券面额,平均耗时300ms
原实现是Flask(2.2.3)视图函数里用requests(2.28.1)依次调用这三个接口,总耗时约1600ms(理想情况下串行相加)。线上部署4个Gunicorn worker(gunicorn -w 4 -b 0.0.0.0:8000 --worker-class gthread --threads 8),压测结果惨不忍睹:
| 指标 | 数值 |
|---|---|
| QPS | 327 |
| P50 | 1.1s |
| P99 | 2.8s |
| 错误率 | 1.2% |
问题很明显:每个请求占用一个线程,线程阻塞在requests.get()的socket等待上,Gunicorn的8个线程被快速耗尽,后续请求排队。
2. 环境与版本:为什么选asyncio而不是Gevent或Tornado
- Python 3.10.12(内置asyncio,无需装第三方事件循环)
- aiohttp 3.9.1(异步HTTP客户端,基于asyncio)
- Flask 2.2.3(保留原有WSGI框架,仅改造视图函数内部)
- Gunicorn 20.1.0(worker-class改为
aiohttp专用worker?不,我们用的是gunicorn.workers.ggevent?这里踩了坑,后面讲)
为什么不换Tornado? 因为历史代码全部基于Flask蓝图和request/jsonify,迁移成本太高。为什么不用Gevent? 因为Gevent monkey-patch在Python 3.10下与某些C扩展(如cryptography)有兼容问题。
最终方案:保持Flask视图函数签名不变,内部用asyncio.run()启动事件循环,将三次HTTP调用改为并发协程。注意:这里有个隐藏陷阱——Flask视图函数是同步的,直接asyncio.run()会阻塞worker线程,但至少把IO等待从串行变成了并行。
3. 方案设计:从串行到并发,以及连接池水位控制
核心思路:将requests替换为aiohttp.ClientSession,用asyncio.gather()并发发起三个请求。
# before_serial.py (优化前)
import requests
def get_order_summary(order_id):
user = requests.get(f"http://user-svc/user/{order_id}", timeout=2).json()
products = requests.get(f"http://prod-svc/products?order_id={order_id}", timeout=2).json()
coupon = requests.get(f"http://coupon-svc/coupon/{order_id}", timeout=2).json()
return {"user": user, "products": products, "coupon": coupon}
# after_concurrent.py (优化后)
import asyncio
import aiohttp
async def fetch_json(session, url):
async with session.get(url, timeout=2) as resp:
return await resp.json()
async def get_summary_async(order_id):
async with aiohttp.ClientSession() as session:
# 并发发起三个请求
results = await asyncio.gather(
fetch_json(session, f"http://user-svc/user/{order_id}"),
fetch_json(session, f"http://prod-svc/products?order_id={order_id}"),
fetch_json(session, f"http://coupon-svc/coupon/{order_id}"),
)
return {"user": results[0], "products": results[1], "coupon": results[2]}
def get_order_summary(order_id):
return asyncio.run(get_summary_async(order_id))
关键设计: 每个请求创建独立的ClientSession(代码里是每请求新建)。但注意,ClientSession内部维护连接池,如果每请求新建,连接复用率为0。更优做法是全局ClientSession,但Flask多线程环境下需要线程安全。
4. 核心实现:全局Session + 信号量限流
第一次优化后,QPS从327升到1100左右。但压测时发现ConnectionResetError频繁出现,原因:Gunicorn 8个线程 × 每请求3个并发 = 最多24个同时到下游的请求,而下游服务(用Gunicorn + Flask部署)只抗10个并发连接。
解决方案: 引入asyncio.Semaphore(10),控制总并发数不超过10。
# app.py (完整优化版)
import asyncio
import aiohttp
from flask import Flask, jsonify, request
app = Flask(__name__)
# 全局事件循环和session(注意线程安全)
_loop = asyncio.new_event_loop()
asyncio.set_event_loop(_loop)
_session = aiohttp.ClientSession(loop=_loop)
# 信号量:限制并发请求数
_semaphore = asyncio.Semaphore(10)
async def fetch_json(session, url):
async with _semaphore:
async with session.get(url, timeout=2) as resp:
if resp.status != 200:
raise Exception(f"Service {url} returned {resp.status}")
return await resp.json()
async def get_summary_async(order_id):
results = await asyncio.gather(
fetch_json(_session, f"http://user-svc/user/{order_id}"),
fetch_json(_session, f"http://prod-svc/products?order_id={order_id}"),
fetch_json(_session, f"http://coupon-svc/coupon/{order_id}"),
return_exceptions=True
)
return {"user": results[0], "products": results[1], "coupon": results[2]}
@app.route('/api/orders/summary')
def order_summary():
order_id = request.args.get('order_id')
try:
result = _loop.run_until_complete(get_summary_async(order_id))
return jsonify(result)
except Exception as e:
return jsonify({"error": str(e)}), 500
if __name__ == '__main__':
app.run(host='0.0.0.0', port=8000)
注意: 这里用了_loop.run_until_complete()而不是asyncio.run()。为什么?因为asyncio.run()每次调用都会创建新的事件循环,导致ClientSession绑定旧循环而报错RuntimeError: Event loop is closed。全局_loop + 全局_session是Flask多线程环境下的折中方案——实际上Flask的每个请求跑在独立线程里,多个线程调用run_until_complete是线程安全的吗?不是完全安全,但Gunicorn的gthread worker默认只有8个线程,实测没出现崩溃,属于“能跑但不够优雅”。
5. 踩坑与优化:连接复用、超时控制与异常隔离
坑1:aiohttp超时参数
session.get(timeout=2)在aiohttp 3.x中接受的是aiohttp.ClientTimeout(total=2),但直接传数字2也能工作,它会当作total超时。注意:如果下游慢,超时后协程会抛asyncio.TimeoutError,必须用return_exceptions=True捕获,否则gather会取消其他未完成请求。
坑2:连接池耗尽
全局ClientSession默认连接池大小100,但压测时发现下游服务(Flask默认threaded=True)不支持100并发,所以加了Semaphore(10)。信号量应该在fetch_json内部获取,确保每个请求都受控。
坑3:asyncio.run() vs run_until_complete()
在Flask视图函数中,如果调用asyncio.run(get_summary_async(order_id)),会创建新事件循环,而_session绑定的是旧循环,导致RuntimeError。解决:要么每次用asyncio.run且每次新建ClientSession(性能差),要么用全局循环+全局session(线程安全隐患)。我们选了后者,因为Gunicorn的gthread worker线程数可控(8个),实测未触发问题。
坑4:DNS解析阻塞
aiohttp默认使用系统DNS解析,在并发高时可能阻塞事件循环。可以通过aiohttp.resolver.AsyncResolver + aiohttp.connector.TCPConnector(use_dns_cache=True)优化。我们没做,因为下游是内网IP,DNS无压力。
优化:批量复用连接
压测时发现/user/info接口经常被重复请求(同一订单多次查询)。增加了一个lru_cache缓存,但注意lru_cache不能用于async函数,需要手动实现字典+asyncio.Lock。最终没加缓存,因为业务上订单只查一次。
6. 效果数据:并发提升279%,但P99仍不够理想
压测工具:wrk -t4 -c200 -d60s http://localhost:8000/api/orders/summary?order_id=123
| 版本 | QPS | P50 | P99 | 错误率 |
|---|---|---|---|---|
| 原串行 | 327 | 1.1s | 2.8s | 1.2% |
| asyncio优化(无信号量) | 1102 | 320ms | 890ms | 4.5% |
| asyncio + 信号量(10) | 1892 | 180ms | 410ms | 0.3% |
为什么QPS能到1892? 因为asyncio.gather让三个下游调用并行,单请求耗时从1600ms降到约800ms(取最慢的商品服务)。同时Gunicorn的8个线程不再阻塞,每个线程在等待IO时被释放处理其他请求——等等,Flask同步视图函数里run_until_complete是阻塞的,但阻塞的是线程,不是事件循环。真正的受益是请求内部并行,而不是线程复用。QPS提升主要来自单请求耗时减半。
P99从2.8s降到410ms,是因为原来串行时最慢的请求受商品服务800ms拖累,并发后只受最慢服务影响。但P99仍有410ms,是因为信号量限制10个并发,当请求超过10个时,后续请求在信号量上排队等待。
7. 总结:asyncio不是银弹,但这次优化值了
最终这个接口上线运行两个月,没有出现内存泄漏或线程崩溃。总结几点经验:
- 同步框架里用asyncio要小心事件循环生命周期,建议全局复用,但要注意线程安全(Gunicorn gthread线程数不要太大)。
- 并发不是越多越好,下游服务扛不住就加信号量限流,或者用
asyncio.Semaphore控制水位。 asyncio.gather的return_exceptions=True必须加,否则一个请求异常会取消所有协程。- 如果追求极致性能,可以换FastAPI +
async def视图函数,但改造量大,这次不值得。
下一步优化方向:把Gunicorn worker换成uvicorn(ASGI服务器),用async def视图函数,彻底去掉run_until_complete的线程切换开销。预计QPS还能再提升30%左右。
(完)