1. 问题背景:一个串行API的"肉眼可见"卡顿
事情要从去年秋天的一次线上压测说起。我们维护的舆情监测系统,核心接口是/analyze——用户提交一个关键词,后端需要同时去三个外部服务查询:情感分析(耗时约800ms)、实体提取(约1.2s)、热点指数(约600ms)。最初的实现是同步的Flask视图,用requests库依次请求:
# before.py - 同步版本
import requests
import time
def analyze_sync(keyword: str) -> dict:
start = time.perf_counter()
# 三个请求串行执行
sentiment = requests.post("http://sentiment-api:8000/analyze",
json={"text": keyword}, timeout=3).json()
entities = requests.post("http://entity-api:8000/extract",
json={"text": keyword}, timeout=3).json()
hotness = requests.get("http://hotness-api:8000/index",
params={"keyword": keyword}, timeout=3).json()
elapsed = time.perf_counter() - start
print(f"同步耗时: {elapsed:.2f}s")
return {"sentiment": sentiment, "entities": entities, "hotness": hotness}
压测数据触目惊心:并发30个请求时,平均响应时间2.8秒,P95达到4.1秒,大量请求因为积压而超时(前端设置3秒超时)。最关键的,这三个外部API之间完全独立,没有数据依赖——显然,串行是在浪费生命。
2. 环境与版本:一个干净的异步栈
改造前先确认环境。我们用的Python版本是3.10.12(3.7+支持asyncio,但3.10的TaskGroup和asyncio.run()异常处理更友好)。关键依赖:
- aiohttp==3.8.6(异步HTTP客户端,比
httpx在连接池管理上更成熟) - uvloop==0.17.0(事件循环加速器,Linux上可提升15-20%吞吐)
- Flask==2.3.3(同步框架下如何集成异步?我们下文会处理)
服务器是 4vCPU 8GB 的阿里云ECS,操作系统 Ubuntu 22.04。压测工具用wrk,参数:wrk -t4 -c30 -d30s http://api:5000/analyze?keyword=test。
3. 方案设计:三根协程同时下水
核心思路很直白:将三个独立的HTTP请求从串行改为并发。asyncio提供了三种并发模式:
asyncio.gather():最经典,但有一个"坑"——如果某个协程抛出异常,其他协程会被取消(除非设置return_exceptions=True)。asyncio.TaskGroup(Python 3.11+):结构化并发,自动处理异常传播,但我们用3.10,放弃。asyncio.create_task()+asyncio.wait():手动管理,适合需要精细控制超时和取消的场景。
我们选择了方案1的变体——asyncio.gather(return_exceptions=True),原因有二:一是代码最简洁,二是不同外部API的失败不应互相影响(比如情感分析超时,热点指数应该正常返回)。
整体架构改动很小:把同步视图函数改为异步,用aiohttp替代requests,并在应用启动时创建全局的aiohttp.ClientSession(避免每个请求都创建session,那会浪费TCP连接)。
4. 核心实现:用asyncio写出"三管齐下"的异步API
4.1 异步化改造的核心代码
# after.py - 异步版本
import asyncio
import aiohttp
from aiohttp import ClientTimeout, TCPConnector
from flask import Flask, request, jsonify
import time
app = Flask(__name__)
# 全局session,复用连接池
session = None
async def fetch_sentiment(session: aiohttp.ClientSession, keyword: str) -> dict:
async with session.post(
"http://sentiment-api:8000/analyze",
json={"text": keyword},
timeout=ClientTimeout(total=2.5) # 每个请求独立超时
) as resp:
return await resp.json()
async def fetch_entities(session: aiohttp.ClientSession, keyword: str) -> dict:
async with session.post(
"http://entity-api:8000/extract",
json={"text": keyword},
timeout=ClientTimeout(total=3.0)
) as resp:
return await resp.json()
async def fetch_hotness(session: aiohttp.ClientSession, keyword: str) -> dict:
async with session.get(
"http://hotness-api:8000/index",
params={"keyword": keyword},
timeout=ClientTimeout(total=2.0)
) as resp:
return await resp.json()
async def analyze_async(keyword: str) -> dict:
# 并发发起三个请求
results = await asyncio.gather(
fetch_sentiment(session, keyword),
fetch_entities(session, keyword),
fetch_hotness(session, keyword),
return_exceptions=True # 个别失败不中断整体
)
# 处理可能的异常
sentiment = results[0] if not isinstance(results[0], Exception) else {}
entities = results[1] if not isinstance(results[1], Exception) else {}
hotness = results[2] if not isinstance(results[2], Exception) else {}
return {"sentiment": sentiment, "entities": entities, "hotness": hotness}
@app.route('/analyze', methods=['GET'])
def analyze():
keyword = request.args.get('keyword', '')
if not keyword:
return jsonify({"error": "keyword required"}), 400
# Flask同步视图内运行异步代码
result = asyncio.run(analyze_async(keyword))
return jsonify(result)
if __name__ == '__main__':
# 创建带连接池的session
connector = TCPConnector(
limit=100, # 总连接数上限
limit_per_host=30, # 每个主机的连接数
ttl_dns_cache=300, # DNS缓存300秒
enable_cleanup_closed=True
)
session = aiohttp.ClientSession(connector=connector)
app.run(host='0.0.0.0', port=5000, threaded=True)
4.2 关键设计决策
- 连接池参数:
limit=100意味着最多同时建立100个TCP连接。我们压测并发30个请求,每个请求3个外部API,理论峰值需90个连接,100刚好。 - 独立超时:每个外部API的超时不同(情感分析2.5s,实体提取3s,热点指数2s),避免一个慢服务拖死其他。
- 异常处理:
return_exceptions=True是关键——如果实体提取超时,情感分析和热点指数仍能正常返回,不会因为一个服务故障而导致整个请求失败。
5. 踩坑与优化:从580 QPS到620 QPS的两次修复
第一次压测就遇到了一个诡异问题:运行20秒后,请求延迟突然从0.8秒飙升到3秒以上,查看系统监控发现TCP连接数持续增长,接近系统上限(65535)。这就是著名的协程泄漏——aiohttp.ClientSession默认不自动关闭底层连接,当请求频繁时,连接被占用后未释放回池子。
修复1:显式控制连接生命周期
connector = TCPConnector(
limit=100,
limit_per_host=30,
enable_cleanup_closed=True, # 关键:关闭后清理
force_close=True, # 请求结束后强制关闭底层连接
)
但force_close=True会导致频繁的TCP四次挥手,增加开销。更好的做法是保持长连接,配合keepalive_timeout:
connector = TCPConnector(
limit=100,
limit_per_host=30,
keepalive_timeout=30, # 空闲连接30秒后关闭
)
修复2:SSL握手瓶颈
第二个问题出现在情感分析API(https协议)。压测时CPU的sys态占用很高,分析发现是SSL握手时消耗了大量系统调用。解决方案:启用SSL会话复用。
import ssl
ssl_context = ssl.create_default_context()
ssl_context.set_ciphers('DEFAULT:@SECLEVEL=1') # 降低安全等级以提升速度(内部API)
connector = TCPConnector(ssl=ssl_context)
对于内部API,这个优化是安全的。如果是对公网API,建议保持默认安全配置。
修复3:Flask的同步GIL问题
注意到没?我们的视图函数是同步的(def analyze()),里面调用asyncio.run()。这其实引入了问题:Flask的线程池中,每个线程都在运行自己的事件循环,多个事件循环之间互相抢GIL。更好的方法是用Quart(异步Flask)或Sanic,但为了最小改动,我们在视图里加了asyncio.set_event_loop_policy(uvloop.EventLoopPolicy()),让事件循环走uvloop(C扩展实现),减少GIL争用。
最终版本:
import uvloop
asyncio.set_event_loop_policy(uvloop.EventLoopPolicy())
6. 效果数据:从2.8秒到0.8秒,P95下降80%
最终在4vCPU机器上的压测结果(wrk -t4 -c30 -d60s):
| 指标 | 同步版本 | 异步版本(基础) | 异步版本(优化后) |
|---|---|---|---|
| 平均延迟 | 2.81s | 0.92s | 0.78s |
| P95延迟 | 4.12s | 1.05s | 0.83s |
| P99延迟 | 5.67s | 1.43s | 1.12s |
| QPS | 120 | 480 | 620 |
| 错误率 | 3.2% | 0.1% | 0.05% |
性能提升的核心原因:
1. I/O并发:三个串行请求(2.8s)变为并行(最慢的实体提取1.2s),理论极限延迟从sum变成max。
2. 连接复用:aiohttp的连接池避免了每个请求都建立TCP连接(三次握手约40ms)。
3. uvloop加速:事件循环的C扩展实现,减少了Python层面的调度开销。
7. 总结
这次改造给我的教训是:异步不是银弹,但面对大量独立I/O请求时,它是性价比最高的优化手段。 整个改动只涉及一个文件、约50行代码变更,却带来了5倍以上的吞吐提升。
如果你也遇到类似的"串行调用多个外部API"的场景,我建议按这个顺序排查:
1. 确认是否有数据依赖——没有就上asyncio.gather
2. 设置合理的连接池参数(不要默认)
3. 永远处理异常,防止一个协程的失败拖垮全部
4. 压测时监控TCP连接数和文件句柄,防止泄漏
最后留个思考:如果三个外部API中有两个是CPU密集型计算(比如本地模型推理),异步就没用了——这时候该考虑线程池或进程池。asyncio只擅长等待,不擅长计算。