- F1/F6 HTTP 契约:C-28 目标 ResponseDto + 请求体上限暂定 10MB,Q3 指向 C-28 - F2 legacy 与 SIS 的优先级限定;F3 超期判据对齐 ENQUEUED_AT;F4 改引 D1 - F5 参数默认值改引 reference/PARAM;F7 补 4 条契约状态词;F8 定义 MAID/MAFL 并明确排除共享航班;F9 补 INV-20 - N1/N2 错配编号改引;N3 迁移链指针;N4/M5 白名单补 D-x/OPS-x;N6 补 INV-16/17 映射;N7 加粗规则 - M1 回填退避事实与论证;M2 Q6/Q7;M4 Q4 假定标注;M7 去复述 纯文档,不改代码/迁移/配置/legacy。
139 lines
12 KiB
Markdown
139 lines
12 KiB
Markdown
# 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_SCHD` + 资源明细表 + `FLIGHT_ROUTE_POINT`,权威口径见 [flight-state.md](flight-state.md))。无 Redis 依赖;ES 历史投影属暂缓的阶段 B。
|
||
|
||
本文描述架构约束,不代表所有能力已实现;实现缺口见 [invariants.md](invariants.md) 的声明边界与 Plane(ACM2)。模块交互、状态机与参数详见 [design.md](design.md),前提与不变量见 [invariants.md](invariants.md),对外契约见 [contracts.md](contracts.md),需求见 [user-stories.md](user-stories.md),参数与指标见 [reference.md](reference.md);操作步骤在上线/切流前另立规程(设计阶段只保留前置条件与红线)。历史报文契约仍以 [SIS 接口规范](legacy/SIS_AODB_RMS-V0.1.md) 和 [XSD](legacy/unisysaodbsis.xsd) 为兼容依据,其他 legacy 资料仅作参考,旧系统行为基线的职责见 [README](README.md) 职责表。
|
||
|
||
现场供库时目标为 Oracle 11g,否则自建 PostgreSQL;当前只有 PG 实现可运行,Oracle 不是已支持的平台。
|
||
航班表结构和处理逻辑见 [运营航班状态设计](flight-state.md)。
|
||
|
||
## 2. 总体架构
|
||
|
||
```text
|
||
CIIMS / AODB 等上游
|
||
│ 写入 XML
|
||
▼
|
||
共享 MySQL:CMINMSGS
|
||
│ 轮询未处理记录
|
||
▼
|
||
┌──────────────── msgexchange-v2(单实例)────────────────┐
|
||
│ ingress:发现报文 → PostgreSQL 持久化入队 │
|
||
│ │ │
|
||
│ processing:取 FIFO 队头 → 解析 / 去重 → 处理器决策 │
|
||
│ └─ PG 单事务:航班变更 + 终态 + 待发事件│
|
||
│ │
|
||
│ jobs:独立维护线程(回填补偿 / 历史归档 / 留痕清理) │
|
||
│ delivery:读取 PG 待发事件 → 投递 / 重试 │
|
||
└─────────────────────────┬──────────────────────────────┘
|
||
├─ Kafka:msg / schd
|
||
└─ 共享 MySQL:COUTMSGS
|
||
|
||
处理结果提交后,再回填 CMINMSGS 的处理标记;失败需补偿。
|
||
查询接口读取航班动态,不参与状态写入。
|
||
```
|
||
|
||
收报、处理、投递与维护作业各使用独立线程,不占用 HTTP 事件循环。**航班当前态的写入只发生在持有 `PIPELINE_LOCK` 的事务内**,由主泵串行驱动。
|
||
|
||
采用 Kotlin + JDK 25、Micronaut 编译期依赖注入和 JDBC 持久化。数据库变更由 Flyway 管理,但只作用于自有 PostgreSQL。具体依赖版本以 `build.gradle.kts` 为准,不在架构文档重复维护。
|
||
|
||
## 3. 模块职责
|
||
|
||
| 模块 | 职责与边界 |
|
||
|---|---|
|
||
| `ingress` | 轮询信箱、持久化入队、补偿重扫及兼容 HTTP 写入;不解析业务报文。 |
|
||
| `codec` | XML 解码,区分非法报文与可修复的解码失败。 |
|
||
| `processing` | FIFO 调度、业务身份绑定与去重、领域决策与落库(SCHD/FLOP/FDEL/ADFT):纯领域逻辑只返回决策;Processor 作为事务协调器,在锁事务内完成状态写入、事件与回填意图登记,不直接触碰 Kafka。 |
|
||
| `delivery` | 消费待发事件,负责按目标保序、`schd` 聚合、投递和失败重试。 |
|
||
| `jobs` | 回填补偿扫描、航班历史清理与留痕保留期清理;独立 job 线程执行(调度见 design.md「维护作业与归档」,红线见 flight-state.md「生命周期与开放项」,与主泵的互斥见 `INV-18`),不参与 FIFO。 |
|
||
| `domain` / `config` | 领域状态、事件和决策模型,以及运行参数。 |
|
||
| `infra` | 仓储(JDBC/stub)、外部适配器、重试、健康检查与日志;通过接口隔离基础设施。 |
|
||
|
||
## 4. 主流程
|
||
|
||
### 收报与处理
|
||
|
||
1. `InboxPoller` 按 `PARAM:msgx.pipeline.poll-interval` 周期按 ID 区间扫描水位 `W` 之后的信箱记录(`ID > W`,**不以处理标记为谓词**),在自有 PG 中建立 `PROC_STATE(PENDING)`;水位与入队在同一事务推进(`INV-2`)。重复扫描不能重复入队,中断后由重扫补建。扫描谓词与水位见 design.md「收报与水位」。
|
||
2. 主泵只处理最小未完成 `MSG_ID`。解析报文、绑定业务身份并去重后,分派给 SCHD/FLOP/FDEL/ADFT 处理器。
|
||
3. 在自有 PG 同一事务内(先取 `PIPELINE_LOCK`)保存航班状态变更(`FLIGHT_SCHD` 与明细表)、`MSG_EVENT` 待发事件、处理终态与回填意图。
|
||
4. 事务提交后,补写共享信箱的处理标记(外部副作用,由回填退避重试与超期强制补写保障)。
|
||
|
||
### 投递
|
||
|
||
`Dispatcher` 从 `MSG_EVENT` 取出待发事件。普通事件按投递目标和 `EVENT_ID` 保序;某个目标失败时,不能跳过其队头投递后续事件。
|
||
|
||
`schd` 是最新状态通知,不逐条发送中间变化:统一由 `flushSchd` 按 `FLID` 聚合,取批次内最新事件后发送。它不提供逐条变更历史,不能与普通事件的 FIFO 语义混为一谈。
|
||
|
||
## 5. 必须保持的约束
|
||
|
||
本节只列约束的**归属**;完整定义与验证映射见 [invariants.md](invariants.md),实现与演进不得违反:
|
||
|
||
- 消息严格 FIFO:`INV-3`、`INV-4`、`INV-5`(发现完整性依赖 `PRE-2`/`PRE-3`,当前不可对外声明,见 CLM-1/CLM-2)。
|
||
- 动态状态单写者与写者集合互斥:`D2`、`INV-18`。
|
||
- 身份去重:`INV-9`(身份组成见 design.md「消息、身份与决策」)。
|
||
- 快照可恢复与运营日不可变:`INV-12`、`INV-13`。
|
||
- 航班当前态的物理清除只发生在历史归档之后:`D1`(红线见 flight-state.md「生命周期与开放项」)。
|
||
|
||
这些约束优先于吞吐量优化。单写者降低了并发复杂度,代价是队头阻塞和吞吐上限;如需并行化,必须先重新定义顺序与状态归属,不能只调整线程数。
|
||
|
||
## 6. 数据归属与一致性
|
||
|
||
| 存储 | 承载内容 | 职责说明 |
|
||
|---|---|---|
|
||
| 自有 PostgreSQL | 单行锁 `PIPELINE_LOCK`、处理状态与回填事实 `PROC_STATE`、消费水位 `INBOX_CURSOR`、待发事件 `MSG_EVENT`、请求跟踪 `REQ_TRACK`、航班当前态 `FLIGHT_SCHD` + 8 张资源明细表 + `FLIGHT_ROUTE_POINT`、留痕 `SCHD_SNAP_LOG` | 本系统唯一业务数据库。消息处理、状态推进、处理终态、回填意图与待发事件在单事务内原子提交;本地事务只在此库。 |
|
||
| 共享 MySQL | `CMINMSGS` 入站信箱、`COUTMSGS` 出站信箱 | 外部系统所有。本系统仅执行约定的信箱读写与处理标记回填,不建表、不迁移 schema、不写历史表;由库方按 `Q9` 执行的清除与历史归档见 contracts.md「保留与清除」。兼容 HTTP 入口可按既有契约写入入站信箱。 |
|
||
|
||
**不使用跨库事务。** PG 事务只能保证“处理结果与待发事件一起提交”(`INV-17`),不能覆盖 MySQL 回填或 Kafka 发送等外部副作用。跨存储依靠幂等、重试和持久化补偿恢复;各中断位置的判定与恢复动作见 design.md「中断恢复」。
|
||
|
||
对外投递按**至少一次**设计,不承诺端到端恰好一次。Kafka 生产者幂等不能消除应用重启或 outbox 重发带来的所有重复。
|
||
|
||
## 7. 关键决策
|
||
|
||
仅保留仍具约束价值、且无法从正文(§4–§6、design.md)直接推出的决策,按 D1–D4 连续编号供正文与 design.md 引用;其余曾编号条目(严格 FIFO、stub 门控、本地事务、UNSUPPORTED 处理等)已在正文以约束形式表达,不再重复列表。状态只反映是否已落地,不代表决策被撤销。
|
||
|
||
| 编号 | 决策及理由 | 当前状态 |
|
||
|---|---|---|
|
||
| D1 | 航班清场只在历史写入成功后进行,未接通时删 0 条;未经 FDEL 的清场须先补发删除事件。ES 历史投影(阶段 B)暂缓。 | 红线已实现于 `HistorySweepJob`;恢复/去重方案未闭合 |
|
||
| D2 | 动态状态单写者,生产只允许一个活动实例;多实例必须先具备可靠的排他保护。 | 事务行锁已实现;实例级排他未完成 |
|
||
| D3 | Kafka 生产要求 `acks=all`、`enable.idempotence=true`、`max.in.flight=1`;不允许通过关闭幂等来满足生产接入。 | 约束未强制:默认值与 D3 不一致(见 reference 参数表),且可用环境变量覆盖 |
|
||
| D4 | 自有库终态记录只归档到 `PROC_STATE_HST`,不侵入共享库的表结构或保留策略。 | 目标表未建,尚无归档作业 |
|
||
|
||
## 8. 部署、切换与运维
|
||
|
||
**部署与安全**
|
||
|
||
- 生产维持单活动实例,停机时停止接收新任务并等待工作线程退出。已有事务级行锁,但消息认领和整个实例的排他保护尚未完成,不能依靠行锁宣称支持双实例 FIFO。
|
||
- 配置、口令和环境端点通过环境变量提供。兼容写接口沿用内网信任模式,缺少鉴权,必须限制网络访问;管理端点不得直接暴露到生产外网。
|
||
- Eureka 用于服务发现,Logstash 接收结构化日志;日志出口故障不应阻塞业务处理。
|
||
|
||
**替换旧系统**
|
||
|
||
采用“影子对拍 → 切流 → 旧系统冻结”。共享信箱不能让新旧系统同时认领和回填;影子输入使用只读水位或回放。影子环境须隔离 PG schema/实例、Kafka topic 和服务注册身份,并禁止误写生产信箱。切流时保证只有一个权威写者。
|
||
|
||
**可观测性要求**
|
||
|
||
使用消息 ID、事件 ID 关联处理与投递日志;健康检查反映依赖实际可用性,而不只是进程存活。运行中重点关注队列积压、队头滞留时间、投递延迟、重试/DEAD 数量和回填补偿积压。死信和一致性异常需要可执行的告警与重放流程,不能只留一条错误日志。
|
||
|
||
## 9. 当前实现与上线门槛
|
||
|
||
当前已实现收报入队与判重、严格 FIFO 主泵、SCHD/FLOP/FDEL/ADFT 处理器与 PG 单事务写入、outbox 与 `schd` 聚合投递骨架、回填补偿与航班历史清理脚手架。**这些只证明机制可用,不证明生产链路已闭环**:默认配置不自动启动管道,真实数据库、信箱与出站适配必须显式开启(`msgx.pipeline.autostart`、`mailbox.shared-mysql.enabled`、真实 `DeliveryPort` 适配)。
|
||
|
||
上线前必须完成并验证:
|
||
|
||
- 真实 PG + 共享 MySQL 信箱的端到端处理、补偿与投递,以及出站信箱适配;未闭合的缺口与不可声明项见 [invariants.md](invariants.md) 的声明边界。
|
||
- 航班状态不变量与恢复证据:`STATE_VERSION` 推进、`OPERATION_DAY` 不可变、故障中断回滚(`INV-12`–`INV-16`,语义见 flight-state.md)。
|
||
- FIFO 越序、身份去重、FDEL/ADFT、清场顺序与投递故障的回归测试(见 [invariants.md](invariants.md) 验证映射)。
|
||
- 单实例排他保护与启动校验、影子隔离、Kafka 生产配置约束——当前配置允许环境变量覆盖 `acks`/幂等/in-flight,且默认 in-flight 值与 D3 不同,切流前必须按 D3 收敛。
|
||
- 死信与一致性异常的告警、可执行的人工重放流程、端到端追踪、积压指标与安全边界。
|
||
|
||
验收与进度由 Plane 跟踪,缺口逐项见 [invariants.md](invariants.md) 的声明边界与 [user-stories.md](user-stories.md);航班状态规则统一以 flight-state.md 为准。
|