跳转至

模型池架构

同一组「逻辑模型」(如「主图生成」)背后挂多个物理模型,加权随机调度 + 熔断故障转移 + 共享 API 账号。突破单账号 RPS、单模型熔断雪崩两个瓶颈。

一、问题与目标

现状(模型池前)

用户 → API → payload.ai_model_id → worker → 单物理模型 → 单上游账号

问题:

  1. 单账号 RPS 上限:火山方舟单账号约 50 RPS。100 用户同时生图 → 单账号被打爆 → 大量 429。
  2. 单模型故障无转移:模型 A 异常 → 该模型所有任务失败重试 → 雪崩。
  3. 无法利用多账号:DB 里 AiModelApiKey 是 1:N,但同模型的多 Key 是顺序兜底,没有 RPS 聚合。
  4. 配置分散:每个模型独立配 Key、独立设并发,运维成本高。

目标

目标 解决方式
突破单账号 RPS 多账号共享(同一账号可被多个模型绑定)
故障自动转移 加权随机 + 熔断器(连续失败自动隔离,定时探测恢复)
配置集中化 引入「池」概念,按业务场景(main_image / detail_image)聚合
向后兼容 ai_model_id 入参继续工作(视为池中单成员)

二、核心概念

模型池(AiModelGroup)

逻辑模型,按业务场景命名(如 main_imagedetail_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)统一编排:

  1. ModelPoolRouter.pick / pick_excluding
  2. 调用方 prepare(selection) 拼装上游参数
  3. (可选)AccountSemaphore + ConcurrencyGuard
  4. NonlinearService.generate(含瞬时错误重试)
  5. 成功 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

集成测试场景(手动)

  1. 创建池 main_image,加入 3 个模型(权重 1/3/6)
  2. 模型 A 连续失败 5 次 → 观察 UI 中 A 变红(OPEN)
  3. 后续 1000 次请求全部落到 B/C,分布近似 1:2(B:C = 3:6)
  4. 60s 后 A 半开 → 探测成功 → A 恢复绿色(CLOSED)
  5. 全部熔断 → NoAvailableModelError → 死信队列

十一、演进路径

阶段 内容 状态
一期 进程内熔断 + 加权随机 + 共享账号 + Admin UI ✅ 已完成
二期 Redis 共享熔断状态(多 Pod 一致) 📋 计划
三期 动态权重(按成功率自动调整 weight) 📋 计划
四期 跨 region 账号路由(火山不同 region RPS 独立) 📋 计划

二期:Redis 共享熔断

一期痛点:N 个 Pod 各自维护熔断器,单模型真实失败次数被分散,N 个 Pod 各自达到阈值才会熔断。

二期方案:BreakerRegistry 改用 Redis 计数(INCR + EXPIRE),所有 Pod 共享同一计数器。