Files
msgexchange-v2/docs/architecture.md
T
windyboy cd5de56937 refactor(flight-schd): 全宽表零子表收敛 16 个集合并定案 PIPELINE_LOCK 行锁 (ACM2-29)
依据 ACM2-29 新定案【全宽表·零子表】,废除「*_TXT 过渡 + 按需升独立子表」
原规划(对 XSD maxOccurs=99 的过度设计),结合报文样例与成都现场地服实际
规律(登机门 1~2、值机柜台 1~3、转盘 1~2、延误单有效、靠撤桥/轮挡各 1 次)
将 16 个明细集合全部收敛为 FLIGHT_SCHD 宽表标量列或紧凑 VARCHAR 字串。

行为变更:
- V1.2.0 迁移:新增集合平铺槽位列(GTDT×2/CKDT×3/CLDT×2/PSDT×2/CHDT×2,
  跨集合同名属性按 B 前缀/CH 前缀消解)、里程碑标量列(DELY_*、ABTM_A/D、
  CHOT_ON/OFF)、异常前缀标量列(FDIV/FRET/FLAB)、紧凑航路字串
  (ROUT_PATH/ERUT_PATH)与无界集合 JSON 字串列(SRVT/VIPF/MAFL_TEXT);
  删除全部 16 个 *_TXT 文本列;不建任何子表、零 CLOB;
- 仓储层:删除子表替换/回查机制,写侧集合键平铺为列(序号属性 0 = 显式
  删除标记,分舱复用条目抹平去重,超界按定案丢弃),读侧由平铺列重建
  16 个集合键,视图与 legacy flightInfo hash 保持同构(KAFKA_SCHD 线格式
  与 FS7 Diff 逐字段比对不受存储形态影响);
- I5 单写者锁修正:PIPELINE_LOCK 单行 SELECT ... FOR UPDATE 取代
  PG advisory lock,PG/Oracle 11g 同构,无 DBMS_LOCK DBA 特权依赖;
- FlightFieldsJson 对契约结构字段解析 JSON tree 输出原生数组/对象,
  杜绝集合被再次编码成字符串的双重转义。

不变量:消息严格 FIFO、单写者互斥、增量=字段级合并(集合键出现=单资源
集合级全量快照替换)、快照=整体替换、按代域化差删均不变。

迁移影响:本地开发库为一次性测试数据,已重置并由 Flyway 全新应用
V1.0.0→V1.2.0(同文件名内容变更,沿用旧库会触发校验和不匹配)。

验证:./gradlew test 全绿(66 个用例,含本地 PG 真实方言集成:集合槽位
替换/清除、里程碑与航路平铺回读、PIPELINE_LOCK NOWAIT 互斥与提交释放)。
2026-09-08 08:48:24 +08:00

12 KiB
Raw Blame History

msgexchange-v2 架构文档

1. 系统定位与范围

msgexchange-v2 是机场 OMMS 的上游报文处理中间件,用于替换旧版 msgexchange-api。 它读取 CIIMS、AODB 等系统写入共享 MySQL 信箱的 XML 报文,按顺序更新航班动态,再将结果提供给下游。

