1. 问题背景:一个慢如蜗牛的订单查询接口
事情要从今年Q2的代码review说起。我们负责的订单服务有一个 /api/orders/batch 接口,用于批量查询订单详情(每次最多20个订单)。逻辑很简单:接收订单ID列表 → 循环调用下游三个服务(用户信息、商品信息、物流信息) → 聚合返回。
压测时发现:订单量一旦超过5个,接口响应时间就线性飙升。20个订单的场景下,平均响应时间高达1.6s,P99更是到了3.2s。而这三个下游服务的响应时间其实都只有30-50ms——问题出在同步串行调用上。
伪代码大概是这样的:
# 同步版本(问题代码)
def get_orders_detail(order_ids):
results = []
for oid in order_ids:
user = requests.get(f"http://user-svc/{oid}").json()
product = requests.get(f"http://product-svc/{oid}").json()
logistics = requests.get(f"http://logistics-svc/{oid}").json()
results.append({...})
return results
20个订单 × 3个服务 = 60次串行HTTP调用,每次等待30ms就是1800ms。这还没算序列化和网络开销。
2. 环境与版本:我们踩过的坑都跟版本有关
先交代当时的运行环境:
- Python 3.9.12(3.9以下asyncio.run有坑,后面讲)
- Flask 2.2.3(同步框架,需要结合ASGI或线程池)
- aiohttp 3.8.4(用于异步HTTP客户端)
- databases 0.7.0(异步数据库驱动,基于SQLAlchemy Core)
- Redis 6.2 + aioredis 2.0.1(后来发现aioredis被合并到redis-py 4.x)
- 服务器:4核8G的阿里云ECS,CentOS 7.9
重要提示:如果你的Python版本低于3.8,asyncio的很多特性(如asyncio.Runner)不可用。我们最初在3.7上尝试,发现asyncio.run()不能嵌套调用,后来升级到3.9才解决。
3. 方案设计:从同步到异步的三个关键改造
核心思路:将Flask的同步处理函数改为协程+异步IO。具体分三步:
- 网络IO异步化:用aiohttp替换requests,实现HTTP调用的并发等待
- 数据库IO异步化:用databases库连接PostgreSQL,异步执行查询
- 并发控制:用asyncio.Semaphore限制并发数,防止下游服务被打挂
架构图简化为:Flask(通过Gunicorn + uvicorn worker运行)→ 协程调度 → aiohttp异步请求 + databases异步查询 → 聚合返回
这里要特别注意:Flask本身不支持异步,所以我们用Gunicorn的uvicorn worker来运行Flask应用。配置如下:
gunicorn -k uvicorn.workers.UvicornWorker -w 4 --timeout 120 app:app
每个worker内部运行事件循环,处理异步请求。
4. 核心实现:代码对比与关键细节
4.1 异步HTTP调用(核心改造)
改造前的同步代码(已简化):
import requests
def fetch_order_detail(order_id):
user = requests.get(f"http://user-svc/{order_id}", timeout=2).json()
product = requests.get(f"http://product-svc/{order_id}", timeout=2).json()
logistics = requests.get(f"http://logistics-svc/{order_id}", timeout=2).json()
return {"user": user, "product": product, "logistics": logistics}
改造后的异步代码:
import asyncio
import aiohttp
from aiohttp import ClientTimeout, TCPConnector
# 全局连接池复用,避免每次创建新连接
connector = TCPConnector(limit=100, limit_per_host=30, ttl_dns_cache=300)
timeout = ClientTimeout(total=10)
async def fetch_order_detail_async(session, order_id, sem):
async with sem: # 并发控制:同时最多20个协程执行HTTP请求
async with session.get(f"http://user-svc/{order_id}", timeout=timeout) as resp:
user = await resp.json()
async with session.get(f"http://product-svc/{order_id}", timeout=timeout) as resp:
product = await resp.json()
async with session.get(f"http://logistics-svc/{order_id}", timeout=timeout) as resp:
logistics = await resp.json()
return {"user": user, "product": product, "logistics": logistics}
async def batch_fetch(order_ids):
sem = asyncio.Semaphore(20) # 控制并发度为20
async with aiohttp.ClientSession(connector=connector) as session:
tasks = [fetch_order_detail_async(session, oid, sem) for oid in order_ids]
results = await asyncio.gather(*tasks, return_exceptions=True)
# 处理异常:如果某个任务失败,记录日志并返回空
return [r if not isinstance(r, Exception) else None for r in results]
关键细节:
- TCPConnector(limit=100): 限制总连接数,防止文件描述符耗尽
- limit_per_host=30: 对每个下游服务的连接数限制,避免压垮一个服务
- ttl_dns_cache=300: DNS缓存5分钟,减少DNS解析开销
- ClientTimeout(total=10): 全局超时10秒,防止个别慢请求拖垮整体
- asyncio.Semaphore(20): 控制并发度20,这是根据下游服务容量测出来的最佳值
4.2 异步数据库查询
同步数据库操作(已简化):
import psycopg2
def get_order_from_db(order_id):
conn = psycopg2.connect(dsn="postgresql://...")
cur = conn.cursor()
cur.execute("SELECT * FROM orders WHERE id=%s", (order_id,))
return cur.fetchone()
改为异步:
import databases
from databases import Database
database = Database("postgresql+asyncpg://user:pass@host/db", min_size=5, max_size=20)
async def get_order_from_db_async(order_id):
query = "SELECT * FROM orders WHERE id = :oid"
return await database.fetch_one(query=query, values={"oid": order_id})
# 在应用启动时连接,关闭时断开
async def startup():
await database.connect()
async def shutdown():
await database.disconnect()
注意:databases库底层使用asyncpg,连接池大小max_size=20需要与并发度匹配。我们一开始设了100,结果把PostgreSQL的连接数打满了,后来调小到20配合Semaphore才稳定。
5. 踩坑与优化:三个差点让我放弃的坑
坑1:Flask异步视图的兼容性问题
Flask 2.0+支持异步视图函数,但前提是使用ASGI服务器(如uvicorn)。我们最初用flask run调试,发现异步函数根本不执行,返回空值。后来查文档才知道,Flask的dev server是Werkzeug,不支持异步。解决方案:开发环境也用uvicorn app:app启动,或者用flask run --with-threads(但不推荐)。
坑2:aiohttp连接池泄漏
上线第一天,下游服务反馈说我们的请求突然变少了。查日志发现aiohttp报错Too many open connections。原因是我们在每个请求里都创建了新的ClientSession,没有复用。修复:将ClientSession实例化为全局变量,加上connector限制连接数。
坑3:asyncio.gather的异常处理
最初代码直接返回await asyncio.gather(*tasks),结果一个订单查询失败抛异常,整个batch都挂了。修复:加上return_exceptions=True,手动处理异常并返回None。同时增加重试逻辑:
retry_count = 3
for attempt in range(retry_count):
try:
return await fetch_order_detail_async(...)
except Exception as e:
if attempt == retry_count - 1:
logger.error(f"Order {order_id} failed after {retry_count} retries: {e}")
return None
await asyncio.sleep(0.1 * (2 ** attempt)) # 指数退避
6. 效果数据:性能提升的量化对比
压测工具:locust 2.15.1,模拟100个并发用户,持续压测5分钟。
| 指标 | 同步版本 | 异步版本(无并发控制) | 异步版本(Semaphore=20) |
|---|---|---|---|
| QPS | 152 | 3200 | 2850 |
| 平均延迟 | 820ms | 48ms | 62ms |
| P50延迟 | 680ms | 35ms | 45ms |
| P99延迟 | 2.1s | 180ms | 120ms |
| 错误率 | 0.3% | 2.1%(连接超时) | 0.05% |
解读:
- 无并发控制的异步版本QPS最高(3200),但错误率也高(2.1%),因为下游服务扛不住突发的并发请求,导致大量连接超时。
- 加上Semaphore=20后,QPS略降(2850),但错误率降到0.05%,P99延迟反而更低(120ms vs 180ms),因为避免了拥塞崩溃。
- 对比同步版本,QPS提升了18倍,P99延迟降低了17.5倍。
资源消耗对比(4核8G服务器):
- 同步版本:CPU 65%,内存 320MB
- 异步版本:CPU 45%,内存 280MB
异步版本不仅性能更好,资源占用反而更低,因为减少了线程切换开销。
7. 总结:异步编程的适用边界
这次改造让我深刻理解了asyncio的适用场景:
- IO密集型任务(HTTP、DB、文件读写)效果显著,计算密集型任务效果有限
- 并发度需要控制,不是越大越好。我们最终定在20,是根据下游服务的TP99和连接池大小计算出来的
- 错误处理要严谨,异步的异常传播更加隐蔽,务必加上return_exceptions和重试机制
- 版本兼容性:Python 3.8+、aiohttp 3.8+、databases 0.7+,低于这些版本容易踩坑
最后给个建议:如果你的Flask应用遇到IO瓶颈,别急着上celery或消息队列,先试试asyncio——成本最低,效果往往出人意料。前提是做好连接池管理和并发控制。
(全文完)