1. 问题背景:一个慢到被运维找上门的API
我们有个内部服务叫/api/v1/orders/batch,前端一次要拉取50个订单的详情。最初实现就是标准的Flask同步视图,内部串行调用下游订单服务的五个接口(详情、商品、物流、优惠、用户)。每次请求平均耗时1.8秒,压测时QPS卡在180左右。
运维同事丢来一张监控截图:CPU才用了30%,但线程数飙到400+,大量请求堆积在等待下游IO。这就是典型的IO密集型瓶颈——线程都在等网络响应,GIL锁导致CPU无法充分利用,而线程切换开销又大。
2. 环境与版本:先交代清楚实验条件
- Python 3.10.12(asyncio在3.10才支持
asyncio.timeout,3.8的旧项目建议用async_timeout包) - Flask 2.3.3(WSGI同步框架,需要配合
asgiref或直接换aiohttp) - aiohttp 3.9.1(异步HTTP客户端)
- gunicorn 21.2.0(生产部署,worker模式后面细说)
- 压测工具:wrk 4.2.0,参数
-t8 -c200 -d30s
测试环境:4核8G Docker容器,下游服务用fastapi模拟,每个接口延迟200ms。
3. 方案设计:为什么选asyncio而不是ThreadPoolExecutor
先看原始代码(简化版):
# before_sync.py
import requests
from flask import Flask, jsonify
app = Flask(__name__)
def fetch_order_detail(order_id):
# 模拟串行调用5个下游接口,每个200ms
result = {}
result['order'] = requests.get(f'http://svc/order/{order_id}').json()
result['product'] = requests.get(f'http://svc/product/{order_id}').json()
result['logistics'] = requests.get(f'http://svc/logistics/{order_id}').json()
result['promo'] = requests.get(f'http://svc/promo/{order_id}').json()
result['user'] = requests.get(f'http://svc/user/{order_id}').json()
return result
@app.route('/api/v1/orders/batch', methods=['POST'])
def batch():
order_ids = request.json['order_ids'] # 50个ID
result = [fetch_order_detail(oid) for oid in order_ids]
return jsonify(result)
这里有两个层次的串行:订单之间串行(50个订单循环),每个订单内部串行(5个下游接口)。
有人会问:用ThreadPoolExecutor加requests不行吗?我试过,能提到500QPS,但有两个问题:1)线程池大小难调,调大了内存暴涨(每个线程默认8MB栈);2)下游服务一旦变慢,线程池排队,整体延迟迅速恶化。
asyncio的优势在于单线程内协程切换,不需要操作系统线程上下文切换,一个进程能支撑数万连接。更重要的是,aiohttp原生复用连接池,避免每次请求都重新握手。
4. 核心实现:全异步化改造
改造分两步:第一步,把requests换成aiohttp,实现fetch_order_detail的异步版本;第二步,把Flask换成aiohttp.web(因为Flask不支持异步视图,虽然3.0支持了async def视图,但底层WSGI还是同步的)。
# after_async.py
import asyncio
import aiohttp
from aiohttp import web
# 全局连接池:限制最大连接数,避免打爆下游
connector = aiohttp.TCPConnector(limit=100, ttl_dns_cache=300)
session = aiohttp.ClientSession(connector=connector)
async def fetch_one(session, url, timeout=0.5):
"""单个下游请求,带超时控制"""
try:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=timeout)) as resp:
return await resp.json()
except asyncio.TimeoutError:
return {'error': 'timeout'}
async def fetch_order_detail(session, order_id):
# 并发请求5个下游,不相互等待
urls = {
'order': f'http://svc/order/{order_id}',
'product': f'http://svc/product/{order_id}',
'logistics': f'http://svc/logistics/{order_id}',
'promo': f'http://svc/promo/{order_id}',
'user': f'http://svc/user/{order_id}',
}
tasks = [fetch_one(session, url) for url in urls.values()]
results = await asyncio.gather(*tasks, return_exceptions=True)
return dict(zip(urls.keys(), results))
async def batch_handler(request):
data = await request.json()
order_ids = data['order_ids'] # 50个
# 控制并发度:同时最多处理20个订单,避免下游过载
sem = asyncio.Semaphore(20)
async def bounded_fetch(oid):
async with sem:
return await fetch_order_detail(session, oid)
# 全部订单并发执行
results = await asyncio.gather(*(bounded_fetch(oid) for oid in order_ids))
return web.json_response(results)
app = web.Application()
app.router.add_post('/api/v1/orders/batch', batch_handler)
if __name__ == '__main__':
web.run_app(app, host='0.0.0.0', port=8000, access_log=None)
关键改动点:
1. asyncio.gather并发所有订单:50个订单同时启动,每个订单内部5个请求也并发,理论最大并发=505=250个连接。
2. Semaphore(20)限流:防止一次性250个请求把下游打满,实测下游连接数超过100就会开始报错。
3. aiohttp.ClientTimeout(total=0.5)*:每个下游请求超时500ms,整体最坏情况约1秒(因为5个并发,取最慢的),比原来1.8秒好。
部署命令也变了,不能用flask run:
# 原先:gunicorn -w 4 before_sync:app
# 现在:aiohttp应用直接用python起,或配合uvloop
python after_async.py
# 生产建议:python -m aiohttp.web -H 0.0.0.0 -P 8000 after_async:app
5. 踩坑与优化:三个坑差点让我放弃
坑1:Event Loop被同步代码阻塞
最初我在fetch_order_detail里不小心调用了time.sleep(0.1)模拟延迟,结果整个事件循环卡住,QPS直接归零。排查半天发现是同事在中间件里用了同步的redis.get。记住:asyncio代码里绝不能有同步阻塞调用,一切IO都要走协程(await)或者丢给线程池(loop.run_in_executor)。
坑2:连接池耗尽导致连环超时
aiohttp.ClientSession默认连接池大小是100。当并发订单数从10调到50时,下游连接数瞬间到250,aiohttp内部排队等待可用连接,导致超时。解决方式就是上面代码里的connector参数,限制limit=100,同时用Semaphore控制应用层并发。
坑3:Python 3.10的asyncio.timeout vs 老项目
3.10之前用asyncio.wait_for,但它在超时后不会取消任务,导致任务继续运行浪费资源。3.11的asyncio.timeout是上下文管理器,推荐用这个。我因为生产环境是3.10.12,用了asyncio.timeout,效果不错。
6. 效果数据:对比压测结果
用wrk压测30秒,200并发,结果如下:
| 指标 | 同步版 (Flask + requests) | 异步版 (aiohttp) |
|---|---|---|
| QPS | 182 | 1547 |
| 平均延迟 | 1080ms | 640ms |
| P99延迟 | 2300ms | 280ms |
| 最大延迟 | 4100ms | 950ms |
| 内存占用 | 420MB (4 worker × 105MB) | 230MB (1进程) |
| CPU使用 | 30% (瓶颈在线程切换) | 85% (有效利用) |
QPS提升约8.5倍,P99从2.3秒降到280ms(8.2倍提升)。最意外的是内存从420MB降到230MB——因为去掉了4个gunicorn worker,单进程事件循环处理所有请求。
补充一个真实场景的观察:同步版在下游服务抖动时(延迟从200ms涨到800ms),QPS会跌到50以下;异步版由于超时控制(500ms),最坏情况P99仍在800ms内,不会无限堆积。
7. 总结与后续优化方向
这次重构让我确信:在Python里做高并发IO密集型服务,asyncio是首选方案。但要注意几点:
- 不是所有项目都适合改造,如果业务逻辑里同步库占大头(比如重度依赖requests、psycopg2同步驱动),改造成本高,收益有限。
- 异步化后要重点监控连接池状态和事件循环延迟(可以用loop.slow_callback_duration设置警告阈值)。
- 后续可以做的优化:用uvloop替代默认事件循环(再提升15%左右)、HTTP/2多路复用(减少连接数)、缓存下游响应(async_lru装饰器)。
最后提一句,如果团队刚接触asyncio,建议先从一个非核心API试点,踩过坑再推广。我这套代码已经在生产跑了两个月,稳定在1200-1500 QPS,运维终于不用半夜找我了。