问题背景:一个接口拖垮了整个服务
上个月接手一个内部数据中台服务,其中一个/api/v1/report/aggregate接口成了众矢之的。这个接口的逻辑很简单:接收一个用户ID,然后并发调用三个下游服务——用户画像服务、订单统计服务、风控标签服务——最后聚合结果返回。
问题在于,原始实现用的是requests库,同步串行调用:
# 原始同步实现(简化版)
def aggregate_report(user_id: str):
profile = requests.get(f"http://profile-svc/user/{user_id}", timeout=3).json()
orders = requests.get(f"http://order-svc/stats/{user_id}", timeout=3).json()
risk = requests.get(f"http://risk-svc/tags/{user_id}", timeout=3).json()
return {"profile": profile, "orders": orders, "risk": risk}
三个下游服务平均响应时间分别是800ms、900ms、600ms,串行加起来就是2.3秒。而网关超时设置是3秒,一旦某个下游抖动到1.5秒以上,整个请求就超时。生产监控显示这个接口的P95稳定在2.3秒,P99直接3秒+,每天触发约200次超时告警。
环境与版本:明确技术栈边界
先交代一下重构时的环境,避免版本问题干扰:
- Python 3.10.12(关键:3.10之后
asyncio.run()成为标准用法,3.7之前的写法不要参考) - Flask 2.2.5(同步框架,但我们可以用
asgiref.sync.async_to_sync桥接) - aiohttp 3.8.4(异步HTTP客户端)
- 部署环境:Docker容器,2核CPU,4GB内存,单实例
核心决策:不换框架(FastAPI/Starlette),因为存量代码和运维体系都是Flask的,换框架成本太高。我们只需要在视图函数内部把同步IO换成异步IO,用asyncio.run()在每次请求时创建事件循环即可——简单粗暴,但极其有效。
方案设计:异步化 + 信号量限流 + 连接池
3.1 核心设计思路
- 全链路异步化:用
aiohttp替代requests,三个下游调用用asyncio.gather()并发执行 - 信号量限流:防止上游突发流量打爆下游服务,
Semaphore(50)限制最大并发数 - 连接池复用:
aiohttp.ClientSession全局复用,避免每次请求都重建连接(TCP握手开销巨大) - 超时控制:
asyncio.wait_for给每个子任务单独设置超时,避免一个慢服务拖垮整体
3.2 代码结构演进
第一版异步代码长这样:
import asyncio
import aiohttp
from flask import Flask, jsonify
app = Flask(__name__)
# 全局session,复用连接池
_session = aiohttp.ClientSession(
connector=aiohttp.TCPConnector(limit=100, ttl_dns_cache=300),
timeout=aiohttp.ClientTimeout(total=2)
)
async def fetch_json(session, url):
async with session.get(url) as resp:
return await resp.json()
async def aggregate_async(user_id: str):
tasks = [
fetch_json(_session, f"http://profile-svc/user/{user_id}"),
fetch_json(_session, f"http://order-svc/stats/{user_id}"),
fetch_json(_session, f"http://risk-svc/tags/{user_id}"),
]
results = await asyncio.gather(*tasks, return_exceptions=False)
return {
"profile": results[0],
"orders": results[1],
"risk": results[2]
}
@app.route('/api/v1/report/aggregate')
def aggregate_report():
user_id = request.args.get('user_id')
try:
result = asyncio.run(aggregate_async(user_id))
return jsonify(result)
except Exception as e:
return jsonify({"error": str(e)}), 500
核心实现:完整代码与关键细节
4.1 最终版代码(含限流+超时+错误隔离)
import asyncio
import aiohttp
import time
import logging
from flask import Flask, request, jsonify
from asgiref.sync import async_to_sync
logger = logging.getLogger(__name__)
# ------ 配置区 ------
DOWNSTREAM_TIMEOUT = 2.0 # 单次下游请求超时(秒)
MAX_CONCURRENCY = 50 # 全局并发信号量
CONNECTION_LIMIT = 100 # 连接池上限
# ------ 全局资源 ------
_semaphore = asyncio.Semaphore(MAX_CONCURRENCY)
_session = aiohttp.ClientSession(
connector=aiohttp.TCPConnector(
limit=CONNECTION_LIMIT,
ttl_dns_cache=300, # DNS缓存5分钟,减少DNS查询
force_close=False, # 保持keep-alive
),
timeout=aiohttp.ClientTimeout(total=DOWNSTREAM_TIMEOUT)
)
async def fetch_with_semaphore(session, url):
"""带信号量限流的GET请求"""
async with _semaphore:
try:
async with session.get(url) as resp:
if resp.status == 200:
return await resp.json()
else:
logger.warning(f"下游返回非200: {url} -> {resp.status}")
return None
except asyncio.TimeoutError:
logger.error(f"下游超时: {url}")
return None
except Exception as e:
logger.error(f"下游请求异常: {url} -> {str(e)}")
return None
async def aggregate_async(user_id: str):
"""并发聚合三个下游,返回dict或None"""
urls = [
f"http://profile-svc/user/{user_id}",
f"http://order-svc/stats/{user_id}",
f"http://risk-svc/tags/{user_id}",
]
# 每个任务独立超时控制,避免一个慢任务拖死全部
tasks = [asyncio.wait_for(fetch_with_semaphore(_session, url), timeout=DOWNSTREAM_TIMEOUT) for url in urls]
results = await asyncio.gather(*tasks, return_exceptions=True)
# 过滤掉异常和None,但保留部分结果
valid_results = [r for r in results if isinstance(r, dict)]
if not valid_results:
raise RuntimeError("所有下游服务均不可用")
return {
"profile": results[0] if isinstance(results[0], dict) else {"error": "profile_unavailable"},
"orders": results[1] if isinstance(results[1], dict) else {"error": "orders_unavailable"},
"risk": results[2] if isinstance(results[2], dict) else {"error": "risk_unavailable"},
"partial": len(valid_results) < 3
}
# 使用async_to_sync桥接,保持Flask视图函数同步风格
@app.route('/api/v1/report/aggregate')
def aggregate_report():
user_id = request.args.get('user_id')
start = time.perf_counter()
try:
# 每次请求创建新的事件循环,避免跨请求状态污染
result = asyncio.run(aggregate_async(user_id))
elapsed = time.perf_counter() - start
logger.info(f"聚合接口耗时: {elapsed*1000:.1f}ms, user_id={user_id}")
return jsonify(result)
except Exception as e:
logger.error(f"聚合失败: {str(e)}")
return jsonify({"error": "internal_error"}), 500
4.2 几个必须注意的细节
asyncio.run()每次调用创建新事件循环,不能在Flask应用级复用(会导致循环关闭后session失效)aiohttp.ClientSession是线程安全的,但不能跨事件循环使用——所以我们放在模块级,每次asyncio.run内部共用asyncio.gather(..., return_exceptions=True)配合isinstance()判断,比裸try/except更优雅地处理部分失败wait_for的超时要小于ClientTimeout,否则ClientTimeout会先触发,导致连接池状态异常
踩坑记录:三个让我debug到凌晨的坑
坑1:asyncio.run() 导致连接池失效
现象:第一个版本我天真地在模块级写了asyncio.get_event_loop(),然后在Flask请求里loop.run_until_complete()。结果跑到第100个请求时,所有下游请求全部超时。
原因:Flask的development server是多线程的,每个线程共享同一个event loop,但aiohttp的session绑定的是创建它时的loop。多线程并发调用同一个loop的run_until_complete()会导致致命错误。
解决:改用asyncio.run(),每个请求独立loop。代价是每次请求有约5ms的loop创建开销,但换来的是线程安全。如果追求极致性能,可以用asgiref.sync.async_to_sync + 持久loop,但复杂度高很多,不建议。
坑2:ClientTimeout 全局超时 vs 单请求超时
现象:设置ClientTimeout(total=2)后,单个下游请求超过2秒会被取消,但gather里的其他任务还在跑,导致整个函数等待最慢的那个任务。
解决:必须用wait_for包裹每个独立任务,这样即便某个任务超时被取消,gather也会立即返回其他已完成的结果。
坑3:DNS解析导致首次请求延迟
现象:压测时发现前50个请求P95很高,之后才降下来。
原因:aiohttp的DNS解析是异步的,但每次解析需要几十毫秒。如果不缓存,每个连接都要重新解析。
解决:TCPConnector(ttl_dns_cache=300),缓存DNS结果5分钟。另外,在Docker环境中,建议用host.docker.internal或固定IP,避免DNS解析成为瓶颈。
效果数据:对比压测结果
压测工具:locust,200并发,持续5分钟,测试环境与生产一致(2核4GB容器)。
| 指标 | 重构前(同步requests) | 重构后(asyncio+aiohttp) | 提升 |
|---|---|---|---|
| P50 延迟 | 2.1s | 420ms | 5.0x |
| P95 延迟 | 2.3s | 480ms | 4.8x |
| P99 延迟 | 3.1s | 650ms | 4.8x |
| 吞吐量(req/s) | 86 | 405 | 4.7x |
| 错误率 | 2.3% | 0.1% | -95% |
关键变化:
- 三个下游的串行调用(800+900+600=2300ms)变成并发调用(max(800,900,600)=900ms),理论上限是900ms,实测480ms说明连接池复用带来了额外收益
- 错误率从2.3%降到0.1%,因为超时控制更精准,不再出现整链超时
- 吞吐量提升4.7倍,CPU使用率从85%降到40%(等待IO的时间被释放出来处理更多请求)
总结:异步改造的适用边界与后续建议
这次重构的核心收益:把IO等待时间从2.3秒压缩到480毫秒,本质上是用并发换延迟。对于IO密集型的下游调用场景,这是性价比最高的优化方式——不需要改架构,不需要换框架,只改一个视图函数。
适用边界:
- 如果下游服务本身响应极快(<50ms),异步收益有限,不值得引入复杂度
- 如果你的代码有大量CPU密集计算,异步没用,应该考虑多进程或C扩展
- 如果下游服务不支持高并发(比如单线程的旧系统),信号量限流必须加上
后续优化方向:
1. 将asyncio.run()替换为持久化的async_to_sync + 统一事件循环,可再省5ms/请求
2. 接入opentelemetry做全链路追踪,观察每个下游的耗时分布
3. 对下游服务做熔断(基于错误率),而不是只靠超时控制
这次重构总共花了2天半时间,其中1天在查踩坑问题。如果早点明确事件循环的生命周期管理,半天就能搞定。希望这篇博客能帮你省下那1天。有问题欢迎评论区交流,特别是关于asyncio.run()和多线程的坑,我可以展开聊更多。