1. 问题背景:一个拖垮数据库的同步接口
上个月我们有个订单导出接口频繁超时,监控显示单次请求平均耗时3.8秒,数据库连接池经常被打满。看代码是典型的同步嵌套——Flask视图里用requests调内部订单服务,拿到JSON后再用open()写临时文件。
# 优化前:同步阻塞版本
@app.route('/api/orders/export')
def export_orders():
order_data = []
for order_id in range(1000):
resp = requests.get(f'http://order-service/api/orders/{order_id}', timeout=2)
order_data.append(resp.json())
with open('/tmp/orders.json', 'w') as f:
json.dump(order_data, f)
return send_file('/tmp/orders.json')
问题很明显:requests.get()是IO阻塞型调用,每个请求平均等200ms,1000个订单就是200秒?不对,这里串行的话应该是1000 * 200ms = 200秒,但实际因为有分页批量查询所以是3.8秒,但并发一上来,GIL加上线程切换,CPU和内存都吃不消。
2. 环境与版本
- Python 3.10.8(asyncio在3.10成熟度最好,3.11的TaskGroup虽然好但团队没升)
- Flask 2.2.3(同步框架,但可以用
asgiref.sync.async_to_sync桥接协程) - httpx 0.23.3(支持异步HTTP客户端,比aiohttp更顺手)
- 部署:4核8G的ECS,Gunicorn 20.1.0,4 workers,每个worker跑一个事件循环
3. 方案设计:三层异步化
先把同步的requests换掉,然后是文件IO,最后是Flask视图函数的协程化。核心思路是:所有IO操作都要让出事件循环。
第一层:HTTP调用改用httpx.AsyncClient,用asyncio.gather并发请求。
第二层:文件写入用aiofiles.open异步写入。
第三层:视图函数本身用async_to_sync包装,让Flask能跑协程,但注意不能直接在协程里调Flask的request/session上下文,需要保持同步入口。
4. 核心实现:asyncio.Semaphore控制并发风暴
直接无脑asyncio.gather所有请求会打爆下游服务。我们加了个信号量限制同时进行的请求数,实测最优值是30。
# 优化后:异步非阻塞版本
import asyncio
import httpx
import aiofiles
from asgiref.sync import async_to_sync
# 信号量控制并发数,防止下游服务雪崩
semaphore = asyncio.Semaphore(30)
async def fetch_order(client, order_id):
async with semaphore:
resp = await client.get(f'http://order-service/api/orders/{order_id}')
return resp.json()
async def export_orders_async(order_ids):
async with httpx.AsyncClient(timeout=3.0, limits=httpx.Limits(max_connections=50)) as client:
tasks = [fetch_order(client, oid) for oid in order_ids]
results = await asyncio.gather(*tasks, return_exceptions=True)
# 过滤异常结果
return [r for r in results if isinstance(r, dict)]
@app.route('/api/orders/export')
def export_orders():
order_ids = get_order_ids_from_db()
# async_to_sync 桥接同步框架和协程世界
data = async_to_sync(export_orders_async)(order_ids)
# 异步写文件
async def write_file():
async with aiofiles.open('/tmp/orders.json', 'w') as f:
await f.write(json.dumps(data))
async_to_sync(write_file)()
return send_file('/tmp/orders.json')
关键点:semaphore定义在函数外面,如果定义在fetch_order内部,每次请求都会创建一个新的信号量,等于没限制。
5. 踩坑与优化:三个隐蔽的性能杀手
坑1:协程泄漏。第一次跑完发现内存涨了300MB不释放。排查半天是httpx.AsyncClient没关闭。用async with包装后解决。另外注意asyncio.gather的return_exceptions=True必须加,否则一个请求挂了,整个协程组都会取消。
坑2:事件循环被CPU密集操作卡死。我们的json.dumps处理1万条订单数据要耗时500ms,这期间事件循环完全卡住。解决办法是把这个操作丢到线程池里:
import asyncio
from concurrent.futures import ThreadPoolExecutor
executor = ThreadPoolExecutor(max_workers=2)
async def write_file(data):
loop = asyncio.get_running_loop()
# 将同步的json.dumps放到线程池执行,不阻塞事件循环
json_str = await loop.run_in_executor(executor, lambda: json.dumps(data))
async with aiofiles.open('/tmp/orders.json', 'w') as f:
await f.write(json_str)
坑3:数据库连接也是阻塞的。我们的get_order_ids_from_db()用的还是同步的psycopg2,这个函数在协程世界里是灾难。后来改成asyncpg,或者用aiosqlite(如果数据量小)。
6. 效果数据:压测结果对比
用wrk压测,场景是1000个订单ID,并发50连接,压测5分钟:
| 指标 | 优化前(同步) | 优化后(asyncio) |
|---|---|---|
| 平均响应时间 | 3.84s | 1.12s |
| QPS | 32 | 103 |
| CPU占用 | 85%(4核满载) | 46% |
| 内存峰值 | 1.2GB | 680MB |
| 数据库连接数 | 峰值40 | 峰值8 |
为什么CPU降这么多?同步版本中,线程切换和GIL竞争导致大量CPU时间浪费在上下文切换上;异步版本单线程事件循环,IO等待时不占CPU。注意我们没有增加worker数量,还是4个Gunicorn worker,只是每个worker内部用异步。
7. 总结与建议
asyncio不是银弹,它适合IO密集型任务,如果你的业务是CPU计算密集(比如图像处理),该用多进程还是用多进程。但Web API绝大多数瓶颈都在IO——数据库、HTTP调用、文件读写。
几点实战建议:
- 别把Flask整个变成异步,Flask的同步模型在3.10下稳得很,只需把耗时的IO操作异步化。
- 信号量必须全局唯一,定义在协程函数外面。
- 所有子协程都要有超时,用
asyncio.wait_for(fetch_order(...), timeout=2),防止下游服务挂了导致你的事件循环被拖死。 - 监控事件循环的延迟,可以加一个每秒检查
loop.time()的协程,如果延迟超过500ms说明有地方阻塞了。
现在这个接口稳定运行了一个月,平均耗时1.1秒,P99在2.3秒。如果你也在用requests+Flask搭接口,试试这个改造路径,收益还是很明显的。