Files
msgexchange-v2/docs/architecture.md
T

18 KiB
Raw Blame History

msgexchange-v2 架构文档

系统角色:机场 OMMS 上游报文处理中间件——消费 CIIMS/AODB 等上游写入共享信箱的 XML 报文,经 FIFO 管道处理后向下游投递;报文源系统。 架构基线为 Plane airport_chengdu_msgexchange_v2ACM2 项目的 ACM2-3(综合架构 v4 脚手架跟踪 ACM2-4,评审与实施计划(U01U30ACM2-10。本文是仓库内的架构速览, 与代码同步维护。有效口径 = ACM2-3 未被取代的内容 + 后续决策 ACM2-12,不能只按 ACM2-3 的历史正文回改本文。 存储边界与阶段 B 以 ACM2-12 为准(ACM2-11 仅为决策史;自有 PostgreSQL + 共享 MySQL 信箱 + Redis 动态/gen;阶段 B 缓做)——ACM2-3 中"同库事务锚 / MySQL 六辅助表 / 阶段 B"表述 已被上述决策修订。 配套设计细节见 design.md

1. 系统定位

机场 OMMS 上游报文处理中间件:消费 CIIMS/AODB 等上游写入共享信箱的 XML 报文, 经严格 FIFO 管道解析、决策、维护航班动态权威态,并向下游(Kafka / 出站信箱 / 查询接口) 投递;报文源系统,替代 CIIMS 或 AODB。替换 legacy msgexchange-api Java 8 / Spring Boot 1.5 / Maven)。过渡策略为双跑三步

影子对拍(共享库水位/回放 + 自有 PG 独立 schema 比对)→ 切流(nextgen 权威)→ 旧仓库冻结
  • legacy 维护不受本仓库影响;本仓库不声明 legacy 旧表 schema且不在共享 MySQL 建任何表 CMINMSGS/COUTMSGS 所在库属他人系统,本系统仅信箱 DML——ACM2-12,见 §6)。
  • wire 契约冻结:消息结构唯一事实源为 docs/legacy/SIS_AODB_RMS-V0.1.md + docs/legacy/unisysaodbsis.xsd; HTTP 端点路径与响应语义沿用现役(如 compat 写路径 POST /cminmsgs/send 返回记录 ID)。
  • 影子对拍口径(ACM2-12:共享库是单信箱无法"同入口双收",影子输入改以共享库只读 水位/回放 + 自有 PG 独立 schema 比对(详见 §8 与 ACM2-12 Checks ⑤)。

1.1 上下游边界(中间件职责)

方向 角色 本系统做什么 本系统做什么
入站(主路径) CIIMS / AODB 等上游 JDBC 轮询共享 MySQL CMINMSGS 发现新信 → 自有 PG 入队 → 解析处理 不生成原始业务报文;不替代 CIIMS 落信
入站(compat 手工工具 / 对拍 HTTP POST /cminmsgs/send 写信箱 + PG 入队 非生产主拓扑
处理 本系统 维护 Redis 航班动态权威态;Handler 纯函数决策 不持有 AODB 主数据编辑权
出站 下游消费者 Kafkamsg/schd)、共享 MySQL COUTMSGS、查询 HTTP 不直接推送至前端(经 Kafka 等中转)

与 SIS / legacy 一致:上游经 CIIMS 等外部系统写入 CMINMSGS 后,本系统以 1s JDBC 轮询 DATE_PROCESSED IS NULL)采集并处理——与 legacy MsgExchangeRunner 同口径。

2. 技术栈

