# 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 阶段 A:Inbox 与业务预留 保留现有耐久 Inbox。业务 Worker 批量领取时继续执行应用、企业、签名、模板、黑名单、日限额、号码频控和正价余额预留。 - 使用 `FOR UPDATE SKIP LOCKED` 在短事务内领取 32 条,验证稳定后最多 64 条; - 领取事务只更新租约,业务查询和外部操作不放在持锁事务内; - 日限、频控和冻结继续保留逐消息稳定业务键及数据库唯一约束; - 同企业账户锁只覆盖最终余额与流水写入。 ### 4.2 阶段 B:批量路由规划与提交事实 按应用、运营商、签名组合批量读取候选,在一个短事务内为每条可发送消息持久化: - 独立 `SmsSubmitRecord`; - 目标通道、`submitId` 和消息 `submit_queued` 状态; - 一条 `GatewaySubmitOutbox` 命令事实。 事务提交前,SQL 再次约束应用/企业 active、通道 active、连接 connected、签名报备 approved。任一消息失效只拒绝该消息,不回滚整批其他消息。供应商网络调用、Redis 写入和任务进度聚合都不进入该事务。 ### 4.3 阶段 C:PostgreSQL Submit Outbox 建议新增表: ```text 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/入口 | 16~24 | 保证客户 Submit 和查询接口可用 | | Inbox/路由规划 | 16~24 | 批量查询和短事务持久化 | | 结果/回执/计费 | 12~16 | 吸收供应商回调峰值 | | Outbox 发布器 | 4~8 | 只领取、发布和回写 | | 运维/迁移保留 | 不少于 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:3001`的`cmpp-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=2925`;20 TPS为199/199、账单64675;30 TPS为299/299、账单97175;50 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档通过。