Files
lislgosms/docs/phase-4-send-pipeline-redesign.md
T

13 KiB
Raw Blame History

CMPP 发送链路重新设计方案

更新日期:2026-08-25 状态:已在测试环境完成实施与正价50 TPS验收 适用范围:耐久 Inbox 完成业务校验后,到 Gateway 供应商 Submit、结果回调、回执和计费结算的完整链路

1. 结论

现阶段不应继续在 BullMQ 发送 Worker 内做小幅并发、缓存或微批补丁。已完成的耐久 Inbox、入口幂等、路由/报备查询收敛、开放会话热点移除、无消费者 BullMQ 副本移除和计费锁缩短应保留。

若以后继续追求完整供应商提交 100~500 条/秒,建议作为独立项目,重新设计为四个隔离阶段:

  1. Inbox 业务处理与计费预留;
  2. 批量路由规划与提交事实持久化;
  3. PostgreSQL Submit Outbox 批量发布;
  4. 独立的提交结果、回执、计费结算与任务进度回调。

核心不是把若干 BullMQ Job 临时拼成数组,而是把“应提交什么”先作为 PostgreSQL 事实一次性落库,再由无业务判断的发布器向 Gateway 供给命令。这样才能同时减少数据库往返、消除数据库到 Redis 的崩溃窗口,并隔离供应商提交与回调写入的资源竞争。

2. 原方案中应保留的部分

  • 正价测试使用 325 分而不是 0 单价,暴露了真实计费锁竞争;
  • pg_stat_statements 和固定低基数阶段指标能区分入口、Worker、Gateway 与回调瓶颈;
  • 路由、在线连接和签名报备合并查询,取消最终通道重复报备查询;
  • CMPP 单号码内部任务不再逐消息执行整批状态 GROUP BY
  • 提交事务不再更新 CmppSubmitSession.submitTotal 热点行;
  • 删除没有消费者的 gateway.submit.queue BullMQ 同步副本,保留 Gateway 实际消费的 Redis Stream
  • 账户 advisory lock 移到事务末尾账务段后,平均等待由约 443.963ms 降至约 5.420ms
  • 验收以完整供应商提交、队列排空、消息唯一性和账务恒等式为准,没有用入口 SubmitResp 冒充完整吞吐。

这些改动降低了单消息成本和锁等待,但没有改变发送链的逐消息架构,因此未把完整供应商吞吐提升到目标档位。

3. 原方案存在的问题

3.1 微批位置选错

原 P2 在 BullMQ Worker 已领取单个 Job 后,用 5ms 窗口聚合最多 20 条消息。这个位置太晚:任务领取、调度和并发槽已经发生,批次规模还受 Worker 并发上限约束,无法稳定形成 32/64 条批次。

批内慢消息还会形成最慢项栅栏,使本来可以独立完成的消息互相等待。实测微批完整供应商提交约 22.20 条/秒,未优于对照,证明这不是正确的批处理边界。

3.2 只批量加载,没有批量处理主成本

原微批只合并了消息读取,运营商识别、路由、在线/报备判断、SmsSubmitRecord 创建、消息更新、Stream 发布、结果/回执/计费和任务进度仍逐消息执行。省下一个 WHERE id IN (...) 查询不足以抵消批次协调开销。

3.3 PostgreSQL 与 Redis 仍不是原子边界

当前链路先提交 PostgreSQL,再发布 Redis Stream。进程若在两步之间退出,会留下数据库已进入待提交状态、但 Gateway 未收到命令的窗口。若 Redis 已成功而数据库未记录发布结果,重试又可能重复发布。

Redis pipeline 只能减少网络往返,不能解决跨介质一致性。必须增加 PostgreSQL Outbox,并以稳定 submitId、命令幂等键和 Gateway 去重实现至少一次发布。

3.4 提交供给与回调写入争用资源

供应商 Submit 结果、分片结果、最终回执、下游投递和计费结算会同时写 PostgreSQL。回调越快,数据库写入越密集,越容易抢占发送 Worker 的连接和 CPU。仅提高 Worker 并发只会把等待转移到数据库池。

3.5 实时业务事实不能用长期缓存替代

通道启停、真实连接、签名报备、应用/企业状态、余额和号码频控都可能变化。正确做法是批次内共享只读快照,并在提交事实落库前用数据库条件最终校验;不能建立跨批次长期缓存决定是否发送。

3.6 压测口径和夹具容易误导

