Files
msgexchange-v2/docs/architecture.md
T
windyboy 8093d20f9b docs(acm2-63): 收敛文档口径并补齐作业健康指标登记
含:MAFL 派生投影与 SRVT/VIPF 无界集合口径、schd EVENT_ID 写代次与 DEAD 代次替换、运营日时区 fail-fast、缺口索引增删(G-MAFL/G-SRVT-VIPF/G-COMPAT-HTTP/G-REQ-OPEN-UNIQUE,移除 G-JOB-HEARTBEAT)、job 心跳/扫描积压/最近失败指标登记、INV-18 验证映射改为实际测试名。
2026-09-12 21:00:01 +08:00

139 lines
12 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 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
共享 MySQLCMINMSGS
│ 轮询未处理记录
┌──────────────── msgexchange-v2(单实例)────────────────┐
│ ingress:发现报文 → PostgreSQL 持久化入队 │
│ │ │
│ processing:取 FIFO 队头 → 解析 / 去重 → 处理器决策 │
│ └─ PG 单事务:航班变更 + 终态 + 待发事件│
│ │
│ jobs:独立维护线程(回填补偿 / 历史归档 / 留痕清理) │
│ delivery:读取 PG 待发事件 → 投递 / 重试 │
└─────────────────────────┬──────────────────────────────┘
├─ Kafkamsg / schd
└─ 共享 MySQLCOUTMSGS
处理结果提交后,再回填 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` 取出待发事件。普通事件按 `FLID``EVENT_ID` 保序(与分区键一致):同一 `FLID` 内不跳队,不同 `FLID` 不互相阻塞(实现仍为目标级全序,收敛见 `CLM-7`)。
`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``SRVT`/`VIPF` 专用明细尚未实现 `[G-SRVT-VIPF]`)、留痕 `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 为准。