本系统负责收报、解析、状态更新和结果投递,不生成上游业务报文,不替代 CIIMS/AODB,也不提供 AODB 主数据编辑能力。

  • 主要入口:轮询共享 MySQL 的 CMINMSGS
  • 兼容入口POST /cminmsgs/send,供现役兼容、手工工具和对拍使用;写入信箱后返回记录 ID,不是生产收报主路径。
  • 输出Kafka 的 msg / schd 消息、共享 MySQL 的 COUTMSGS 出站信箱,以及查询 HTTP 接口;不直接推送前端。
  • 当前范围(阶段 A:运营航班权威状态落自有 PostgreSQL(表 FLIGHT_SCHDSCHD_GEN,ACM2-28 定案采纳选项 C),Redis 彻底退出动态权威与全部写路径。ES 历史投影及相关清场流程属于暂缓的阶段 B。

本文描述架构约束,不代表所有能力已实现;实现缺口见第 9 节。模块交互、状态机和参数详见 design.md,需求见 user-stories.md。历史报文契约仍以 SIS 接口规范XSD 为兼容依据,其他 legacy 资料仅作参考。

2. 总体架构

CIIMS / AODB 等上游
        │ 写入 XML
        ▼
共享 MySQLCMINMSGS
        │ 轮询未处理记录
        ▼
┌──────────────── msgexchange-v2(单实例)────────────────┐
│ ingress:发现报文 → PostgreSQL 持久化入队                │
│                         │                              │
│ processing:取 FIFO 队头 → 解析 / 去重 → Handler 决策    │
│                         └─ PG 单事务:航班变更 + 终态 + 待发事件│
│                                                        │
│ jobs:在主泵空闲或消息退避窗口内执行维护作业              │
│ delivery:读取 PG 待发事件 → 投递 / 重试                  │
└─────────────────────────┬──────────────────────────────┘
                          ├─ Kafkamsg / schd
                          └─ 共享 MySQLCOUTMSGS

处理结果提交后,再回填 CMINMSGS 的处理标记;失败需补偿。
查询接口读取航班动态,不参与状态写入。

收报、处理和投递各使用一条专用线程,不占用 HTTP 事件循环。只有主泵可以写航班动态及快照版本,维护作业也必须遵守这一规则。

采用 Kotlin + JDK 25、Micronaut 编译期依赖注入和 JDBC 持久化。数据库变更由 Flyway 管理,但只作用于自有 PostgreSQL。具体依赖版本以 build.gradle.kts 为准,不在架构文档重复维护。

3. 模块职责

模块 职责与边界
ingress 轮询信箱、持久化入队、补偿重扫及兼容 HTTP 写入;不解析业务报文。
codec XML 解码,区分非法报文与可修复的解码失败。
processing FIFO 调度、业务身份绑定与去重、Handler 决策、快照处理及状态提交。Handler 只返回决策,不直接访问数据库或 Kafka。
delivery 消费待发事件,负责按目标保序、schd 聚合、投递和失败重试。
jobs 持久化维护作业,由主泵在允许的窗口执行;阶段 B 作业暂不启用。
reference 静态参考数据和上游请求跟踪。
domain / config 领域状态、事件和决策模型,以及运行参数。
infra 仓储、外部适配器、重试、健康检查与日志;通过接口隔离基础设施(Redis 已退出核心写路径)。

4. 主流程

收报与处理

  1. InboxPoller 默认每秒扫描 DATE_PROCESSED IS NULL 的信箱记录,在自有 PG 中建立 PROC_STATE(PENDING)。重复扫描不能重复入队;入队失败留待重扫。
  2. 主泵只处理最小未完成 CMINMSGS_ID。解析报文、绑定业务身份并去重后,调用对应 Handler 生成决策。
  3. 在自有 PG 本地事务中同时保存航班状态变更(FLIGHT_SCHD / SCHD_GEN)、处理结果和 MSG_EVENT 待发事件。
  4. 事务提交后,回填共享信箱的处理标记(外部副作用,补偿保障)。

投递

DispatcherMSG_EVENT 取出待发事件。普通事件按投递目标和 EVENT_ID 保序;某个目标失败时,不能跳过其队头投递后续事件。

schd 是最新状态通知,不逐条发送中间变化:统一由 flushSchdFLID 聚合,取批次内最新事件后发送。它不提供逐条变更历史,不能与普通事件的 FIFO 语义混为一谈。

5. 必须保持的约束

  • 消息严格 FIFO:队头失败并退避时,后续消息仍不能越过它。只有队头完成或按失败策略进入终态后,队列才继续推进。收报重扫和水位设计必须防止较小 ID 漏入队而被后续消息越过。
  • 动态状态单写者FLIGHT_SCHD 运营航班表和快照 SCHD_GEN 只由主泵单线程写入。事务通过 PIPELINE_LOCKFLIGHT_SCHD_WRITER 行执行 SELECT ... FOR UPDATE 互斥;该方案在 PG/Oracle 11g 均无需数据库扩展或额外 DBA 特权。不能通过增加实例或处理线程直接扩容。
  • 身份去重:同一业务身份只能绑定一条有效处理记录,重复报文不应再次产生业务副作用。具体身份组成和重放规则见设计文档。
  • 快照可恢复:快照覆盖、旧数据清理和版本推进需要原子性与重放保护;不能在恢复时把旧代数据重新写回。
  • 作业不与消息混排PUMP_JOB 是独立队列,只在没有消息队头或队头处于退避窗口时执行。执行窗口与饥饿边界需要明确验证。

这些约束优先于吞吐量优化。单写者降低了并发复杂度,代价是队头阻塞和吞吐上限;如需并行化,必须先重新定义顺序与状态归属,不能只调整线程数。

6. 数据归属与一致性

存储 承载内容 职责说明
自有 PostgreSQL 处理状态 PROC_STATE、待发事件 MSG_EVENT、维护作业 PUMP_JOB、请求跟踪 REQ_TRACK、静态主数据 REF_MASTER、运营航班表 FLIGHT_SCHD 与日代 SCHD_GEN 本系统唯一业务数据库。消息处理、快照推进与待发事件在单事务内原子提交;本地事务只在此库。
共享 MySQL CMINMSGS 入站信箱、COUTMSGS 出站信箱 外部系统所有。仅执行约定的信箱读写和处理标记回填,不建表、不迁移 schema、不写历史表。兼容 HTTP 入口可按既有契约写入入站信箱。
不使用跨库事务。 PG 事务只能保证“处理结果与待发事件一起提交”,不能覆盖 Redis 更新、MySQL 回填或 Kafka 发送。跨存储依靠幂等、重试和持久化补偿恢复:
中断位置 恢复要求
信箱已有报文,PG 入队失败 重扫补建,并按信箱 ID 去重。
PG 提交失败 事务原子回滚,无中间态残留;消息重试时整体重放。
PG 已提交,信箱回填失败 持久化记录补偿任务并重试回填,不能重新执行已完成的业务处理。
下游已接收,本地尚未标记发送成功 允许重发;下游或出站适配协议必须具备去重能力。

对外投递按至少一次设计,不承诺端到端恰好一次。Kafka 生产者幂等不能消除应用重启或 outbox 重发带来的所有重复。

7. 关键决策索引

保留 D1–D12 编号,便于设计文档和工程历史引用;以下是决策摘要,而非完成清单。

编号 决策及理由
D1 业务报文严格 FIFO,维护作业窗口执行,优先保护航班状态的时序正确性。
D2 阶段 B 暂缓。历史写入成功后才可生成删除事件;顺序调用本身不保证原子性,恢复与去重方案需在启用前补齐。
D3 schd 只从 flushSchd 聚合发送,减少已被覆盖的中间状态通知。
D4 未实现的报文类型按 UNSUPPORTED 可恢复失败处理,不当作非法报文直接丢弃;补齐能力后按重放规则恢复。
D5 在持有具体消息或批次上下文的位置记录失败和退避;不吞掉线程中断或 JVM 严重错误。
D6 基础设施通过接口注入,时间通过 Clock 注入,便于确定性测试顺序、重试与超时。
D7 内存 stub 仅显式开启时装配,生产禁止使用,避免把未持久化的数据误当作已落库。
D8 使用编译期依赖注入,并以启动冒烟测试验证关键 Bean 装配。
D9 自有 PG 内完成本地事务,共享 MySQL 仅作信箱;跨存储采用补偿,不使用 XA。
D10 动态状态单写者,生产只允许一个活动实例;多实例必须先具备可靠的排他保护。
D11 Kafka 生产要求 acks=allenable.idempotence=truemax.in.flight=1,切流前验证 Broker 兼容性;不允许通过关闭幂等来满足生产接入。
D12 仅将自有库终态记录归档到 PROC_STATE_HST,不侵入共享库的表结构或保留策略。

8. 部署、切换与运维

部署与安全

  • 生产维持单活动实例,停机时停止接收新任务并等待工作线程退出。运行期排他保护尚未完成,当前不能依靠程序自动阻止双实例写入。
  • 配置、口令和环境端点通过环境变量提供。兼容写接口沿用内网信任模式,缺少鉴权,必须限制网络访问;管理端点不得直接暴露到生产外网。
  • Eureka 用于服务发现,Logstash 接收结构化日志;日志出口故障不应阻塞业务处理。

替换旧系统

采用“影子对拍 → 切流 → 旧系统冻结”。共享信箱不能让新旧系统同时认领和回填;影子输入使用只读水位或回放。影子环境须隔离 PG schema/实例、Kafka topic 和服务注册身份,并禁止误写生产信箱。切流时保证只有一个权威写者。影子期采用 FlightStoreDiffTool 跨存储对拍。

可观测性要求

使用消息 ID、事件 ID 关联处理与投递日志;健康检查反映依赖实际可用性,而不只是进程存活。运行中重点关注队列积压、队头滞留时间、投递延迟、重试/DEAD 数量和回填补偿积压。死信和一致性异常需要可执行的告警与重放流程,不能只留一条错误日志。

9. 当前实现与上线门槛

当前已有管道骨架、重试机制、部分 JDBC 适配和开发环境 stub 冒烟能力,不能据此认定生产链路已闭环。默认配置关闭管道自动启动及真实数据库/信箱适配。

上线前至少需要完成并验证:

  • 真实 PG 事务、信箱水位与补扫、回填补偿、出站信箱,以及所需业务 Handler。
  • FLIGHT_SCHD 事务原子性、快照 SQL CAS 版本校验、故障中断回滚恢复与影子对拍 cross-store diff 工具。
  • FIFO、身份去重、作业窗口和投递故障下的回归测试。
  • 生产启动校验、单实例排他保护、影子隔离和 Kafka 配置约束;当前配置仍允许 Kafka 参数覆盖,且默认 in-flight 值与 D11 要求不同。
  • 死信告警、人工重放、端到端追踪、积压指标及安全边界。

具体进度由配套设计(design.md §9 已知缺口表)与 Plane ACM2-10 实施计划维护,本文不记录测试数量、临时补丁版本或逐项工单进展。架构基线沿用 ACM2-3,存储与事务边界以 ACM2-12 的修订为准。