1. 问题背景:线程池被打满,CPU却在看戏
上个月接手一个订单服务,核心接口/api/v1/orders需要聚合用户信息、商品快照、物流状态三个下游服务。最初实现是同步Flask视图:
@app.route('/api/v1/orders')
def get_orders():
user = requests.get(f"http://user-svc/users/{uid}", timeout=2).json()
products = requests.post("http://product-svc/batch", json=pid_list, timeout=2).json()
logistics = requests.get(f"http://logistics-svc/tracks/{oid}", timeout=2).json()
return jsonify(assemble(user, products, logistics))
压测结果惨不忍睹:单实例QPS 120,P95延迟380ms,最要命的是ThreadPoolExecutor的20个线程全部卡在requests.get的socket等待上,CPU占用不到30%。这显然是IO密集型瓶颈,同步IO把线程全占死了,机器资源完全浪费。
2. 环境与版本:Python 3.10 + aiohttp
- Python 3.10.8(原生asyncio,没有额外装uvloop)
- aiohttp 3.8.4(替代requests做HTTP调用)
- Flask 2.2.3(保留框架,但视图改为异步)
- gunicorn 20.1.0 + uvicorn workers(注意后面有坑)
- wrk 4.2.0压测工具
为什么不用FastAPI?因为历史包袱,路由和中间件都是Flask的,迁移成本太高。好在Flask 2.2支持异步视图,只要装了asgiref兼容层就能跑。
3. 方案设计:全链路协程化 + 信号量限流
目标:把同步IO全部替换为异步IO,让单进程能同时处理上千个并发请求。
设计要点:
1. 用aiohttp.ClientSession代替requests,复用TCP连接
2. 用asyncio.gather并行发起三个下游调用
3. 用asyncio.Semaphore(100)控制并发量,防止打爆下游
4. 每个下游调用单独设置超时,防止整体超时
核心数据结构是asyncio.Queue和Semaphore,配合asyncio.timeout(Python 3.11之前用asyncio.wait_for)。
4. 核心实现:改造后的异步视图
import asyncio
import aiohttp
from flask import Flask, jsonify, request
app = Flask(__name__)
# 全局session,复用连接池
session = None
# 信号量限制最大并发100个下游请求
semaphore = asyncio.Semaphore(100)
TIMEOUT = aiohttp.ClientTimeout(total=2.0)
async def fetch_json(client, url, payload=None):
async with semaphore:
try:
if payload:
async with client.post(url, json=payload) as resp:
return await resp.json()
else:
async with client.get(url) as resp:
return await resp.json()
except asyncio.TimeoutError:
return {"error": "timeout"}
except aiohttp.ClientError as e:
return {"error": str(e)}
async def get_orders_async(uid, pid_list, oid):
async with aiohttp.ClientSession(timeout=TIMEOUT) as client:
user_task = fetch_json(client, f"http://user-svc/users/{uid}")
product_task = fetch_json(client, "http://product-svc/batch", payload={"ids": pid_list})
logistics_task = fetch_json(client, f"http://logistics-svc/tracks/{oid}")
user, products, logistics = await asyncio.gather(user_task, product_task, logistics_task)
return assemble(user, products, logistics)
@app.route('/api/v1/orders', methods=['GET'])
async def get_orders():
uid = request.args.get('uid')
pid_list = request.args.getlist('pid')
oid = request.args.get('oid')
result = await get_orders_async(uid, pid_list, oid)
return jsonify(result)
# 应用启动时创建全局session
@app.before_first_request
def init_session():
global session
session = aiohttp.ClientSession(timeout=TIMEOUT)
注意:before_first_request在Flask 2.3已被移除,生产环境我用的是app.before_request加标志位。另外ClientSession必须复用,每次新建会浪费大量TCP握手时间。
5. 踩坑与优化:两个隐秘的性能杀手
坑1:Flask同步视图下的asyncio会失效
最初我直接在同步视图里写asyncio.run(get_orders_async(...)),结果QPS反而降到80。因为asyncio.run每次都会创建新的事件循环,线程还是被阻塞在事件循环的run_until_complete上。必须用async def定义视图,让gunicorn的worker直接驱动协程。
坑2:DNS解析阻塞EventLoop
压测时发现延迟抖动剧烈,用py-spy dump抓到栈,发现大量协程卡在getaddrinfo上。aiohttp默认使用线程池解析DNS,但线程池只有5个,一旦打满就阻塞整个loop。解决:
import socket
# 强制使用异步DNS解析
resolver = aiohttp.AsyncResolver(nameservers=["8.8.8.8", "1.1.1.1"])
connector = aiohttp.TCPConnector(resolver=resolver, use_dns_cache=True, ttl_dns_cache=300)
session = aiohttp.ClientSession(connector=connector, timeout=TIMEOUT)
坑3:gunicorn的worker类型
必须用uvicorn.workers.UvicornWorker,否则gunicorn的同步worker会把asyncio视图当普通函数跑,完全不起作用。启动命令:
gunicorn -w 4 -k uvicorn.workers.UvicornWorker -b 0.0.0.0:8000 app:app
6. 效果数据:QPS翻8倍,延迟降75%
压测环境:4核8G云主机,wrk -t8 -c500 -d30s。
| 指标 | 同步版本 | 异步版本 |
|---|---|---|
| QPS | 120 | 980 |
| P50延迟 | 250ms | 42ms |
| P95延迟 | 380ms | 89ms |
| P99延迟 | 610ms | 210ms |
| CPU占用 | 28% | 75% |
| 异常率 | 0.5% | 0.1% |
异步版本把并发能力从20个线程扩展到了上千协程,CPU终于跑满了。下游服务压力测试显示,商品服务的QPS从300涨到1200,说明之前是线程池瓶颈,不是下游瓶颈。
7. 总结与建议
- 不要用
requests做并发IO,线程池是有限资源,协程是廉价资源 asyncio.gather要配合return_exceptions=True使用,否则单个任务异常会取消所有任务- 生产环境必须监控EventLoop延迟,用
loop.time()打点,超过100ms就要查慢任务 - 如果下游API不支持批量接口,考虑用
asyncio.Queue做流量整形,避免突发请求打崩下游
这次重构收益远大于成本,代码只改了200行,但性能提升了8倍。如果你的服务也是IO密集且CPU闲置,建议立刻尝试asyncio迁移。记住:线程是给CPU密集型任务准备的,IO密集型请交给协程。