共计 2144 个字符,预计需要花费 6 分钟才能阅读完成。
问题定义:原始架构的性能瓶颈
在初始同步调用架构下,当 RPS(Requests Per Second)超过 100 时,系统出现明显性能劣化。通过 APM 工具观测到以下典型问题:

- 延迟突增:95 分位响应时间从 120ms 飙升至 1.2s
- 内存泄漏:每次模型推理后 Python 进程 RSS 增长 2 -3MB
- 阻塞调用:数据库查询阻塞事件循环导致级联延迟
架构演进:从同步到事件驱动
通过 ab 工具对两种架构进行对比测试(4 核 8G 云主机):
| 架构类型 | 最大 QPS | 平均延迟 | CPU 利用率 |
|---|---|---|---|
| 同步阻塞 | 850 | 210ms | 92% |
| 事件驱动(本方案) | 3200 | 68ms | 78% |
关键改进点在于:
1. 使用 uvicorn 替代 gunicorn 作为 ASGI 服务器
2. 将模型调用封装为独立线程池任务
3. 采用 Redis Stream 实现背压 (backpressure) 控制
核心实现
异步网关层实现
from fastapi import FastAPI, Depends
from fastapi.security import OAuth2PasswordBearer
app = FastAPI()
oauth2_scheme = OAuth2PasswordBearer(tokenUrl="token")
@app.post("/v1/predict")
async def predict(
text: str,
token: str = Depends(oauth2_scheme)
):
"""处理推理请求的异步端点"""
# JWT 验证逻辑
payload = verify_jwt(token)
# 异步调用模型
result = await model_dispatcher.dispatch(payload['model'], text)
return {"result": result}
模型负载均衡算法
采用动态加权轮询算法,关键实现逻辑:
- 通过心跳机制维护健康节点列表
- 根据最近 5 次响应时间计算动态权重
- 异常节点自动熔断 15 秒
class ModelPool:
def __init__(self):
self._nodes = []
self._weights = []
async def dispatch(self, model_id: str, text: str) -> str:
"""选择最优节点进行推理"""
node = self._select_node()
start = time.time()
try:
resp = await node.predict(text)
self._update_weights(node, time.time() - start)
return resp
except Exception as e:
self._handle_failure(node)
raise
会话状态缓存设计
Redis 数据结构采用 Hash 存储会话上下文:
HSET session:{session_id}
"context" "{json_encoded_context}"
"timestamp" "{unix_time}"
TTL 优化策略:
– 基础 TTL 设置为 300 秒
– 每次访问延长 60 秒
– 凌晨 4 点执行全局扫描清理
生产验证
压力测试数据(Locust)
| 并发用户数 | RPS | 平均延迟 | 95 分位延迟 |
|---|---|---|---|
| 100 | 1200 | 83ms | 142ms |
| 500 | 3100 | 161ms | 198ms |
幂等处理实现
def idempotency(key_func=None):
"""防止重复请求的装饰器"""
def decorator(f):
@wraps(f)
async def wrapper(*args, **kwargs):
request = kwargs.get('request')
if not request:
return await f(*args, **kwargs)
redis_key = f"idem:{key_func(request)}" if key_func else f"idem:{request.url}"
with redis.lock(redis_key, timeout=10):
if redis.get(redis_key):
raise HTTPException(409, "Duplicate request")
redis.setex(redis_key, 60, "1")
return await f(*args, **kwargs)
return wrapper
return decorator
避坑指南
-
TCP 参数优化:
# 必须设置的 sysctl 参数 net.ipv4.tcp_keepalive_time = 60 net.ipv4.tcp_keepalive_intvl = 10 net.ipv4.tcp_keepalive_probes = 6 -
冷启动预热:
- 启动时并行加载高频模型
- 预填充缓存热点数据
-
梯度增加流量(从 10% 开始)
-
结构化日志规范:
{ "timestamp": "ISO8601", "level": "ERROR", "trace_id": "uuid4", "module": "model_pool", "metrics": { "queue_size": 12, "model_latency": 0.45 } }
开放问题
当前架构在单 AZ 运行良好,但如何设计跨可用区的容灾方案?需要考虑:
– 会话状态的多 AZ 同步
– 模型权重的一致性保证
– 区域故障的自动检测切换
欢迎在评论区分享你的分布式系统设计经验。
正文完
