1. 问题背景:一个串行调用3个下游接口的报表服务
我们的订单报表接口 /api/orders/summary 需要聚合三个数据源:
- 用户服务(user-service):返回用户等级、地域(平均耗时800ms)
- 订单服务(order-service):返回订单金额、数量(平均耗时700ms)
- 库存服务(inventory-service):返回库存预警信息(平均耗时600ms)
原实现用 requests 库同步调用,代码直白,但生产环境一压测就露馅:
# before.py - 同步版本
import requests
from flask import Flask, jsonify
app = Flask(__name__)
@app.route('/api/orders/summary')
def summary():
user = requests.get('http://user-service/api/user', timeout=3).json()
order = requests.get('http://order-service/api/order', timeout=3).json()
stock = requests.get('http://inventory-service/api/stock', timeout=3).json()
return jsonify({
'user_level': user['level'],
'order_amount': order['amount'],
'stock_warning': stock['warning']
})
性能实测(wrk -t4 -c200 -d30s):
- 平均延迟:2150ms
- P99:3100ms
- 吞吐率:180 req/s
- CPU:85%(GIL锁竞争 + 线程池切换开销)
三个请求完全独立,却串行等待,纯粹是代码习惯问题。
2. 环境与版本:明确依赖,避免玄学
- Python 3.10.12(原生支持
asyncio.run(),不需要nest_asyncio) - Flask 2.3.3(注意:Flask 是 WSGI 同步框架,不能直接在视图函数里
await) - aiohttp 3.9.1(替代 requests 的异步客户端)
- gunicorn 21.2.0(生产部署,worker 模式改为
gevent或gthread,见下文坑)
3. 方案设计:用 asyncio 并发等待,而不是并行计算
核心思路:把三个 HTTP 调用变成协程,用 asyncio.gather 并发执行。但有两个关键约束:
- Flask 视图函数是同步的,不能直接
await。解决方案:用asyncio.run()包装一个异步函数,但注意不能在每个请求里都asyncio.run()(会重复创建事件循环,性能差)。 - 必须限流:如果下游服务扛不住并发,直接把对方打挂。用
asyncio.Semaphore(10)限制最大并发数为10。
架构图(文字版):
Flask 同步视图
↓ asyncio.run()
异步主函数
↓ 创建 Semaphore(10)
↓ asyncio.gather(协程1, 协程2, 协程3)
↓ 每个协程内部: async with semaphore: await aiohttp.get()
4. 核心实现:异步重构代码
第一步:初始化 aiohttp ClientSession(全局复用,避免每次请求创建连接池)
# async_client.py
import asyncio
import aiohttp
# 全局复用 session,连接池默认 100 个连接
session = None
async def init_session():
global session
if session is None:
# connector 调大连接池,TCPConnector 默认 limit=100
connector = aiohttp.TCPConnector(limit=100, ttl_dns_cache=300)
session = aiohttp.ClientSession(connector=connector)
async def get_json(url, semaphore):
async with semaphore:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=3)) as resp:
return await resp.json()
第二步:重写 Flask 视图函数
# after.py - 异步版本
import asyncio
from flask import Flask, jsonify
from async_client import init_session, get_json, session
app = Flask(__name__)
# 全局信号量:限制最大10个并发下游请求
semaphore = asyncio.Semaphore(10)
def async_summary():
"""异步聚合三个服务数据"""
async def _fetch_all():
await init_session() # 确保 session 已创建
urls = [
'http://user-service/api/user',
'http://order-service/api/order',
'http://inventory-service/api/stock'
]
# gather 并发执行,return_exceptions=True 避免单个失败拖垮全部
results = await asyncio.gather(
*(get_json(url, semaphore) for url in urls),
return_exceptions=True
)
# 解析结果,处理异常情况
user, order, stock = results
if isinstance(user, Exception):
user = {'level': 'N/A'}
if isinstance(order, Exception):
order = {'amount': 0}
if isinstance(stock, Exception):
stock = {'warning': 'unknown'}
return {
'user_level': user['level'],
'order_amount': order['amount'],
'stock_warning': stock['warning']
}
# 每个请求创建新的事件循环(注意:Flask 多线程下必须这样)
return asyncio.run(_fetch_all())
@app.route('/api/orders/summary')
def summary():
return jsonify(async_summary())
生产部署注意:asyncio.run() 每次创建新事件循环,如果 init_session() 反复执行,会警告 unclosed session。解决:把 init_session() 放在模块导入时做一次,但 Flask 的 debug 模式会重载模块,所以用全局 if session is None 判断即可。
5. 踩坑与优化:三个真实遇到的坑
坑1:asyncio.run() 在 gunicorn 多 worker 下的坑
- 现象:部署后报 RuntimeError: asyncio.run() cannot be called from a running event loop
- 原因:gunicorn 默认 sync worker 是单线程的,没问题。但如果你用了 gevent worker,gevent 会 patch 标准库,导致 asyncio 检测到已存在事件循环。
- 解决:gunicorn worker 类改用 gthread(线程池模式),配置 --worker-class gthread --threads 4。或者干脆不用 gunicorn,用 uvicorn 跑 ASGI 框架。但我为了兼容现有 Flask 代码,选择了 gthread。
坑2:aiohttp 连接池耗尽
- 现象:压测时大量 TimeoutError: Connect timeout,查看下游服务 CPU 正常。
- 原因:aiohttp 默认 TCPConnector(limit=100),但我们的信号量是10,按理说不会超过10个并发连接。排查发现是 DNS 解析阻塞——aiohttp 默认用线程池解析 DNS,高并发下线程池满了。
- 解决:TCPConnector(ttl_dns_cache=300),缓存DNS 5分钟。
坑3:asyncio.Semaphore 作用域
- 如果信号量定义在函数内部,每个请求会创建新的信号量,限流形同虚设。必须定义为模块级全局变量。
6. 效果数据:压测对比
用 wrk -t4 -c200 -d60s http://localhost:5000/api/orders/summary 压测,结果:
| 指标 | 同步版本 | 异步版本 | 提升 |
|---|---|---|---|
| 平均延迟 | 2150ms | 420ms | 5.1x |
| P99延迟 | 3100ms | 780ms | 3.9x |
| 吞吐率 | 180 req/s | 920 req/s | 5.1x |
| CPU占用 | 85% | 55% | -35% |
| 错误率 | 0.2% | 0.1% | - |
额外验证:把下游服务模拟延迟提高到2秒(用 time.sleep(2)),同步版本延迟直接飙到6.1秒(串行),异步版本仍然稳定在2.1秒(并行),说明并发红利完全取决于下游延迟。
7. 总结与进一步优化建议
- 适用场景:I/O密集型(HTTP调用、数据库查询、文件读取),不适用于CPU密集型任务(asyncio不解决GIL问题)。
- 不要过度设计:如果下游只有1个接口,同步就够。3个以上独立调用才值得改造。
- 后续可做:
- 用
asyncio.Queue做请求合并(类似 Hystrix 的 request collapsing),批量调用下游。 - 超时控制用
asyncio.wait_for,替代 aiohttp 的timeout参数,更精细。 - 如果想完全摆脱 Flask 同步限制,可以直接迁移到
FastAPI(原生 async 支持),但迁移成本需评估。
最后说一句:异步编程不是银弹,但它能让你的服务在高并发下从「勉强能跑」变成「游刃有余」。关键是把业务拆分成可并发的子任务,然后用最简单的 gather + Semaphore 组合拳,别一上来就上消息队列。