一、问题背景:一个慢接口拖垮了整个服务
上个月接手一个内部工单系统,用户反馈“导出报表经常超时”。查看监控发现/api/v1/report/export接口在高峰时段P99延迟高达1.2s,而该接口内部逻辑很简单:查数据库获取1000条工单ID,然后循环调用外部CRM系统的/batch/status接口查询状态,最后聚合返回。
问题出在循环调用上——用了requests.post同步请求,每次等待200-500ms,1000次串行就是200-500秒。后来加了线程池(ThreadPoolExecutor)勉强能跑,但线程切换开销大,且连接复用差。压测数据触目惊心:50并发下QPS仅120,线程数飙到300+,CPU闲置但IO等待严重。
二、环境与版本:别用最新,用最稳
- Python 3.10.12(
pyenv管理,3.11+的asyncio有性能提升但团队统一3.10) - aiohttp 3.9.1(
pip install aiohttp==3.9.1,3.10版本有连接池bug) - uvloop 0.17.0(可选,但强烈推荐,事件循环性能提升约30%)
- 压测工具:
wrk4.2.0(--threads=4 --connections=50 --duration=30s)
关键配置:aiohttp的TCPConnector需要显式设置limit=100(默认100,但容易耗尽),并启用keepalive_timeout=30。数据库连接用asyncpg 0.27.0(比aiomysql更稳定)。
三、方案设计:从同步阻塞到异步非阻塞
核心思路:将requests.post替换为aiohttp.ClientSession.post,配合asyncio.gather并发发起请求。但直接无限制并发会打死CRM系统,所以用asyncio.Semaphore(20)限制最大并发数(根据CRM压测,20并发是安全阈值)。
架构对比:
Before: Flask(线程池) → requests.post(同步阻塞) → CRM
After: Flask(单线程) → asyncio.run_coroutine_threadsafe → aiohttp(异步IO) → CRM
注意:Flask是WSGI同步框架,不能直接跑async def视图函数。我采用asyncio.run_coroutine_threadsafe将协程提交到独立的事件循环线程中运行,主线程保持响应。
四、核心实现:before/after代码对比
Before(同步阻塞版)
import requests
from concurrent.futures import ThreadPoolExecutor
from flask import Flask, jsonify
app = Flask(__name__)
executor = ThreadPoolExecutor(max_workers=20)
@app.route('/api/v1/report/export', methods=['GET'])
def export_report():
# 模拟查询数据库获取工单ID
ticket_ids = list(range(1, 1001))
def fetch_status(tid):
resp = requests.post(
'https://crm.internal/api/batch/status',
json={"ticket_id": tid},
timeout=2
)
return resp.json()['status']
results = list(executor.map(fetch_status, ticket_ids))
return jsonify({"count": len(results), "statuses": results})
After(异步非阻塞版)
import asyncio
import aiohttp
from flask import Flask, jsonify
import uvloop
app = Flask(__name__)
# 全局事件循环(在独立线程中运行)
loop = uvloop.new_event_loop()
asyncio.set_event_loop(loop)
# 全局aiohttp会话,复用连接池
connector = aiohttp.TCPConnector(limit=100, keepalive_timeout=30)
session = aiohttp.ClientSession(connector=connector, timeout=aiohttp.ClientTimeout(total=3))
async def fetch_status(sem, tid):
async with sem:
async with session.post(
'https://crm.internal/api/batch/status',
json={"ticket_id": tid}
) as resp:
data = await resp.json()
return data['status']
async def run_async(ticket_ids):
sem = asyncio.Semaphore(20) # 限制并发20
tasks = [asyncio.create_task(fetch_status(sem, tid)) for tid in ticket_ids]
return await asyncio.gather(*tasks)
@app.route('/api/v1/report/export', methods=['GET'])
def export_report():
ticket_ids = list(range(1, 1001))
# 通过run_coroutine_threadsafe将协程提交到事件循环线程
future = asyncio.run_coroutine_threadsafe(run_async(ticket_ids), loop)
results = future.result(timeout=10) # 主线程阻塞等待(但事件循环不阻塞)
return jsonify({"count": len(results), "statuses": results})
# 应用退出时关闭资源
@app.teardown_appcontext
def close_session(exception=None):
if hasattr(session, 'close'):
asyncio.run_coroutine_threadsafe(session.close(), loop).result()
关键点:
1. asyncio.Semaphore(20)控制并发,防止打爆CRM
2. aiohttp.ClientSession全局复用,避免每次请求创建新连接(这是性能提升的核心)
3. uvloop替换默认事件循环,减少系统调用开销
五、踩坑与优化:两个血泪教训
坑1:事件循环策略导致RuntimeError
第一次直接在视图函数里写asyncio.run(),结果报错“This event loop is already running”。因为Flask的请求处理线程和事件循环线程不是同一个,必须用run_coroutine_threadsafe。后来在应用启动时创建全局loop,并设置为uvloop。
坑2:连接池泄漏导致内存暴涨
最初没有复用ClientSession,每个请求async with aiohttp.ClientSession() as session,压测30分钟后内存从200MB涨到2GB。原因是ClientSession内部维护连接池,创建/销毁频繁导致文件描述符泄漏。改为全局单例后,内存稳定在350MB。
性能优化:
- 将json={"ticket_id": tid}改为json={"ticket_ids": [tid]}批量请求,但CRM接口不支持,只能作罢
- 使用asyncio.gather的return_exceptions=True,避免单个请求失败导致全部取消
- 设置aiohttp.ClientTimeout(total=3),防止慢请求拖垮整体
六、效果数据:QPS提升15倍
使用wrk -t4 -c50 -d30s http://localhost:5000/api/v1/report/export压测结果:
| 指标 | Before(同步) | After(异步) | 提升 |
|---|---|---|---|
| QPS | 120 | 1850 | 15.4x |
| 平均延迟 | 420ms | 54ms | 7.8x |
| P99延迟 | 800ms | 210ms | 3.8x |
| 线程数 | 300+ | 6(主线程+事件循环) | - |
| CPU占用 | 15% | 45%(IO等待变计算) | 利用更充分 |
压测环境:4核8G虚拟机,CRM服务部署在另一台机器(千兆内网)。延迟主要来自CRM端点(单次请求约50ms),异步并发20个请求时,总耗时约1000/20 * 50ms = 2.5s,但QPS提升是因为Flask主线程不再阻塞,可同时处理多个请求。
注意:QPS提升15倍的前提是CRM服务不成为瓶颈。实际生产中将Semaphore调至30后,CRM响应时间从50ms涨到120ms,说明已到极限。
七、总结:异步不是银弹,但IO密集场景是神器
这次重构让我深刻理解:同步代码适合CPU密集,异步适合IO密集。本场景中80%时间在等待网络IO,异步非阻塞将等待时间让渡给其他请求,资源利用率大幅提升。
后续建议:
1. 如果接口需要数据库操作,配合asyncpg + SQLAlchemy 1.4+的异步扩展
2. 监控事件循环的loop.slow_callback_duration,设置0.1s警告
3. 生产环境部署时,用gunicorn + uvicorn workers(但Flask需迁移为Quart或FastAPI才能完全异步化,当前方案是妥协)
最后想说:asyncio的坑不少,但一旦跑通,性能回报非常可观。如果你们也有类似同步循环调外部API的场景,值得一试。有问题欢迎评论区交流。