跳转至

Messaging / Typed Worker 架构核查报告

范围:8 类 typed worker 的 Redis Stream 消费机制、队列稳定性与健壮性、并发设计。 结论来源:代码走查(messaging/ + backend/worker_main.py + deploy/k3s/base/backend/workers/deployments.yaml)。 状态:2026-07-31 初版;P0/P1/P2 修复计划见文末落地清单。

1. 数据流

flowchart LR
    subgraph API[API 进程]
        pub["publish_template_batch_llm_job<br/>task_type → worker_type"]
    end
    subgraph redis[Redis]
        s1["stream:llm:{wt}:{user_id}<br/>每用户独立 stream"]
        a1["active:llm:{wt}:normal (SET)"]
        a2["active:llm:{wt}:priority (ZSET)"]
        r1["stream:retry:* 旁路计数"]
    end
    subgraph worker[8× typed worker]
        fc["FairConsumer<br/>选用户→PEL→新消息→流水线"]
        gd["ConcurrencyGuard<br/>用户+模型 Semaphore"]
        as["AccountSemaphore<br/>Redis 分布式 max_rps"]
    end
    pub --> s1
    pub --> a1
    pub --> a2
    fc --> s1
    fc --> a1
    fc --> a2
    fc --> r1
    fc --> gd
    fc --> as

关键设计:

  • 每用户独立 streamstream:llm:{worker_type}:{user_id}):避免单 stream 热点与队头阻塞。
  • 活跃用户注册表active:llm:{wt}:normal|priority):提供 O(1) 的「哪些用户有消息」索引,避免 SCAN 全库。
  • 80/20 公平抽取:priority 槽位 80%、normal 20%,高等级优先且普通用户不饿死。
  • 流水线并发run_once 把在途任务挂到 asyncio.create_task 后立即返回,按 concurrency - inflight 补位,避免「首张占满整轮 gather」。
  • 失败保留 PEL + 旁路计数:失败不 ACK,下轮 id=0 回收;用 stream:retry:* 旁路 key 记录重试次数(stream 字段不可变)。
  • 空 stream 清理:即时(消费空即 unregister)+ 定时兜底(每 10 分钟 SCAN 清理,Redis 7 lag=0 判断)。

2. 并发设计现状

实现 范围 备注
Worker 并发 FairConsumer._concurrencyMESSAGING_WORKER_CONCURRENCY 单进程 在途任务集 + 用户级 in-flight 计数
用户配额 ConcurrencyGuard + UserConcurrencyService 进程内 asyncio.Semaphore 每用户 concurrent_tasks;配额 5 分钟缓存
模型保护 ConcurrencyGuard.register_model 进程内 asyncio.Semaphore AiModel.max_concurrency / 池成员并发之和
账号 RPS AccountSemaphore Redis 分布式(Lua) 跨 Pod 共享,防止同一账号超 max_rps
视频上游 RedisSemaphoresemaphore:video:* Redis 分布式(Lua) I2I/I2V/ffmpeg 三档,Temporal 内使用
公平调度 ActiveUserRegistry(SET/ZSET + ZRANDMEMBER/SRANDMEMBER) Redis priority 权重未真正加权(ZRANDMEMBER 均匀)

3. 问题清单

P0-1 多副本下用户/模型级并发失效

ConcurrencyGuard 用进程内 asyncio.Semaphore,但 K8s 中 worker-main-image / worker-detail-image / worker-hot-replicatereplicas: 2。结果:

  • 用户 concurrent_tasks 限制翻倍(2 个 Pod 各放一份额度);
  • 模型级上游保护失效(每 Pod 各放 max_concurrency,总并发 = N × 上限)。

P0-2 活跃注册表与 stream 的一致性窗口

发布是「先 XADD 后 SADD」,两步非原子。若 registry.register 失败(Redis 抖动),消息进了 stream 却没进活跃集合 → FairConsumer 永远不知道 → 消息永久滞留。现有清理逻辑只删「空 stream」,无法发现这类漂移。

P1-3 PEL 重试无退避

失败保留 PEL 后,下一轮 id=0 立刻回收重试,0 间隔。上游持续失败时几秒内耗尽 max_retries 进死信,无退避、无抖动,可能放大上游压力。

P1-4 死信无落库 / 无监控

Stream 层死信 = ACK + 日志(业务层 llm_image_processor 有 job_store 二次兜底),但:无死信 stream、无 Prometheus 指标(RunOnceStats 只进日志)。

P1-5 空转 Redis 开销

无消息时每 0.1s 对 8 类各做 ~4 次 Redis 往返(ZCARD/SCARD/ZRANDMEMBER/SRANDMEMBER),空载约几百 QPS。pick_users 无 pipeline。

P2-6 视频 typed worker 空转

视频试跑已迁 Temporal(video_job_publisher.py 直接抛错禁止 publish)。grass_video / drama_video 两类 worker(K8s 各 1 副本、1~2Gi 内存)无任何消息源,空占资源。

P2-7 priority 权重未真正生效

zrandmember 均匀随机,不按 score 加权(旗舰=10 vs 专业=5 无差别),与「加权抽取」预期不符。

P2-8 重试计数非原子

increxpire 两步,expire 失败则计数 key 永不过期。

4. 优化清单(已落地/待落地)

优先级 优化 状态
P0-1 ConcurrencyGuard 支持 Redis 分布式信号量(用户级 + 模型级),K8s 副本收敛为 1 待落地
P0-2 发布侧 register 失败重试补偿;周期任务做 stream↔注册表漂移对账 待落地
P1-3 PEL 重试指数退避(min(2^(retry-1), 60)s 待落地
P1-4 死信 XADD 到 stream:llm:dead-letter + Prometheus 指标 待落地
P1-5 活跃集合全空时退避 0.5s;pick_users pipeline 化 待落地
P2-6 删除 grass_video / drama_video 空转 Deployment 与本地启动项 待落地
P2-7 priority 加权轮盘赌采样 待落地
P2-8 重试计数改 Lua 原子 待落地

5. 验证方式

  • 单元:poetry run pytest tests/unit/messaging -q --no-cov
  • 本地:起 Redis + 单 typed worker,压入任务观察消费/退避/死信/指标
  • K8s:make k3s-deploykubectl get pods -n harness-app -l component=backend-worker 应只剩 6 类各 1 副本;curl :9101/metrics 验证指标