模型池架构¶
同一组「逻辑模型」(如「主图生成」)背后挂多个物理模型,加权随机调度 + 熔断故障转移 + 共享 API 账号。突破单账号 RPS、单模型熔断雪崩两个瓶颈。
一、问题与目标¶
现状(模型池前)¶
用户 → API → payload.ai_model_id → worker → 单物理模型 → 单上游账号
问题:
- 单账号 RPS 上限:火山方舟单账号约 50 RPS。100 用户同时生图 → 单账号被打爆 → 大量 429。
- 单模型故障无转移:模型 A 异常 → 该模型所有任务失败重试 → 雪崩。
- 无法利用多账号:DB 里
AiModelApiKey是 1:N,但同模型的多 Key 是顺序兜底,没有 RPS 聚合。 - 配置分散:每个模型独立配 Key、独立设并发,运维成本高。
目标¶
| 目标 | 解决方式 |
|---|---|
| 突破单账号 RPS | 多账号共享(同一账号可被多个模型绑定) |
| 故障自动转移 | 加权随机 + 熔断器(连续失败自动隔离,定时探测恢复) |
| 配置集中化 | 引入「池」概念,按业务场景(main_image / detail_image)聚合 |
| 向后兼容 | 老 ai_model_id 入参继续工作(视为池中单成员) |
二、核心概念¶
模型池(AiModelGroup)¶
逻辑模型,按业务场景命名(如 main_image、detail_image)。同池内的多个物理模型视为可互换。
| 字段 | 含义 |
|---|---|
group_key |
业务唯一键,worker payload 用此路由 |
circuit_failure_threshold |
单物理模型连续失败多少次触发熔断(默认 5) |
circuit_recovery_seconds |
熔断后多久进入半开探测(默认 60s) |
物理模型(AiModel)¶
新增 3 列:
| 字段 | 含义 |
|---|---|
group_id |
所属池(外键,SET NULL) |
weight |
池内加权随机的权重(≥1,默认 1) |
priority |
同池内优先级分层(高优先用满才用低优先,默认 0) |
共享账号(AiModelApiAccount)¶
独立的上游账号,可被多个物理模型绑定(多对多)。突破单账号 RPS。
| 字段 | 含义 |
|---|---|
account_key |
API Key(应用层加密,UI 脱敏) |
provider |
账号归属的供应商 |
max_rps |
单账号并发上限(火山单账号 50 RPS) |
key_version |
加密版本号(用于轮换密钥) |
绑定关系(AiModelAccountBinding)¶
模型 ↔ 账号 多对多,带 priority:同模型多账号时按优先级用,失败降级。
三、调度算法¶
pick(group_key)¶
1. 查询 group 下所有 is_active 物理模型
2. 用 BreakerRegistry 过滤熔断中的模型(breaker.allow() == False)
3. 按 priority 分层(高优先层用满才用低优先层)
4. 层内按 weight 加权随机:random.choices(layer, weights=[m.weight for m in layer])
5. 选中模型 → 按 binding priority 取第一个可用账号
6. 全部熔断 → 抛 NoAvailableModelError(worker 转死信,不重试)
故障转移示例¶
池 main_image 有 3 个模型,权重 [1, 3, 6]:
T0: 正常调度 → 90% 选 C, 30% 选 B, 10% 选 A
T1: A 连续失败 5 次 → A 熔断 → 后续只在 B/C 间加权随机
T2: C 也熔断 → 只用 B
T3: B 也熔断 → NoAvailableModelError → worker 转死信
T4: 60s 后 A 半开探测 → 成功 → A 重新进池
四、熔断器状态机¶
连续失败 ≥ threshold
┌────────────────────────┐
│ ▼
CLOSED ◄─── 探测成功 ─── HALF_OPEN ──── 经过 recovery_seconds ──── OPEN
│ │ │
│ │ 探测失败 │
│ └────────────────────────────────────────┘
│
└─ 成功立即清零失败计数
| 状态 | 行为 |
|---|---|
| CLOSED | 正常放行所有请求 |
| OPEN | 拒绝所有请求(被池调度过滤) |
| HALF_OPEN | 仅放行一次探测,期间拒绝其他;探测成功 → CLOSED,失败 → 重新 OPEN |
实现:backend/core/services/circuit_breaker.py,单进程内 asyncio.Lock 保护状态转换。
五、并发信号量¶
双层结构(保持原有)¶
用户请求
│
▼
[用户级 sem] ─── 配额来自 MemberPlan.concurrent_tasks
│
▼
[模型/池级 sem] ─── 容量来自 AiModel.max_concurrency(兼容)或 sum(池成员 max_concurrency)(池化)
│
▼
上游 API
池级聚合¶
ConcurrencyGuard.register_pool(pool_key, total_concurrency) 创建池级信号量,容量为池内成员 max_concurrency 之和。
# 池内 3 个模型,分别 20/30/10 并发 → 池级 sem 容量 = 60
guard.register_pool("main_image", total_concurrency=60)
# worker 内:
async with guard.acquire(user_uuid, "pool:main_image"):
result = await nonlinear.generate(...)
六、数据流¶
flowchart LR
User["用户请求<br/>task_type + ai_model_id"] --> API["backend API"]
API --> Svc["TemplateBatchGenerateService"]
Svc --> Pub["publish_llm_job<br/>payload.model_group"]
Pub --> Stream[("stream:llm")]
Stream --> Worker["Worker.run_once"]
Worker --> GW["ModelGateway.generate"]
Canvas["Canvas 同步生图"] --> GW
GW --> Router["ModelPoolRouter.pick"]
Router --> Break["BreakerRegistry<br/>过滤熔断模型"]
Break --> Pick["加权随机<br/>+ 选账号"]
Pick --> Sem["AccountSemaphore<br/>+ ConcurrencyGuard"]
Sem --> Upstream["火山/OpenAI/ZenMux"]
Upstream -.->|"成功/失败"| GW
GW -.->|"record_result / force_open"| Break
ModelGateway 调用链¶
内部 facade(model_gateway.py)统一编排:
ModelPoolRouter.pick/pick_excluding- 调用方
prepare(selection)拼装上游参数 - (可选)
AccountSemaphore+ConcurrencyGuard NonlinearService.generate(含瞬时错误重试)- 成功
record_result(True);失败区分信号量拥堵 / fatal / 普通失败后换成员
Worker 批量生图与 Canvas 文生图均经此入口;Admin mock-generate-probe 仍只选路、不走 Gateway。
七、Admin UI¶
模型池页(/admin/model-pools)¶
- 卡片列表:每个池显示「组名 / 成员数 / 总并发 / 熔断阈值 / 恢复时间」
- 点击池 → 展开成员表(实时熔断状态彩色徽章 + 连续失败次数,5s 轮询刷新)
- 编辑池:可调 circuit_failure_threshold / circuit_recovery_seconds
共享账号页(/admin/api-accounts)¶
- 账号列表(脱敏 Key):标签 / 供应商 / RPS / 每日上限 / 已绑定模型数
- 绑定模型:选择模型 + 优先级
- 解绑
AI 模型页(已改造)¶
表格新增「所属池」「权重」「优先级」三列,编辑行可设置 group_id / weight / priority。
八、关键决策与权衡¶
| 决策 | 选择 | 理由 |
|---|---|---|
| 池粒度 | 按 model_group |
与现有 image_models 配置语义一致,运维熟悉 |
| 选择策略 | 加权随机 | 比 round-robin 更灵活,权重可在不重启下动态调;比 hash 更适合故障转移 |
| Key 共享 | 独立 api_accounts 表 + 多对多 binding |
突破单账号 RPS;老 api_keys 表保留向下兼容 |
| 熔断器 | 进程内单例(一期) | 实现简单;多 Pod 时改 Redis 共享状态(二期) |
| 加权随机粒度 | 按 priority 分层,层内 weight | 兼顾「主备」(priority)和「负载均衡」(weight)两种语义 |
| 兼容性 | payload 同时塞 ai_model_id + model_group |
老 payload 无 model_group 时 worker 回退到单模型路径 |
九、配置与运维¶
创建池(SQL 示例)¶
INSERT INTO ai_model_groups (id, group_key, name, circuit_failure_threshold, circuit_recovery_seconds)
VALUES (gen_random_uuid(), 'main_image', '主图生成池', 5, 60);
-- 把现有模型加入池
UPDATE ai_models SET group_id = ?, weight = 1, priority = 0 WHERE id = ?;
创建共享账号¶
INSERT INTO ai_model_api_accounts (id, account_key, label, provider, max_rps)
VALUES (gen_random_uuid(), 'volc-xxx', '火山账号 A', 'volcengine_ark', 50);
-- 绑定到模型
INSERT INTO ai_model_account_bindings (model_id, account_id, priority)
VALUES (?, ?, 0);
调整熔断参数¶
无需重启:编辑池的 circuit_failure_threshold / circuit_recovery_seconds 即可。
已创建的熔断器参数不会自动更新,需要调用 BreakerRegistry.reconfigure(model_id, ...) 或重启 worker。
手动恢复熔断¶
# 重置单个模型
breaker_registry.reset(model_id)
# 重置全部
breaker_registry.reset()
十、压测与验证¶
单元测试¶
test_circuit_breaker.py(9 个):状态机全覆盖test_pool_weighted_random.py(7 个):分布近似、优先级分层、空候选test_model_pool_failover.py(9 个):熔断隔离、池级 sem 阻塞、记录回调
poetry run pytest tests/unit/backend/core/services/test_circuit_breaker.py \
tests/unit/backend/core/services/test_pool_weighted_random.py \
tests/unit/backend/core/services/test_model_pool_failover.py -v
集成测试场景(手动)¶
- 创建池
main_image,加入 3 个模型(权重 1/3/6) - 模型 A 连续失败 5 次 → 观察 UI 中 A 变红(OPEN)
- 后续 1000 次请求全部落到 B/C,分布近似 1:2(B:C = 3:6)
- 60s 后 A 半开 → 探测成功 → A 恢复绿色(CLOSED)
- 全部熔断 → NoAvailableModelError → 死信队列
十一、演进路径¶
| 阶段 | 内容 | 状态 |
|---|---|---|
| 一期 | 进程内熔断 + 加权随机 + 共享账号 + Admin UI | ✅ 已完成 |
| 二期 | Redis 共享熔断状态(多 Pod 一致) | 📋 计划 |
| 三期 | 动态权重(按成功率自动调整 weight) | 📋 计划 |
| 四期 | 跨 region 账号路由(火山不同 region RPS 独立) | 📋 计划 |
二期:Redis 共享熔断¶
一期痛点:N 个 Pod 各自维护熔断器,单模型真实失败次数被分散,N 个 Pod 各自达到阈值才会熔断。
二期方案:BreakerRegistry 改用 Redis 计数(INCR + EXPIRE),所有 Pod 共享同一计数器。