选型 说明
语言/运行时 Kotlin 2.3 + JDK 25 JDK 21 不可行(Micronaut 5.1 系要求 JVM 25+ACMA-9 实测)
框架 Micronaut platform BOM 5.1.3core 系实际解析 5.1.13,classpath 混用;版本重钉属 U02/U05 编译期 DIKSPkotlin-ksp + micronaut-inject-kotlin 5.1.3)生成 *$Definition
持久化 自有 PostgreSQL(全部内部状态)+ 共享 MySQL 信箱 自有库:消息管道 PROC_STATE/MSG_EVENT + PUMP_JOB/REQ_TRACK + 21 类 REF_MASTER(迁移 db/migration);共享库主契约为 CMINMSGS/COUTMSGS DML。JDBC 仓储与 InboxPoller 已有初版,事务/补偿/出站适配仍属 U05(ACM2-12
权威存储 Redis(航班动态 flightInfo + 快照 gen 仅主泵线程写(I5);Lua 原子覆盖/版本推进(gen 协议重设计属 U09)
投递 Kafkaacks=all + 幂等;切流前须确认 Broker 支持 InitProducerId(22),严禁非幂等降级) outbox 模式,经 MSG_EVENT 表中转
投影(阶段 B Elasticsearch(历史) 阶段 B 缓做(ACM2-12FLIGHT_STATE 不落表,Redis 永续动态权威;历史投影链路不变
注册中心 EurekaMicronaut 原生键) 服务名契约 msgexchangeapi(影子 msgexchangeapi-shadow)——U17 未落地:当前注册名仍取 micronaut.application.name=msgexchange-nextgen),msgx.service-name 无运行时消费方(见 §8 与 design.md §9
可观测 logstash TCPAsync 包装)+ MDC traceId + 自定义健康指示器 design.md §8

3. 总体拓扑

 CIIMS/上游 ──外部写(他人系统)──▶ 共享 MySQL CMINMSGS(信箱)
                                        │
                                        │ JDBC 轮询/重扫(① 发现 DATE_PROCESSED IS NULL 新信)
                                        ▼
                    ┌──────────────────────────────────────────────────┐
                    │                msgexchange-nextgen               │
                    │                (单实例 · 单写者)                │
                    │ ingressInboxPoller + 补偿重扫,U05           │
                    │   └──② 自有 PG PROC_STATE(PENDING)              │
                    │      (跨库非同事务;② 失败→① 重扫补建)          │
                    │   compatPOST /cminmsgs/send ──▶ 信箱 insert   │
                    │      + PG 入队(现役 HTTP 写路径,U16 对拍)        │
                    │                                                  │
                    │ processingmsgx-pump 线程,严格 FIFO 队头)       │
                    │   Pump ──tick──▶ MessageProcessor                │
                    │      │                │ decodeXmlCodec        │
                    │      │                │ identity 绑定(I3       │
                    │      │                │ Handler.decide(纯函数)   │
                    │      │─Schd RESP/DNLD▶ SnapshotFlow(流程4      │
                    │      │─PUMP_JOB───▶ JobExecutor(作业窗口,决策1 │
                    │      │                                           │
                    │      ├────Redis Lua──▶ Redis flightInfoA权威)  │
                    │      ├──自有 PG 事务2──▶ MSG_EVENT + SUCCEEDED     │
                    │      └──共享 MySQL 回填 DATE_PROCESSED(外部副作用)│
                    │                                                  │
                    │ deliverymsgx-dispatcher 线程,每 target FIFO   │
                    │   Dispatcher ──逐条──▶ Kafka(msg)                │
                    │      └─flushSchd 聚合─▶ Kafka(schd)              │
                    │     (阶段 B 追加:ES flight_hts → Redis 投影删除) │
                    └──────────────────────────────────────────────────┘
          共享 MySQL(信箱)◀──① 轮询发现 / 回填──▶  自有 PostgreSQL ◀──② 管道状态
          Redis(动态+gen)◀── Lua 写 ── processing
                                        │                │
                                        ▼                ▼
                                   下游 Kafka topic    Eureka / logstash

要点:

  • 收报主路径(与现役/SIS 一致):上游经 CIIMS 等外部系统写入共享 MySQL CMINMSGS (本系统不建表);ingressJDBC 轮询DATE_PROCESSED IS NULL1s 节律,与 legacy MsgExchangeRunner 同口径)发现新信 → 自有 PG 建 PROC_STATE(PENDING)。 PG 入队失败时以共享库水位重扫补建U05)。POST /cminmsgs/send 为现役 HTTP 路径(手工/对拍),非上游报文到达的主拓扑。

  • 三条专用 daemon 单线程msgx-inbox-poller / msgx-pump / msgx-dispatcher)由 PipelineLifecycleServerStartupEvent 后拉起,不占用 Netty event loop;停机 requestStop + interrupt + joinU07)。仅当 msgx.pipeline.autostart=true 时装配——生产默认关, 当前属有意脚手架门禁(生产可运行需先完成 U05 数据层实装,见 §7)。

  • 单写者约束(I5:阶段 A 全部 Redis 写集中在主泵线程;实例数必须为 1 (运行期租约/选主保护属 U26,尚未实装,当前靠部署拓扑约束)。

  • 消息严格 FIFO,作业采用窗口语义(决策 1 经 ACM2-12 修订):消息按最小未完成 CMINMSGS_ID 保序;PUMP_JOB 不与消息构成统一全序,仅在无消息队头或队头处于退避窗口时执行。 当前 jobBefore 的 head-state 判定已近似该语义,但窗口、饥饿边界和测试仍待 U15 固化。

  • 存储边界(ACM2-12:本框图内事务库为自有 PostgreSQLPROC_STATE/MSG_EVENT/ PUMP_JOB/REQ_TRACK/21 类);CMINMSGS/COUTMSGS 在共享 MySQL 信箱(上游外部写、 本系统 JDBC 轮询读 + 处理回填写);Redis 除 flightInfo 还承载快照 gen。图中跨库步骤 (① JDBC 发现 / ② PG 入队 / 外部回填)与 ACM2-3 原「同库事务锚」语义不同,见 §6: 本地事务只在自有 PG,信箱交互为外部读/写 + 最终一致。

4. 模块职责

职责 对应 ACMA-8 主要类
ingress/ 收报:JDBC 轮询共享信箱发现新信 → 自有 PG 建 PENDING+ 补偿重扫;HTTP 写路径 compat);不解析报文 流程 1I3 InboxPollerU05InboxController InboxService
processing/ 主泵:FIFO 领取、ignoreMsg、identity 绑定、纯函数决策、RESP/DNLD 快照、自有 PG 事务2 流程 2/4I1/I2/I5 Pump MessageProcessor SnapshotFlow Identity Handler(Registry)
delivery/ 投递:每 target 严格 FIFO、schd 聚合 流程 3 Dispatcher SchdAggregation
jobs/ 泵作业:清场/归档/投影重建(PUMP_JOB 自有 PG,作业窗口执行) 流程 4/5/7I4 JobExecutor HistorySweepJob ArchiveJob ProjectionRebuildJob
codec/ XML 解码 + 失败分类(MALFORMED vs CODEC_ERROR 决策 4 前置 XmlCodec DecodeResult
domain/ 状态机枚举、事件/决策模型、Phase 开关 I1I5 ProcState MsgEvent Decision MsgKind
infra/ 仓储接口、重试策略、Redis Lua、stub、健康、日志 数据模型节 见 design.md
config/ PipelineProps 参数表(ACMA-8 参数初值) PipelineProps

5. 关键架构决策

# 决策 落点
D1 消息保持严格 FIFO;PUMP_JOB 为独立持久队列,不与消息组成统一全序,仅在无消息队头或队头退避窗口执行 Pump.tick / JobExecutor
D2 阶段 B(缓做,ACM2-12):ES 投递成功后同线程同步 enqueue 删除事件(不轮询 ack Dispatcher.tick(定案 2;阶段 B 评估后启用)
D3 schd 唯一出口是 flushSchd 批量聚合(逐条循环显式排除 KAFKA_SCHD) Dispatcher.tickU06/N03
D4 未实装 ≠ 非法:无 handler / staging 未实装 → FAILED(UNSUPPORTED) 可重放,绝不写终态 MessageProcessor SnapshotFlowU10/N21
D5 失败迁移在持有具体 head/batch 的边界完成;loop 只作最后防线,不吞 InterruptedException/Error MessageProcessor DispatcherU08
D6 接口驱动 + 假仓储单测;时间一律经可注入 Clock infra/persistence FailureScheduler
D7 stub 装配门禁:msgx.stubs=true 才装配内存实装;与 autostart 组合支撑 dev 冒烟 infra/stubU07/U01
D8 编译期 DI(KSP)+ 启动期冒烟测试锁定 BeanDefinition 生成 build.gradle.ktsU01

6. 数据边界(ACM2-12 最终口径)

  • 自有 PostgreSQL(本系统唯一自有数据库):全部内部状态,本地事务只在此库成立。 db/migration/V1.0.0__own_pg_pipeline.sqlPG 方言)建:PROC_STATE(取消息侧:处理 状态/重试/毒丸)+ MSG_EVENT(发消息侧:outbox+ PUMP_JOB(作业调度)+ REQ_TRACK 15 类请求)+ REF_MASTER21 类静态,SOURCE 审计)。
  • 共享 MySQLcdairport,他人系统库)——本系统不建任何表/schema,仅信箱 DML 上游外部写 CMINMSGS;本系统 JDBC 轮询读 + 处理完成回填 DATE_PROCESSED/STATUS 出站写 COUTMSGS(他人系统读取发送)。与信箱的交互是外部副作用,非本系统事务的一部分: 收报主路径=JDBC 轮询发现新信 → 自有 PG 建 PENDING 入队;compat HTTP 写= 信箱 insertRaw 成功(返回 CMINMSGS_ID)→ PG 入队;PG 建行失败以共享库 DATE_PROCESSED IS NULL 重扫补建;回填 = 处理成功后 PG 本地事务外异步/持久化补偿,最终一致(ACM2-12 / ACM2-19)。
  • Redis:航班动态 flightInfo(阶段 A 权威,I5+ 快照 SCHD_GEN(Lua 内原子 「覆盖+按代差删+版本推进」,协议重设计属 U09)。
  • 阶段 BFLIGHT_STATE):缓做,不落表;ES 仍承载历史航班(判史/查询), 相关投影/删除事件链待阶段 B 重评估后启用。
  • legacy 旧表 schema 归 legacy 仓库维护;CMINMSGS/COUTMSGS 的结构与保留策略由共享库方管理。 ARCHIVE 归档定案(ACM2-17:严禁向共享 MySQL 写入 CMINMSGS_HST(共享库严格保持 CMINMSGS 读/回填、 COUTMSGS 写入两表契约);终态入站消息归档目标确定为自有 PG PROC_STATE_HST

7. 权威(阶段 A Redis;阶段 B 缓做)与当前就绪度

阶段 权威 投递目标 状态
Amsgx.phase=A Redis flightInfo+ gen KAFKA:msg、KAFKA:schd 管道骨架+重试闭环已实装;PG/JDBC 信箱轮询已有初版,但水位、事务、回填补偿与出站信箱尚未闭环;gen→Redis 协议属 U09
B(缓做,重新评估后命名/启用) Redis 仍为动态权威 ES 历史写入与成功集清场 FLIGHT_STATE 不落表;HISTORY_SWEEP 标为 DEFERRED,不作为阶段 A 切流门禁

就绪度(2026-09-07 文档审查口径):现有 JUnit XML 报告记录 39 项测试全绿,README 记录 dev stub 进程级冒烟曾通过; 本轮因沙箱无法写用户级 Gradle 缓存,未重新证明该结果。当前已有 JDBC PG 仓储、共享信箱适配器 与 InboxPoller 初版,但没有覆盖跨库补偿和 PG 本地事务的集成验证。dev stub 冒烟路径可端到端 ./gradlew run 无外部依赖启动 → compat HTTP 写路径返回 200 → /health UP,修复记录见 README「进程级 dev 冒烟」); 生产默认配置不会启动处理管道autostart=false,数据源/信箱也默认 disabled)。生产就绪前置: U05(事务、补偿、出站与集成验证)、U07 fail-fast 定案、U09(快照恢复协议)、 U13(投递毒丸补全)、U15(统一序号)。逐项状态见 ACM2-10「定稿实施计划」。

