现场环境仅提供 Oracle 11g(无任何 JSON 能力),废除 FLTR_JSON JSONB 整文档存储: - V1.1.0 迁移重写:FLIGHT_SCHD 改为一行一航班宽表(SCHD.FLTR 标量字段列 + 1:N 明细集合序列化文本列),SCHD_GEN.FLIDS_JSON 展开为 SCHD_GEN_FLID(FDAY, FLID); 已应用过旧版迁移的环境需重建 schema 重放 - 契约层:FlightChange/Handler 在线视图改为字段集映射(与 legacy flightInfo hash 同构),仓储白名单拒绝未知字段;增量写=字段级合并(hmset 同语义), 快照=整体替换 - 事件载荷:快照投递与投影重建经 FlightFieldsJson 确定性序列化(键序稳定) - 影子对拍:FlightStoreDiffTool 改为 PG 宽表列值 vs legacy FLTR JSON 逐字段 归一化比对,新增嵌套集合比对用例 - 测试:61 用例全绿(不变量门槛 1/2/3、UTC 方言、真实 PG 方言集成、清场五场景) - 文档:decision-flight-state.md 追加 §8 修订记录(不改写定案历史),design.md 表说明同步 11g 方言移植(ON CONFLICT→MERGE、advisory lock 替代、Flyway/驱动矩阵)、 集合列升子表、类型化列提升等未尽事项另立 Plane issue 跟踪。
16 KiB
msgexchange-v2 设计文档
1. 阅读说明
本文说明模块如何协作、状态如何流转,以及失败后如何恢复。系统范围、存储归属和部署约束见 architecture.md,不在这里重复。
阶段 A 采纳 ACM2-28 选项 C 定案:运营航班权威状态落自有 PostgreSQL(表 FLIGHT_SCHD 与 SCHD_GEN),Redis 彻底退出动态权威与全部写路径。阶段 B 的历史投影和清场暂不启用。
2. 数据与领域模型
2.1 持久化记录
所有内部表都属于自有 PostgreSQL;共享 MySQL 只保留约定的信箱读写边界。
| 记录 | 用途 | 关键约束 |
|---|---|---|
PROC_STATE |
入站消息的处理状态、身份、重试次数和错误原因 | CMINMSGS_ID 主键防止重复入队;IDENTITY_KEY 唯一约束防止业务重复;按最小未完成消息 ID 取队头。 |
MSG_EVENT |
等待投递的事件(outbox) | EVENT_ID 决定投递顺序;TARGET 区分目标;PARTITION_KEY 在 schd 中为 FLID。 |
PUMP_JOB |
持久化维护作业 | 状态为 QUEUED / RUNNING / DONE / FAILED;不与业务消息共用排序序号。 |
REQ_TRACK |
上游请求及应答关联 | 保存请求类型、参数、出站信箱 ID、发送和完成时间;同类只允许一个开放请求。 |
REF_MASTER |
静态参考数据 | (RTYPE, RKEY) 唯一,SOURCE 记录数据来源。 |
FLIGHT_SCHD |
运营航班当前权威状态(SCHD 快照 + FLOP/ADFT 增量合并),一行一航班宽表 | FLID 主键,FDAY 所属日代(可空),SCHD.FLTR 标量字段列 + 1:N 明细集合序列化文本列(11g 修订,废除 FLTR_JSON JSONB 整文档),TIMESTAMPTZ 时间戳。 |
SCHD_GEN |
日计划代版本(差删依据与版本 CAS) | FDAY 主键,VERSION 版本号,TIMESTAMPTZ。 |
SCHD_GEN_FLID |
当前代有效航班 FLID 集合(原 FLIDS_JSON 展开为行) |
(FDAY, FLID) 主键,FDAY 外键级联删除。 |
PROC_STATE_HST |
终态处理记录的归档目标 | 属于目标设计,当前迁移尚未建表;不得改写为共享库历史表。 |
字段与索引定义以 src/main/resources/db/migration/ 为准(含 V1.1.0__flight_schd.sql)。报文原文仍从共享信箱读取,因此必须协调原文保留期,不能在消息尚需处理或重放时提前清理。
Redis 已彻底退出动态权威与写路径;FLIGHT_SCHD 与 SCHD_GEN 在自有 PG 中由主泵单线程独占写入。
2.2 消息、身份与决策
XmlCodec 将 XML 解码为 DecodedMessage,包含 SNDR / TYPE / STYP / SEQN / DTTM 元数据和业务载荷。MsgKind 区分 SCHD 与 FLOP 子类型;Handler 查找与日志类型标识使用同一套映射。
业务身份统一由 Identity.of 生成:
SNDR | TYPE | STYP | SEQN
接收时只按信箱 ID 去重;解码后才首次绑定业务身份。重试保留原有绑定,不能把自己判为重复消息。身份被另一条记录占用时,当前消息转为 SKIPPED,记录 duplicate-of:<id>。是否加入日期边界取决于上游序号重置规则,默认关闭;上线后不能随意更换身份算法。
Handler 是纯函数:
Handler.decide(flightView, message) → Decision
Decision = 航班变更 + msg 通知 + schd 状态 + 出站意图 + 静态数据变更
Handler 不写 Kafka 或数据库。主泵负责应用决策;各类变更在自有 PG 单事务内提交。
2.3 状态与错误分类
处理:PENDING / FAILED → SUCCEEDED(成功)
→ SKIPPED(忽略或重复)
→ FAILED(等待重试)
→ DEAD(非法报文或重试耗尽)
投递:PENDING → SENT
→ PENDING(退避后重试)
→ DEAD(重试耗尽)
SUCCEEDED / SKIPPED / DEAD 是处理终态,不再阻塞后续消息;FAILED 不是终态,仍占据队头。DEAD 表示需要处置,不等于业务成功。
| 错误类别 | 处理方式 |
|---|---|
MALFORMED |
报文非法,直接 DEAD,不在原记录重放白名单内。 |
CODEC_ERROR |
解码能力问题,退避重试;修复后允许重放。 |
UNSUPPORTED |
Handler 或快照能力未实现,按可恢复失败处理,不直接当作非法报文;仍受重试上限约束。 |
INFRA |
基础设施或执行异常,退避重试。 |
EXHAUSTED |
重试耗尽或滞留超时,转 DEAD,人工复核后允许重放。 |
3. 收报与主泵
3.1 收报
InboxPoller 默认每秒读取未处理信箱记录,在 PG 建立 PENDING,不解析业务载荷。PG 插入必须按信箱 ID 幂等,失败由后续扫描补建。
水位优化分为两条路径:快路径读取水位之后的新记录,补偿路径重扫遗漏的未处理记录。只有本批 PG 入队全部确认后才能推进水位。水位不是已处理标记,也不能单独证明较小 ID 已收齐;迟提交和补扫场景的顺序保证需要在启用前验证。
兼容 HTTP 入口执行“写入共享信箱 → PG 入队”。两步不在同一事务中:信箱成功而 PG 失败时,原文不能丢失,由轮询补建;客户端失败重试可能再次写信箱,业务身份去重仍然必需。
3.2 主泵调度
每次 Pump.tick:
- 读取最小未完成消息,必须包含
FAILED,不能只查当前可执行的记录。 - 无消息,或队头仍在退避窗口内时,允许执行一个维护作业;消息已可执行时优先处理消息。
- 队头达到重试或滞留上限时转
DEAD(EXHAUSTED);未到重试时间则等待,不领取后续消息。 - 其余情况调用
MessageProcessor.processOne。
作业执行时长和饥饿边界需要限制,不能用长期作业阻塞已到期消息。滞留超时应基于稳定的起始时刻,不能用每次失败都会刷新的 updatedAt 代替。
3.3 单条处理
读取原文 → 解码 → 忽略规则 → 首次绑定身份
├─ RESP / DNLD:快照流程
└─ 其他:Handler 决策
↓
PG 单事务:FLIGHT_SCHD 增量更新 + 待发事件 + 处理终态
↓
提交后补偿回填信箱
- 原文缺失当前归为
MALFORMED;读取异常不能伪装成“缺失”,应进入基础设施重试。 - 忽略规则覆盖约定的
LDM / REGN / RSTA / EROR,转SKIPPED并审计;不能产生业务副作用。 - 航班变更(
FLIGHT_SCHD)与待发事件、处理终态必须在同一 PG 事务原子提交。跨存储双写窗口已根除。 - 所有终态都需要回填信箱,包括成功、忽略、重复和死信;非终态禁止回填。回填必须在 PG 提交后执行,并有持久化补偿、退避与告警;影子环境禁写。
4. 日计划快照与请求匹配
4.1 快照发布
SCHD-RESP 和 SCHD-DNLD 都进入 SnapshotFlow:
- 暂存校验:流式解析后完成整包校验和航班规范化;失败前不修改权威状态。暂存数据可在崩溃后从原文重建。
- 应答守卫:
RESP必须匹配开放的RQFD请求;无匹配、已过期或报文时间早于发送时间时,转SKIPPED并审计,不更新快照。 - SQL 原子替换:在自有 PG 单事务内批处理写入新代
FLIGHT_SCHD全量、按FDAY域内差删旧代航班,并通过 SQL CAS 推进SCHD_GEN版本。删除集为“旧代航班集合 − 新代航班集合”,仅作用于属于该代的记录,ADFT(FDAY=NULL)与已迁移航班天然存活。 - 提交结果:在同一 PG 事务中保存
MSG_EVENT、将消息置为SUCCEEDED,并将匹配RESP的请求置为DONE;事务提交后执行信箱回填。
单事务保证崩溃后整体回滚,无跨存储中间态;重放时识别版本,绝不发生二次自增。
4.2 上游请求与静态数据
RequestCoordinator 管理请求生命周期:
REGISTERED → SENT → WAITING → DONE
└──→ EXPIRED
注册同类新请求前使旧开放请求过期。只有 COUTMSGS 写入确认后才标记 SENT 并关联出站记录;写信箱成功但本地未确认的情况需要补偿与去重,不能无条件重新发送。
应答优先按已确认的回显字段精确匹配。回显契约未确认时,按同类开放请求和 DTTM ≥ sentAt 判断的降级方式存在跨代误配风险,必须明确接受并审计,不能宣称精确关联。比较前统一时区和时间单位。
静态应答写入自有 PG 的 REF_MASTER,日计划应答走快照流程。请求完成必须在相应数据处理成功之后;超时和迟到应答不能修改已关闭请求对应的状态。
5. 事件投递
5.1 普通事件
Dispatcher 按 TARGET 读取最小未发送 EVENT_ID。队头退避未到期时,该目标停止推进;发送确认后才标记 SENT,失败记录次数和下次执行时间。所有外部调用需要有界超时,避免阻塞整个投递线程。
投递是至少一次:下游已接收但本地未标记成功时可能重发。Kafka 生产约束沿用架构决策 D11,但生产者幂等不替代应用层事件去重;跨重启的事件身份和下游去重契约仍需落实。共享出站信箱也必须单独解决重复写入,不能假设 Kafka 的保证适用于 MySQL。
5.2 schd 聚合
KAFKA_SCHD 不进入逐条投递循环,只由 flushSchd 发送:
- 队头可执行后,按
EVENT_ID顺序领取有界批次。 - 按
FLID分组,保留批次内最大EVENT_ID对应的状态,组成 FLTR JSON 数组发送。 - 成功后将本批被代表的事件一起标记完成,推进
lastFlush;失败则整批增加次数并退避,达到上限整批转DEAD。
默认聚合周期 3 秒、批上限 500。它提供最新状态通知,不保留每次中间变化;批次不能绕过尚在退避的队头。
6. 失败恢复与维护作业
6.1 失败、重试与重放
ProcFailure 和 FailureScheduler 统一处理侧的失败落账;投递侧按单条或聚合批次执行同样的次数与退避规则。默认最多 5 次,退避档位为 1、2、4、8、16 秒,单档封顶 60 秒。
失败必须在持有具体消息或批次的位置记录,外层循环只做兜底日志和等待,不重复增加次数。线程中断应恢复中断标记并向上传递;不捕获 JVM Error 作为普通业务失败。所有时间判断通过注入的 Clock 完成。
ReplayService 只允许 CODEC_ERROR / UNSUPPORTED / INFRA / EXHAUSTED 从 FAILED / DEAD 回到 PENDING,重置次数和下次执行时间,保留身份与错误审计。旧消息进入终态后,后续消息可能已经执行;因此重新入队不等于恢复历史顺序,人工重放前必须评估状态覆盖和版本保护,不能直接批量重放到生产。
6.2 归档与阶段 B 作业
ARCHIVE 仅归档自有库的终态记录,默认按接收时间保留 1 天,保留期可配置为 1~7 天。归档表、接收时间依据和关联事件处理尚需落地;不得归档未完成记录,也不能因移走身份记录而意外失去业务去重能力。共享信箱保留策略由库所有方管理。
HISTORY_SWEEP 与 PROJECTION_REBUILD 暂缓。未来清场必须只删除已确认成功写入历史存储的集合;历史写入未接通时默认删除零条。历史写入与删除事件入队之间仍需恢复方案,顺序调用不构成原子提交。
7. 接口与运行配置
兼容入口为 POST /cminmsgs/send,请求体为原始 XML,当前成功响应为 HTTP 200、text/plain 格式的信箱记录 ID。这只表示接收结果,不表示业务处理成功。请求媒体类型、字符集和失败响应仍需与现役逐项对拍;查询与其他兼容端点不能因列入需求就视为已提供。
运行配置以 application.yml、application-dev.yml 和 .env.example 为准,设计上重点区分:
pipeline.autostart与msgx.stubs:分别控制管道启动和内存适配器;生产禁止 stub,默认不自动启动。pipeline.poll-interval / max-attempts / head-deadline:控制轮询、重试上限和队头滞留;不能改变 FIFO。schd.flush-period / flush-limit:控制状态通知的聚合延迟与批量大小。identity.include-day-boundary:影响去重语义,不能作为普通调优项切换。phase:阶段 A 是当前范围,不应把切为 B 当成已具备历史投影能力。
日志关联消息 ID、事件 ID 和批次;失败记录错误分类、次数、下次执行时间。健康检查反映依赖实际可用性;队头滞留、积压、死信和补偿失败需要指标及告警。日志出口故障不得阻塞业务线程。
8. 验证要求
单元测试使用内存仓储和可推进的 Clock,不依赖睡眠或在线中间件。以下不变量必须有回归测试,接口级单测不能替代主泵调度测试:
| 场景 | 必须验证的结果 |
|---|---|
| 重复扫描、入队中断、较小 ID 迟到 | 不重复入队、不丢记录、不让后续消息越序。 |
| 队头失败、退避及作业竞争 | 消息不越队;到期后恢复;作业不使消息无限饥饿。 |
| 同身份多条记录、失败后重试、归档后重复 | 只产生一次有效业务处理,不把自身重试判为重复。 |
| PG 事务失败、快照重复或迟到 | 整体回滚重试、不重复推进版本、不回退状态、不误删增量航班。 |
| PG 提交失败、信箱回填失败 | 事件与处理结果一起回滚;已提交结果只补偿回填。 |
| 投递确认丢失、批次失败、次数耗尽 | 允许可识别的重发、保持目标顺序、整批退避并保留死信。 |
| 请求超时、无匹配 RESP、时间单位不一致 | 不误用迟到应答,不提前完成请求。 |
| stub 误配置、重复实例、停机中断 | 生产拒绝不安全启动,工作线程能正确退出。 |
真实适配器还需 PG 事务与补偿集成测试、快照 SQL CAS 崩溃恢复测试、Kafka 故障投递验证;常规验证命令为 JDK 25 下执行 ./gradlew test。
9. 实现入口
生产代码根目录为 src/main/kotlin/com/gzzn/omms/msgexchange/:
| 关注点 | 主要入口 |
|---|---|
| 收报与兼容接口 | ingress/InboxPoller.kt、InboxService.kt、InboxController.kt |
| 调度与处理 | processing/Pump.kt(含 MessageProcessor)、Handler.kt、Identity.kt |
| 快照与请求 | processing/SnapshotFlow.kt、reference/RequestCoordinator.kt |
| 投递与作业 | delivery/Dispatcher.kt、SchdAggregation.kt、jobs/JobExecutor.kt |
| 持久化与恢复 | infra/persistence/、infra/redis/、infra/retry/ |
| 启停与配置 | PipelineLifecycle.kt、config/PipelineProps.kt |
10. 当前实现差异
以下缺口直接影响上述设计是否成立,不能以类或接口已存在作为完成依据:
- 事务与外部副作用:ACM2-28 定案后,已完成
FLIGHT_SCHD增量更新与快照全量写入、待发事件与PROC_STATE在自有 PG 单事务原子提交,提交后异步回填信箱。 - 收报与调度:轮询仍从
afterId=0扫描,持久水位与补扫策略未完成;作业以“队头为 FAILED”近似窗口,未区分是否已到期。主泵直取系统时间,滞留判据使用更新时刻,不能保证设计要求的超时升级。 - 快照与业务能力:快照 SQL 事务化已落地,支持 JDBC 批处理、域内差删与 SQL CAS 推进,超限 10000 熔断保护;忽略规则及所需 Handler 仍需在后续阶段铺开。
- 请求、静态数据与归档:请求在真实出站前就标记发送,时间匹配与审计仍需修正;静态数据存在进程内过渡实现,归档表与关联保留策略尚未落地。
- 生产与运维:缺少完整启动校验、运行期单写者保护、影子隔离和 Kafka 强制配置校验;投递滞留升级、死信告警、端到端追踪与指标未闭环。影子对拍 cross-store diff 工具已就绪。
进度与验收项见 Plane ACM2-10 实施计划(U01–U30),业务契约与待确认事项见 user-stories.md。本文件不维护工单流水账、测试数量或历史方案全文。