一、问题背景:一个「看起来没问题」的查询接口
那是周三下午,运营反馈「订单导出功能卡死」。我查了监控,发现/api/v1/orders/batch这个接口在200并发下:
- 平均RT:305ms(其中数据库查询占270ms)
- QPS:327
- PostgreSQL连接数:飙到150+,接近max_connections=200
当时的代码逻辑很「标准」——Flask同步视图,内部用psycopg2连接PostgreSQL查询订单主表,然后循环调用第三方ERP接口获取物流状态。问题就出在这个循环上:每个订单要串行请求一次ERP,平均耗时80ms。如果有50个订单,光这层就是4秒。
同步模型下,线程池炸了,连接池也炸了。我当时的内心OS:这代码没写错,就是架构错了。
二、环境与版本
改造前先交代环境,方便读者复现:
Python: 3.10.11 (需要3.8+才能完整使用asyncio.run)
Web框架: Flask 2.2.5 (同步,通过gunicorn运行)
HTTP客户端: requests 2.28.2 (同步)
数据库驱动: psycopg2-binary 2.9.6 (同步)
压测工具: wrk 4.2.0
服务器: 4C8G 阿里云ECS,PostgreSQL 14 (max_connections=200)
gunicorn配置: gunicorn -w 8 -k gthread --threads 16 app:app
改造后引入:
- aiohttp 3.8.4 (或httpx 0.24.1,本文用httpx)
- asyncpg 0.27.0
- uvloop 0.17.0 (事件循环替换)
三、方案设计:异步化不是「换框架」,是「换I/O模型」
很多人一谈异步就想到换FastAPI。但你完全可以保留Flask,用asyncio.run()在视图内部跑协程。思路分三层:
- 数据库层:psycopg2 → asyncpg。连接池从SQLAlchemy的ThreadPool切换为asyncpg内置Pool,限制max_size=10。
- 外部API调用层:requests → httpx.AsyncClient。用
asyncio.gather并发请求ERP,而不是for循环串行。 - 事件循环:安装uvloop,替换默认的asyncio事件循环。实测uvloop在纯网络I/O场景下能减少约15%开销。
核心设计图(文字版):
同步版本: 请求进来 -> 线程池分配线程 -> psycopg2阻塞查询 -> for循环requests阻塞调用ERP -> 返回
异步版本: 请求进来 -> 主线程asyncio.run -> asyncpg异步查询 -> gather并发调用ERP -> 返回
关键点:异步化要解决的是「等待时间」问题,而不是「计算时间」问题。如果你的逻辑是CPU密集,异步没用。
四、核心实现:before/after代码对比
4.1 同步版本(改造前)
# app_sync.py — 改造前代码
from flask import Flask, jsonify, request
import psycopg2
import requests
from psycopg2.pool import ThreadedConnectionPool
app = Flask(__name__)
pool = ThreadedConnectionPool(5, 20, host='localhost', dbname='orders')
@app.route('/api/v1/orders/batch', methods=['GET'])
def batch_query():
order_ids = request.args.get('ids', '').split(',')
if len(order_ids) > 50:
return jsonify({'error': 'max 50 ids'}), 400
conn = pool.getconn()
try:
# 1. 查询主表
cur = conn.cursor()
cur.execute("SELECT id, amount FROM orders WHERE id = ANY(%s)", (order_ids,))
rows = cur.fetchall()
finally:
pool.putconn(conn)
# 2. 串行调用ERP — 这里是最痛点
erp_results = []
for oid in order_ids:
resp = requests.get(f'http://erp.example.com/status/{oid}', timeout=3)
erp_results.append(resp.json())
return jsonify({'orders': rows, 'erp': erp_results})
if __name__ == '__main__':
app.run(threaded=True)
这个版本的性能瓶颈一目了然:
- 数据库查询阻塞线程(线程池16个,但200并发全排队)
- ERP调用是「串行」的,50个订单 = 50次RTT
4.2 异步重构版本(改造后)
# app_async.py — 改造后核心代码
import asyncio
import uvloop
from flask import Flask, jsonify, request
import httpx
import asyncpg
app = Flask(__name__)
asyncio.set_event_loop_policy(uvloop.EventLoopPolicy())
# 全局初始化连接池 (在app启动时调用)
async def init_db_pool():
return await asyncpg.create_pool(
host='localhost', database='orders',
min_size=5, max_size=10, # 关键:限制连接数,避免打爆数据库
command_timeout=5
)
async def fetch_erp_status(client: httpx.AsyncClient, oid: str, sem: asyncio.Semaphore):
"""并发获取ERP状态,使用信号量限制并发数"""
async with sem: # 限制ERP并发不超过20
try:
resp = await client.get(f'http://erp.example.com/status/{oid}', timeout=3)
return {'oid': oid, 'status': resp.json().get('status')}
except httpx.TimeoutException:
return {'oid': oid, 'status': 'timeout'}
async def handle_batch(order_ids):
# 数据库异步查询
pool = await init_db_pool()
async with pool.acquire() as conn:
rows = await conn.fetch(
"SELECT id, amount FROM orders WHERE id = ANY($1)", order_ids
)
# 并发ERP调用 — 核心优化点
async with httpx.AsyncClient(timeout=3) as client:
sem = asyncio.Semaphore(20) # 防止打爆ERP
tasks = [fetch_erp_status(client, oid, sem) for oid in order_ids]
erp_results = await asyncio.gather(*tasks)
return {'orders': [dict(r) for r in rows], 'erp': erp_results}
@app.route('/api/v1/orders/batch', methods=['GET'])
def batch_query():
order_ids = request.args.get('ids', '').split(',')
if len(order_ids) > 50:
return jsonify({'error': 'max 50 ids'}), 400
# Flask同步视图内运行事件循环
result = asyncio.run(handle_batch(order_ids))
return jsonify(result)
# 重要:需要关闭已有的asyncpg pool,避免连接泄漏
# 生产用uvicorn或hypercorn跑ASGI,这里为演示保留Flask
注意几个细节:
- asyncio.run() 每次创建新事件循环,不能用于高并发生产(每次都要重建连接池)。实际部署请用FastAPI/Quart或AIOHTTP,或者用gunicorn的uvicorn.workers.UvicornWorker。我这里为了对比才用这种「混合模式」。
- asyncio.Semaphore(20)必须有,否则100个订单会同时打ERP,直接给人家打挂。
- asyncpg的fetch返回Record对象,需要dict(r)转JSON。
五、踩坑与优化:三个让我失眠的bug
5.1 坑1:连接池泄漏导致「偶尔卡死」
第一次写完跑测试,前5000请求正常,然后全部超时。排查半天发现是asyncio.run()每次都会新建pool,而init_db_pool()在函数内被反复调用,旧连接没被关闭。最终方案:把pool声明为global,启动时初始化一次,请求时复用。
# 正确做法 —— 在app.py顶层:
_pool = None
async def get_pool():
global _pool
if _pool is None:
_pool = await asyncpg.create_pool(min_size=5, max_size=10)
return _pool
5.2 坑2:uvloop和asyncpg的SSL冲突
我们数据库开了SSL,uvloop下asyncpg报错ssl.SSLError。查了issue发现是asyncpg 0.27.0的已知bug,升级到0.28.0解决(或者禁用uvloop)。最终选择升级asyncpg,因为uvloop带来的性能收益太香。
5.3 坑3:wrk压测时「连接数看起来没降」
改完压测,看PostgreSQL pg_stat_activity傻眼了——连接数还是150。后来发现是压测客户端自己开了太多连接。调整wrk参数-c 200是指并发连接数,而gunicorn线程模式也有影响。最终用gunicorn -w 4 -k uvicorn.workers.UvicornWorker跑ASGI模式才达到理想效果。
六、效果数据:压测对比(wrk -t8 -c200 -d30s)
| 指标 | 同步版本 | asyncio版本 | 提升 |
|---|---|---|---|
| QPS | 327 | 1892 | 4.78倍 |
| 平均RT | 305ms | 52ms | 5.87倍 |
| P99 RT | 890ms | 143ms | 6.2倍 |
| PostgreSQL连接数 | 150+ | 10 | -93% |
| 错误率 | 4.5% | 0% | — |
关键数据说明:
- PostgreSQL连接数从150降到10,是因为asyncpg的max_size=10强制复用连接。数据库CPU从85%降到32%。
- ERP调用并发从「串行50次」变成「并发20路」,单接口耗时从4秒(5080ms)降到约400ms(50/2080ms),瓶颈变为数据库查询。
七、总结与建议
异步化适合「I/O密集 + 并发高」的场景。如果你遇到类似问题,按这个顺序排查:
- 先确认瓶颈在I/O(数据库查询慢、外部API调用慢),而非CPU计算
- 从最耗时的依赖开始替换——我的案例里是ERP串行调用,而不是数据库
- 控制并发度:Semaphore和连接池max_size是防身符
- 别用
asyncio.run处理真实请求——它每次新建事件循环,代价太高。生产环境用ASGI框架或worker
最后说点实在的:如果你在维护老Flask项目,没必要全量迁移FastAPI。Flask视图内嵌asyncio.run是过渡方案,但长期来看,建议将高并发接口单独用FastAPI写,或者整体切换到Quart(Flask的异步版本)。
数据不会骗人——当你的QPS从300提升到1800,那种爽感是写同步代码永远体会不到的。