8. 部署与安全姿态

  • 实例数 = 1(主泵单写者前提);双实例误配当前无运行期防护(U26:租约/DB 锁 + 拒启,未实装)。
  • 影子隔离(目标态;U17/U26 未落地,勿按现状引用):服务名(msgexchangeapi-shadow+ 独立 schema + Redis key 前缀 + 独立 topic 三层隔离;当前代码仅 msgx.register-eureka=false 生效—— Kafka topic 写死字面量 "msg"/"schd"Dispatcher)、FlightRedisClient.eval 无 key 前缀参数、 服务名未接 msgx.service-name(§2)。影子数据库:自有 PG 开独立 schema/实例;共享信箱无法 双写,影子输入=只读水位/回放口径(ACM2-12 Checks ⑤)。
  • 网络信任模型/cminmsgs/send 无鉴权(沿用现役内网信任姿态);eureka default-zone 回退 127.0.0.1:8761;口令/端点全部环境变量外置(零入库)。安全节细化属 U28。
  • 管理端点Micronaut 5.1 下 /env 默认禁用/beans 默认 enabled+sensitivedev/影子经 顶层 endpoints.* micronaut.endpoints.*——实测前缀错误时不生效)放开 /env/beans 与 health 明细。工程未引入 micronaut-securitysensitive 的实际拦截行为待 U28 定案。
  • 同名单风险:影子与生产同名同路径时,compat HTTP 写可能误写生产信箱——切流前必须核对服务名三隔离。

9. 可观测性

  • 日志logstash TCP JSON 通道(Async + neverBlock 降级,logstash 不可达不阻塞业务线程); 结构化生命周期日志(入队/PENDING/SUCCEEDED/SKIPPED/FAILED/DEAD/毒丸/flush 批次);MDC traceId (当前 = cminmsgsId/eventId,处理片段;贯穿入队→投递属 U12 遗留)。
  • 健康/health 聚合 redis-flight-store / kafka-delivery 自定义指示器——真实 ping 判定(false/异常→DOWN,缺 bean→DOWN),非仅 bean 存在。
  • 指标缺口micrometer 队列深度/投递延迟 gauge 未引入(版本对齐待 U05 批次); DEAD/DLQ 告警出口与一致性哨兵实装(U25)未落地——告警当前以 ERROR 日志为落点。