入口接收、Inbox 完成、供应商首次 Submit 和最终回执是不同阶段。入口 100 条/秒成功,不代表完整供应商提交达到 100 条/秒。历史待投递回执、重复号码频控、缺失运营商规则,以及模拟器重启后数据库连接状态尚未恢复,都会污染结论。

4. 目标架构

4.1 阶段 AInbox 与业务预留

保留现有耐久 Inbox。业务 Worker 批量领取时继续执行应用、企业、签名、模板、黑名单、日限额、号码频控和正价余额预留。

  • 使用 FOR UPDATE SKIP LOCKED 在短事务内领取 32 条,验证稳定后最多 64 条;
  • 领取事务只更新租约,业务查询和外部操作不放在持锁事务内;
  • 日限、频控和冻结继续保留逐消息稳定业务键及数据库唯一约束;
  • 同企业账户锁只覆盖最终余额与流水写入。

4.2 阶段 B:批量路由规划与提交事实

按应用、运营商、签名组合批量读取候选,在一个短事务内为每条可发送消息持久化:

  • 独立 SmsSubmitRecord
  • 目标通道、submitId 和消息 submit_queued 状态;
  • 一条 GatewaySubmitOutbox 命令事实。

事务提交前,SQL 再次约束应用/企业 active、通道 active、连接 connected、签名报备 approved。任一消息失效只拒绝该消息,不回滚整批其他消息。供应商网络调用、Redis 写入和任务进度聚合都不进入该事务。

4.3 阶段 CPostgreSQL Submit Outbox

建议新增表:

GatewaySubmitOutbox
- id / submitId(唯一)
- messageRecordId / channelId
- payload / schemaVersion
- status: pending | publishing | published | dead
- attemptCount / nextAttemptAt
- leaseOwner / leaseExpiresAt
- streamEntryId / publishedAt / lastError
- createdAt / updatedAt

发布器只做:

  1. FOR UPDATE SKIP LOCKED 原子领取一批 pending 行;
  2. 使用 Redis pipeline XADD,每条命令携带唯一 submitId
  3. 批量回写 published 或可重试失败。

“Redis 已成功、PostgreSQL 回写前崩溃”会产生重复发布,因此 Gateway 必须按 submitId 去重,结果事件继续按 resultEventId 幂等。Outbox 扫描器负责恢复数据库已提交但未发布的命令。

4.4 阶段 D:回调与计费隔离

提交结果、回执和下游投递使用独立回调进程/连接池:

  • 单分片只处理聚合结果,多分片保留逐片事实与最终聚合;
  • accepted 结算复用冻结事实,保持 released/charged 幂等且净余额正确;
  • rejected/timeout 以唯一 retryOfSubmitRecordId 认领组内补发;
  • 最终失败退款使用稳定幂等键;
  • 任务进度从消息事实异步刷新,CMPP 单消息任务继续直接写计数;
  • 回调积压不得阻塞 Outbox 向 Gateway 供给命令。

4.5 连接池预算

角色 建议初始数据库槽 说明
API/入口 1624 保证客户 Submit 和查询接口可用
Inbox/路由规划 1624 批量查询和短事务持久化
结果/回执/计费 1216 吸收供应商回调峰值
Outbox 发布器 48 只领取、发布和回写
运维/迁移保留 不少于 8 健康检查、诊断和恢复

实际值必须结合 PostgreSQL CPU、max_connections、PgBouncer 模式和 pg_stat_statements 复测,不能把表中数字直接当生产配置。

5. 实施步骤

第 0 步:固定基线和测试夹具

固定专用号段、三运营商比例、隔离应用和六供应商账号;每次故障注入后同时确认模拟器 6/6 和数据库六通道 connected;建立 PostgreSQL、运行源码和环境/systemd 恢复资产。

第 1 步:Outbox 迁移与影子写入

只新增表、索引和指标;原 Stream 发布仍生效,同时影子写 Outbox 但不发布;对比每个 submitId 的原命令与影子 payload。回退只需关闭影子写,不删除历史表。

第 2 步:Outbox 发布器影子验证

发布到不连接真实 Gateway 的影子 Stream,验证领取、租约恢复、重复发布、死信和批量大小;做发布前/后崩溃注入,证明可恢复且幂等。

第 3 步:单应用灰度切换

一个隔离应用切换到正式 Outbox;保留旧直接发布回退开关,但同一消息任何时刻只能有一个发布路径;先做正价 smoke 和全部发送拦截/补发回归。

第 4 步:批量路由与持久化

