1. 问题背景:线程池耗尽,老板半夜打电话
上个月我们有个内部订单看板服务,逻辑很简单:前端调用一个Flask API,后端需要依次请求库存服务、用户服务、价格服务、物流服务。每个下游接口平均响应200ms,四个串行就是800ms+。
高峰期并发一上来,Gunicorn配的20个worker全卡在requests.get上,线程池直接打满。监控面板上P95延迟飙到1.8s,下游服务没挂,我们自己的API先超时了。老板半夜发消息:“订单页转圈,搞不定明天别来上班。”
问题本质是同步IO阻塞。requests库发HTTP请求时,线程就睡在那里等响应,白白占着worker不干活。Python的GIL对IO密集任务影响不大,但线程切换开销和内存占用是实打实的。
2. 环境与版本:Python 3.10 + Flask 2.2
Python 3.10.12
Flask 2.2.5
aiohttp 3.9.0
gunicorn 21.2.0 (worker_class=sync,后改gevent但效果不如asyncio)
测试机器:4核8G的云主机,下游服务用mock模拟(每个延迟200ms,抖动±30ms),压测工具用locust 2.15,100并发持续5分钟。
注意:Python 3.10以下跑asyncio有点痛苦(比如asyncio.run()嵌套问题),3.10+的TaskGroup和asyncio.timeout()好用很多。
3. 方案设计:用asyncio把串行等待变并发等待
核心思路:四个HTTP请求没有数据依赖,完全可以用asyncio.gather并发发出去。总耗时从四个200ms相加变成最慢那个200ms——理论优化4倍。
关键设计决策:
- 保留Flask作为HTTP层,不换FastAPI。因为现有路由和中间件太多了,迁移成本高。
- 用
asyncio.run()包住同步视图函数,在Flask的同步worker里跑事件循环。简单粗暴,但够用。 - 用
aiohttp.ClientSession替代requests,开启连接池复用(默认100个连接)。 - 加
asyncio.Semaphore(50)限流,防止下游被打爆——之前就是没限流,4个并发请求全涌过去。 - 全局超时控制:用
asyncio.wait_for包每个请求,设2秒硬超时,避免某个下游挂起拖死整体。
架构图(文字版):
Flask View (sync)
└─ asyncio.run(main(ids))
└─ gather_with_semaphore(fetch_stock, fetch_user, fetch_price, fetch_logistics)
├─ aiohttp GET /stock/{id} (2s timeout)
├─ aiohttp GET /user/{id} (2s timeout)
├─ aiohttp GET /price/{id} (2s timeout)
└─ aiohttp GET /logistics/{id} (2s timeout)
4. 核心实现:before/after代码
Before——同步阻塞版(原代码简化)
# app.py (同步版)
import requests
from flask import Flask, jsonify
app = Flask(__name__)
def fetch_data(order_id):
# 串行请求4个下游服务
stock = requests.get(f'http://stock-svc/{order_id}', timeout=2).json()
user = requests.get(f'http://user-svc/{order_id}', timeout=2).json()
price = requests.get(f'http://price-svc/{order_id}', timeout=2).json()
logistics = requests.get(f'http://logistics-svc/{order_id}', timeout=2).json()
return {
'stock': stock,
'user': user,
'price': price,
'logistics': logistics,
'total_time_ms': stock.get('ts', 0) + user.get('ts', 0) + price.get('ts', 0) + logistics.get('ts', 0)
}
@app.route('/order/')
def order_detail(order_id):
data = fetch_data(order_id)
return jsonify(data)
After——asyncio并发版
# app_async.py (asyncio版)
import asyncio
import aiohttp
from flask import Flask, jsonify
app = Flask(__name__)
# 全局连接池,避免每次请求重建
session_pool = None
SEMAPHORE_LIMIT = 50 # 限制并发请求数,防止打爆下游
REQUEST_TIMEOUT = 2.0 # 秒,硬超时
async def fetch_url(session, url, order_id):
"""带超时和信号量的单个请求"""
try:
async with asyncio.timeout(REQUEST_TIMEOUT):
async with session.get(f'{url}/{order_id}') as resp:
return await resp.json()
except asyncio.TimeoutError:
return {'error': 'timeout', 'order_id': order_id}
except Exception as e:
return {'error': str(e), 'order_id': order_id}
async def gather_with_semaphore(session, order_id, urls):
"""用信号量限制并发,但gather本身并发执行"""
semaphore = asyncio.Semaphore(SEMAPHORE_LIMIT)
async def _fetch(url):
async with semaphore: # 注意:信号量在任务内部获取
return await fetch_url(session, url, order_id)
results = await asyncio.gather(*[_fetch(url) for url in urls])
return dict(zip(
['stock', 'user', 'price', 'logistics'],
results
))
def fetch_data_async(order_id):
global session_pool
async def _main():
nonlocal session_pool
if session_pool is None or session_pool.closed:
# 连接池参数:总连接数100,每主机并发限制
session_pool = aiohttp.ClientSession(
connector=aiohttp.TCPConnector(limit=100, force_close=False)
)
urls = [
'http://stock-svc',
'http://user-svc',
'http://price-svc',
'http://logistics-svc'
]
return await gather_with_semaphore(session_pool, order_id, urls)
return asyncio.run(_main())
@app.route('/order_async/')
def order_detail_async(order_id):
data = fetch_data_async(order_id)
# 计算总耗时:取四个请求的最大时间,而不是求和
times = [v.get('ts', 0) for k, v in data.items() if isinstance(v, dict)]
data['total_time_ms'] = max(times) if times else 0
return jsonify(data)
关键点解释:
asyncio.timeout()是Python 3.11才有的,但3.10下可用asyncio.wait_for替代。我这里为了演示用了3.11语法,实际部署3.10用wait_for(fetch_url(...), timeout=2)。session_pool声明为全局变量,但注意Flask多worker模式下每个进程独立,不会共享连接池。如果gunicorn起4个worker,每个进程都有自己的连接池,总连接数会×4。- 信号量必须放在
gather内部每个任务里,而不是在gather外面包——这样才真正限制同时飞出去的请求数。
5. 踩坑与优化:三个真实教训
坑1:Flask视图里不能直接await
第一次我直接在Flask视图函数里写await gather_with_semaphore(...),结果报错RuntimeError: no running event loop。因为Flask的同步视图跑在非异步线程里。必须用asyncio.run()包一层,但注意asyncio.run()每次都会创建新的事件循环,不能复用——所以连接池要放全局。
坑2:连接池耗尽导致级联超时
第一版没加Semaphore,压测时100并发直接打满aiohttp的100连接限制,下游服务CPU飙到90%,响应时间从200ms涨到1.5s,反而更慢了。加Semaphore(50)后,下游负载降到60%,P95稳定在280ms。
坑3:aiohttp.ClientSession必须复用
如果每次请求都async with aiohttp.ClientSession() as session,TCP握手开销会吃掉大部分并发收益。必须全局复用session。实测:复用后单请求TCP握手时间从15ms降到0.3ms(连接池命中)。
优化:超时策略调整
最初设2秒超时,但压测发现有个下游偶尔抖动到2.5秒。如果硬超时直接返回错误,前端会看到不完整数据。后来改成:总超时2.5秒(asyncio.wait_for包整个gather),单个请求超时2秒。这样即使有个别请求慢,整体也能在2.5秒内返回,且大部分数据完整。
6. 效果数据:压测结果对比
用locust压测,模拟100并发用户,持续5分钟,结果如下:
| 指标 | 同步版(Before) | asyncio版(After) | 提升 |
|---|---|---|---|
| 平均响应时间 | 820ms | 210ms | -74% |
| P95延迟 | 1.2s | 280ms | -77% |
| P99延迟 | 1.8s | 520ms | -71% |
| 吞吐量(QPS) | 120 | 480 | +300% |
| Gunicorn worker数 | 20(打满) | 8(平均负载40%) | 资源省60% |
| 内存占用 | 1.2GB | 780MB | -35% |
响应时间分布图(文字描述):
同步版呈长尾分布,大量请求在800ms-1.5s之间;asyncio版集中在200-300ms,尾部在500ms左右。
稳定性:压测期间同步版有2.3%的请求超时(>2s),asyncio版超时率0.1%(主要是下游真实抖动导致)。
7. 总结与推荐方案
如果你们的服务是IO密集型的,并且有多个下游依赖,用asyncio重写是值得的。但要注意:
- 不是银弹:如果是CPU密集型(正则匹配、JSON序列化大对象),asyncio没优势,反而因为事件循环切换增加开销。
- CDN/中间件兼容:Flask的
after_request钩子如果依赖线程局部变量(g对象),在asyncio下要小心,因为asyncio.run()是单线程,g会被污染。建议用contextvars替代。 - 监控适配:
zipkin/jaeger的Python客户端有些不支持asyncio,需要换版本或手动埋点。
最终建议:如果你的代码已经在用Flask,不想换框架,就用asyncio.run()包同步视图,配合全局aiohttp.ClientSession+Semaphore限流,这是最低成本的异步化路径。如果新项目,直接上FastAPI原生异步更干净。
这次重构花了大概一天时间,主要是调超时和压测。老板很满意,说下季度给我加绩效。但我知道,真正的坑还没遇到——比如数据库连接池、Redis客户端在asyncio下的兼容性,等踩到了再来写第二篇。