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
关键设计:
- 每用户独立 stream(
stream: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._concurrency(MESSAGING_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 |
| 视频上游 | RedisSemaphore(semaphore: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-replicate 为 replicas: 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 重试计数非原子¶
incr 与 expire 两步,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-deploy后kubectl get pods -n harness-app -l component=backend-worker应只剩 6 类各 1 副本;curl :9101/metrics验证指标