初始批量 32,按应用、运营商、签名共享查询;批量插入 Submit/Outbox,逐消息保留唯一键和失败结果;使用 EXPLAIN (ANALYZE, BUFFERS)pg_stat_statements 验证真实 SQL。

第 5 步:回调与连接池隔离

拆分 Submit/Receipt 回调连接池;集合式处理可安全合并的流水、账单和任务进度,同时保留逐消息幂等键和确定账务顺序。

第 6 步:容量验收

smoke → 20 → 30 → 50 → 100 → 200 → 300 → 500 逐档执行。每档完全排空后再升档;只有完整供应商提交达到目标、无丢重、账务一致且数据库稳定才继续。

6. 功能开关与回退

建议提供:

  • SEND_SUBMIT_OUTBOX_SHADOW_ENABLED
  • SEND_SUBMIT_OUTBOX_PUBLISH_ENABLED
  • SEND_ROUTE_BATCH_ENABLED
  • SEND_ROUTE_BATCH_SIZE
  • SEND_CALLBACK_POOL_ENABLED

回退时先停止新领取,等待在途事务完成,关闭正式 Outbox 发布并恢复旧直接发布;按 submitId 对账 Outbox 与 Stream PEL,不能删除 Outbox 行或清空 Redis 来回退。

7. 验收标准

  • 客户 SubmitResp 零丢失,MessageId 与业务号码唯一;
  • Inbox、Submit Outbox、Gateway/结果 Stream 和 BullMQ 最终排空;
  • 每条消息最多一个有效首次 Submit,补发必须关联原提交;
  • 供应商完整提交速率达到该档目标,不只看入口速率;
  • 签名、模板、余额、应用/企业状态、号码频次、报备和通道停用继续实时拦截;
  • 冻结、释放、扣费、退款和 SmsBillingRecord 净额一致;
  • 无持续锁等待、无 idle in transaction、连接数不超过预算;
  • 服务重启和 Outbox 重放不重复提交、不重复计费。

8. 风险与工作量判断

这是跨 Prisma migration、发送状态机、Gateway 命令协议、回调处理、部署配置和压测工具的架构改造,不适合作为当前阶段的继续小修。主要风险是双发布、Outbox 重放重复提交、补发认领冲突、实时通道状态过期和账务顺序错误。

建议以后单独立项,先完成 Outbox 影子对账和崩溃恢复,再进入批量路由。没有完成影子验证前,不应直接替换现有正式发布路径。

9. 2026-08-25 完整实施结果

  • GatewaySubmitOutbox、独立发布器和正式发布路径已完成;发送Worker进一步实现最大32条、3ms聚合窗的有界批次。批内一次加载消息、号段、路由/连接及签名报备事实,按通道执行限速,再以一个短事务createMany写Submit、集合式更新消息并createMany写Outbox。重试及异常分支保留逐消息状态机和幂等键。
  • Submit结果、分片结果、回执、上行、受限协议日志和死信改由仅绑定127.0.0.1:3001cmpp-gateway-callback进程接收,使用独立12槽PostgreSQL连接池;指标绑定127.0.0.1:9468。客户HTTP Webhook仍由主API的独立BullMQ Worker处理,避免慢客户回调占用Gateway事实回写池。连接状态控制仍写主API,供应商事件才写回调进程。
  • 初次发布smoke暴露并修复了回调URL职责过宽问题:连接状态误发到回调端点导致数据库显示connecting、路由失败。修复为控制面/事件面双URL后,六连接均为connected;该失败样本未发生供应商提交或计费,不纳入性能结果。
  • 修复后正价smoke为9/9受理,账单9×325=292520 TPS为199/199、账单6467530 TPS为299/299、账单9717550 TPS为499/499、账单162175。四个成功窗口共1006条,账单1006笔、单价均325、合计326950。
  • 50 TPS档499条非补发首次供应商Submit覆盖9.930秒,即50.25条/秒;相对同日改造前Outbox阶段的33.07条/秒提高约52%。该档客户端P50/P95/P99为35/76/135ms;首提499条唯一,全部尝试及Outbox各542条唯一,补发均关联原提交。
  • 压测结束Inbox、Submit Outbox和两条Redis Stream全部排空,数据库无重复Submit ID、无等待锁、无idle in transaction;回调池max=12,total=1,idle=1,waiting=0。隔离应用单价已恢复0,三条临时运营商规则已删除,真实账务事实保留审计。100/200/300/500未继续执行,当前验收结论限定为50 TPS档通过。