Files
msgexchange-v2/docs/design.md
T
windyboyandCursor 7ccdd00a31 docs: 统一中间件定位与 JDBC 轮询主路径口径
Align README, architecture, design, config comments, and ingress docs
with the upstream message-processing middleware model: external CIIMS write
to CMINMSGS, JDBC poll as production ingress, HTTP /cminmsgs/send as compat.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-07 10:27:44 +08:00

20 KiB
Raw Blame History

msgexchange-v2 设计文档

系统角色:机场 OMMS 上游报文处理中间件——消费 CIIMS/AODB 等上游经共享 MySQL 信箱 (CMINMSGS)投递的 XML 报文,解析处理后维护 Redis 动态并向 Kafka / 出站信箱投递; 报文源系统。生产主路径 = JDBC 轮询发现新信;HTTP POST /cminmsgs/send = compat 写路径。 本文对应仓库当前实现,给出模块级设计语义与依据;架构总览见 architecture.md,权威架构为 Plane ACM2-3,实施计划与逐项验收为 ACM2-10U01U30。文中标注「TODO/未实装」的条目均为已知开放项,不属文档遗漏。

0. 系统边界速览

上游(CIIMS/AODB…) ──外部写──▶ 共享 MySQL CMINMSGSDATE_PROCESSED IS NULL
                                    │ JDBC 轮询/重扫(InboxPollerU05
                                    ▼
                         自有 PG PROC_STATE(PENDING) ──▶ 主泵 FIFO 处理
                                    │
                    ┌───────────────┼───────────────┐
                    ▼               ▼               ▼
              Redis 动态      Kafka msg/schd   COUTMSGS 出站
              (阶段 A 权威)   (下游订阅)      (他人读取发送)

compatPOST /cminmsgs/send ──▶ insertRaw + PG 入队(手工/对拍,非主拓扑)
  • 中间件定位:本系统位于 CIIMS 与下游消费者之间,负责采集 → 解析 → 决策 → 投递; 报文原文由上游写入共享信箱,本系统只读(主路径)或 compat 写(辅助)。
  • 与 legacy 对齐legacy MsgExchangeRunner 同样以 1s 轮询 CMINMSGS 为处理入口; legacy HTTP 收报接口在 nextgen 中保留为 compat,不改变生产主拓扑。

1. 领域模型

1.1 状态机与错误分类

ProcStatusPROC_STATE.STATE,消息处理侧):
  PENDING ──处理成功──▶ SUCCEEDED(终态;回填共享库 CMINMSGS 为外部副作用,最终一致)
     │ ──同 identity 已绑定──▶ SKIPPED(终态,lastError=duplicate-of:<id>
     └──失败──▶ FAILED(非终态,attempts+1 + nextAttemptAt 退避)
                   │ attempts ≥ maxAttempts 或 队头滞留超 head-deadline
                   ▼
                DEAD(终态/DLQERROR_CLASS=EXHAUSTED 规范化)

EventStatusMSG_EVENT.STATE,投递侧):
  PENDING ──▶ SENT;失败退避回 PENDINGattempts 耗尽整批/单条 → DEAD(DLQ)

ErrorClass(两侧共用):
  MALFORMED     报文非法        → 直接 DEAD,永不重放
  CODEC_ERROR   可随 codec 修复  → FAILED 可重放(白名单内)
  UNSUPPORTED   能力未实装      → FAILED 可重放(白名单内)
  INFRA         基础设施抖动     → FAILED 可重放(白名单内)
  EXHAUSTED     重试耗尽(终态规范化类)→ 人工复核后可重放(白名单内)
  • 「未实装 ≠ 非法」是通用规则(D4):无 handler、快照 staging 未实装都写 FAILED(UNSUPPORTED),绝不写终态——阶段 2 前接入流量不会把报文变砖。
  • 显式重放入口 ReplayService:仅白名单 {CODEC_ERROR, UNSUPPORTED, INFRA, EXHAUSTED} 可从 FAILED/DEAD 回 PENDING ATTEMPTS=0、NEXT_ATTEMPT_AT=NULLerrorClass/lastError 保留审计); 含 MALFORMED 的请求对该类静默忽略。运维接口(controller/runbook)属 U11 遗留。

1.2 报文模型(sealed 分派)

  • MetaFields(sndr, type, styp, seqn, dttm)——实名沿用 legacy META.java。
  • MsgKind sealedSchd(RESP|DNLD|ADFT) + Flop(29 类 STYP)typeTag 产出 SCHD-XXX / FLOP-xxx,与 HandlerRegistry.keyOf 同源(查表键=日志类型,禁止分叉)。
  • 解码失败二分:DecodeResult.Err(MALFORMED) → DEADErr(CODEC_ERROR) → FAILED 退避 T06/U11)。
  • Handler 为纯函数decide(flightView, msg) → DecisionflightChanges / msgNotifies / schdPush / outboundIntents / refUpserts),不触碰 Redis/Kafka——副作用全部由泵边界执行。
  • Handler 实装:0/32(骨架),翻译属阶段 2/3,逐条对照 ACM2-4 行为基线与 KEEP/FIX 矩阵。

2. 数据模型(自有 PostgreSQL · ACM2-12

db/migration/V1.0.0__own_pg_pipeline.sql(PG 方言,自有库;legacy 旧表与共享库表不在 本仓库声明,见 architecture.md §6):

角色 关键列/约束
PROC_STATE 取消息侧:处理伴生状态/重试/毒丸(与共享库 CMINMSGS_ID 对应) uk_proc_identity(IDENTITY_KEY) 唯一约束=I3 依据;idx_proc_head(STATE, CMINMSGS_ID)=队头
MSG_EVENT 发消息侧:统一投递 outbox EVENT_ID 自增=全序;idx_evt_head(TARGET, STATE, EVENT_ID)=每 target 队头
PUMP_JOB 泵作业调度(作业不插队,队头空闲/退避窗口执行) kindARCHIVE/HISTORY_SWEEP/PROJECTION_REBUILD
REQ_TRACK 15 类请求状态机 REGISTERED/SENT/WAITING/DONE/EXPIREDCOUTMSGS_ID BIGINTU18 修正)
REF_MASTER 21 类静态主数据 (RTYPE,RKEY) PKSOURCE=ADMINAPI/AODB/PIPELINEREFRESHED_AT

已知 DDL 缺口(U18,随本库 PG 化修正/收窄):REQ_TRACK.COUTMSGS_ID 已按 BIGINT 时间列须 DATETIME(6)/显式 UTC 口径在 U05 数据层实现时定;FLIGHT_STATE 因缓做不在本库。

存储边界(ACM2-12 定案)

  • 自有 PostgreSQL = 上表全部(消息管道 + 调度 + 请求 + 21 类)。本地事务只在此库: 处理侧「MSG_EVENT 插入 + PROC_STATE→SUCCEEDED」同事务;其余跨存储一律外部副作用。
  • 共享 MySQLcdairport,他人系统)仅信箱 DML、不建表:上游外部写 CMINMSGS; 本系统 JDBC 轮询读 + 处理回填;出站写 COUTMSGS(他人读取发送)。见 §3.1/§3.2 的事务模型。
  • Redis:航班动态 flightInfo + 快照 SCHD_GENgen——Lua 内原子「覆盖+按代差删+ 版本推进」;重放幂等由 Lua 承接(协议重设计属 U09),RefDataRepository 为目标实现的 过渡占位接口。
  • FLIGHT_STATE(阶段 B 权威):缓做不落表Redis 永续动态权威)。
  • 实现状态:迁移 SQL 已按 PG 落地(V1.0.0);Repository 接口归属注释已对正(自有 PG / 信箱封装 / gen→Redis 占位 / FlightState 缓做);Micronaut Data 实装与信箱适配层 CminmsgMailbox/OutboxMailbox)属 U05 批次。

3. 核心流程设计

3.1 流程 1:收报(JDBC 轮询 + HTTP compat

主路径(生产/SIS 口径):上游经 CIIMS 等外部系统写入共享 MySQL CMINMSGS DATE_PROCESSED IS NULL);本系统 ingressJDBC 轮询发现新信(与 legacy MsgExchangeRunner.getNewMsgsAfterId 同语义,1s 节律),自有 PG 入队:

  1. pollNew()U05InboxPoller):JDBC 查共享库 CMINMSGS_ID > watermark 且 DATE_PROCESSED IS NULL(及/或 PG 无对应 PROC_STATE 的补偿重扫);
  2. procState.insert(id):自有 PG 建 PENDING 行入队;本步失败 → 下轮重扫补建;
  3. 不解析报文、接收层无唯一约束(I3);wakePump() 为 TODO 空操作——主泵 1s 轮询兜底。

compat 路径(现役 HTTP 写)InboxService.acceptPOST /cminmsgs/send= 共享信箱 insertRaw + 自有 PG 入队(跨库,非同一事务);「已持久化」响应语义与现役 对拍(U16)。用于手工注入/影子对拍,上游报文到达的主拓扑。

3.2 流程 2:主泵 tickPump.tick

每 tick 按序判定:

  1. headQueued() 取作业、headUnfinished() 取最小未完成 CMINMSGS_ID含 FAILED, 消息间严格保序:队头退避未到期即 sleep 至到期点,后方消息永不越队)。
  2. jobBefore(job, head)head 为空或队头 FAILED 时作业先行——当前为 head-state 近似 (非入队时间排序),与「统一 FIFO」注释存在已知偏离(U15 未实装);ACM2-12 口径为 作业窗口执行(队头空闲/退避窗口),不追求与消息统一全序(Checks ④)。
  3. 队头 FAILED 且退避未到期:poisoned() 判定(attempts≥maxAttempts 或滞留超 head-deadline 10m)→ DEAD(EXHAUSTED) 毒丸升级(Pump 侧;投递侧同语义属 U13,未实装); 实现注headDeadline 判据以 updatedAt 为锚,而 updatedAt 与 nextAttemptAt 同一次 FAILED 写入、退避 ≤60s(封顶)→ 有 nextAttemptAt>now 必有 nowupdatedAt≤60s<10m滞留超时分支实际 不可达,仅 attempts 维度生效;锚点语义随 U05/U13 定案并补测试(Pump 未注入 Clock,主泵级 门禁无单测锁定,见 §7);否则 sleep 至 nextAttemptAt。
  4. 正常队头 → MessageProcessor.processOne 入口守卫(FAILED 且已 exhausted → DEAD)→ rawOf 缺失 → DEAD(MALFORMED) → decodeMALFORMED→DEAD / CODEC_ERROR→FAILED)→ identity 首绑 tryBindIdentity 失败 → SKIPPEDI3)→ Schd DNLD → SnapshotFlow → 其余 → Handler.decide(redis.hgetAllFlightInfo(), msg) → 阶段 ARedis 先写(I2 happens-beforeTODO redisApply)→ 自有 PG 事务 2 MSG_EVENT 插入 + PROC_STATE→SUCCEEDED(同库原子,@Transactional); CMINMSGS 回填(DATE_PROCESSED/STATUS)为共享信箱外部回填:目标态为 PG 提交后 异步/补偿执行,失败重试+告警(最终一致,ACM2-12);当前实现为同步内联占位 insertAll → backfillOnSuccess → SUCCEEDED,无 @TransactionalU05 定案并改)。
  5. 异常边界(U08):processOne 内 try/catch → ProcFailure.fail(INFRA)attempts+1、 退避、达上限 DEAD);InterruptedException 恢复中断位后上抛loop 仅 catch Exception 作最后防线,Error 任其终止进程(异常必可见)。

幂等键(I3):SNDR|TYPE|STYP|SEQNIdentity.of 唯一入口);「含日边界」可配置且 默认关闭SEQN 重置作用域 CONFIRM 前,上线后不改幂等键)。

3.3 流程 3:投递(Dispatcher

  • 每 target 严格 FIFOI1 双层同策略):headUnsent 取队头;队头退避未到期 → 等待不跳过。
  • 逐条循环显式排除 KAFKA_SCHDD3/U06):schd 唯一出口 flushSchd—— claimBatchORDER BY EVENT_ID)→ 按 FLID 分组取 max(EVENT_ID)FIX:现役 buffer 无去重会重发旧值)→ FLTR JSON 数组一次发出;失败整批 attempts+1 退避、队首未到期不 claim、lastFlush 仅成功后推进;达上限整批 DEAD(DLQ)。周期/批上限取参数表 3s / 500)。
  • 轮询间隔取参数表(下限 50ms,N18);无 200ms 硬编码。
  • 阶段 B(定案 2/D2,ACM2-12 缓做):ES 投递成功 → 同线程同步 insertSync 删除事件 deleteOfJackson 结构化序列化,refs 可空恒合法 JSON——U14)。当前不启用。

3.4 流程 4:日计划快照(SnapshotFlow

staging(流式解析+整包校验,TODO 阶段2;未实装→FAILED(UNSUPPORTED)
  → Redis Lua SNAPSHOT_REPLACE(同一 hash 原子「覆盖新代+按代差删」,删除集=旧代flids−新代)
  → gen 版本推进(ACM2-12gen 随 flightInfo 同在 RedisLua 内原子版本 CAS——
     目标实现;现 RefDataRepository/putGenIfVersion 为过渡占位)
  → 自有 PG SUCCEEDED

已知缺口(U09,未定案,ACM2-12 后重设计为 Redis 内协议)gen 与 Lua/SUCCEEDED 不再 分属两存储即可同原子(全部在 Redis Lua);真正跨存储的窗口收窄为「Lua 已完成、PG SUCCEEDED 未写」——重放判据(版本不二次自增)与按代差删在 Lua 内以版本 CAS 承接,恢复协议待定案并补 测试。现有代码的「CAS 重放二次自增」缺陷(版本 1→2)与实现注随协议重设计一并消除。

3.5 泵作业(JobExecutorPUMP_JOB 自有 PG,作业窗口执行)

  • HISTORY_SWEEP3:30 清场,I4 同步链):判史 → 同步写 ES → 仅删成功集。 占位门禁(U10/T07 修订)ES saveSync 接线前 pickHistory 恒空集、删除量恒 0, 禁止「全量可删」fail-open 默认;现役五条判史规则 golden 通过后才允许接线。
  • ARCHIVE3:00):1 天前且仅终态(SUCCEEDED/SKIPPED/DEAD)迁 CMINMSGS_HST(共享库, 外部副作用;TODO U05 批次)。
  • PROJECTION_REBUILD(阶段 B 缓做,ACM2-12;重新评估后再启用)。

4. 失败与重试统一设计(U08

组件 迁移语义
ProcState(处理/快照) ProcFailure + FailureScheduler attempts+1 → exhausted ? DEAD(EXHAUSTED) : FAILED+nextAttemptAt
MsgEvent 逐条(投递) Dispatcher.retryOrDead 同上(DLQ 保留行,attempts 审计)
MsgEvent 批量(schd flushSchd 整批 队首未到期不 claim;整批退避;达上限整批 DEAD
  • 退避表 [1s,2s,4s,8s,16s],单档封顶 backoff-cap-ms=60sattempt≤0 兜底首档(N28)。
  • 时间一律经可注入 java.time.ClockTimeFactory;测试用 MutableClock,无真实睡眠)。
  • loop 兜底 catch 不做状态迁移(迁移已在边界完成),仅防线程静默死亡。

5. 不变量与实现落点

不变量 语义 落点 状态
I1 单写者严格 FIFO + HOL 阻塞 + 毒丸升级 headUnfinished/headUnsent 队头语义、poisoned() 实装(attempts 毒丸生效;head-deadline 判据不可达待修,见 §3.2 注;job 为 head-state 近似 / 作业窗口属 U15
I2 Redis 先写、后于事件创建(happens-before processOne 阶段 A 分支 TODO redisApply(流程占位已留)
I3 identity 首绑幂等;接收层无唯一约束;SUCCEEDED 回填 Identity/tryBindIdentity/backfillOnSuccess 实装
I4 清场仅删 ES 成功集;按代差删 HistorySweepJob/SNAPSHOT_REPLACE delFields 门禁实装,ES 接线 TODO
I5 阶段 A Redis 写仅主泵线程;Delivery 不写 Redis 单线程拓扑 + Targets.phaseA 实装(拓扑约束,U26 运行期保护未做)

6. 配置参数(msgx.*ACMA-8 参数表初值)

默认 说明
phase A 阶段总开关(A/B 权威切换)
service-name msgexchangeapi 契约冻结;影子= msgexchangeapi-shadow
pipeline.poll-interval 1s 泵轮询节律(KEEP 现役)
pipeline.max-attempts 5 处理/投递同值
pipeline.backoff-ms / backoff-cap-ms 1s..16s / 60s 指数退避表与封顶
pipeline.head-deadline 10m 队头滞留上界(毒丸升级)
pipeline.autostart false 生命周期门禁:true 才装配 Pump/Dispatcher 线程(dev+stubs 开)
schd.flush-period / flush-limit 3s / 500 schd 聚合节律与批上限
identity.include-day-boundary false 幂等键日边界(CONFIRM 前禁开)
consistency-check.on-startup / daily-sample-ratio true / 0.01 一致性哨兵(实装属 U25

基础设施键位口径(Micronaut 5.1U03):datasources.default.*自有 PostgreSQL ACM2-12)、flyway.datasources.default.*mailbox.shared-mysql.*(共享信箱,仅 DML)、 kafka.producers.default.*eureka.client.*;logback 独立于本文件, 环境变量前缀 MSGX_LOGSTASH_*

7. 测试策略

  • 接口驱动 + 假仓储:管道语义全部离线单测(无 DB/Redis/Kafka),时间用 MutableClock
  • 不变量测试FIFO/HOL、schd 批退避与 DLQ、重试上限、重放白名单、 聚合最新态、配置绑定、DI 装配冒烟(PipelineSmokeTest:收报→FAILED(CODEC_ERRORstub codec 恒 CODEC_ERROR)→重放→schd 聚合发出)。 范围注:以上覆盖的是边界级(processOne/flushSchd/仓储)语义;主泵 tick 级 HOL/毒丸/退避 门禁因 Pump 未注入 Clock 而无单测锁定(§3.2 注),DispatcherTickTest 仅锁批退避与「队首未到期 不推进」。
  • 现状 37 测试全绿(./gradlew testwrapper 钉 9.6.1 + JDK 25;内网构建设 GRADLE_USER_HOME/TMPDIR 指向可写目录)。
  • U29 门禁(规划):不变量清单化入 CI 红即阻塞;「新增逻辑必伴生不变量测试」入贡献约定。

8. 可观测性设计(U12 已落地部分)

  • logstash TCP 经 AsyncAppenderqueueSize 4096 / neverBlock / discardingThreshold 0): logstash 不可达丢弃日志而非阻塞业务线程(N31b 顺序约束:先降级再补日志)。
  • MDC traceId = cminmsgsId/eventIdTraceLog.withTrace),覆盖处理/投递日志片段; 收报→投递全链贯穿与 micrometer gauge(队列深度/投递延迟)未实装。
  • 生命周期结构化日志:收报 INFO、SUCCEEDED/SKIPPED INFO、FAILED WARN、DEAD/毒丸 ERROR、 flush 批次 INFO、整批 DLQ ERROR——告警暂以 ERROR 日志为落点(DEAD 告警出口属 U13/U25)。

9. 已知缺口与定案待办(对照 ACM2-10)

缺口 计划
U05 自有 PG 数据层无实装 + InboxPollerJDBC 轮询/重扫) 未实装(当前仅 compat HTTP InboxService.accept);信箱适配层 CminmsgMailbox/OutboxMailbox@Transactional/allopen 未引入 → 生产 DI 装配失败;application-dev.yml 仍残留 MySQL datasource URL 占位(stub 下 enabled=false 不建连,U05 清理) U05 批次(Micronaut Data JDBC on PG + InboxPoller + allopen + 信箱外部副作用与补偿回归)
U07/U26 autostart 默认关=有意门禁,但生产无 fail-fast;双实例无运行期防护 fail-fast 定案 + 租约/DB 锁拒启
U09 gen→Redis 协议未重设计(Lua 内原子版本推进;崩溃窗口=「Lua 完成/PG SUCCEEDED 未写」) Redis 内版本 CAS + 恢复协议 + 测试(ACM2-12
U13 投递侧无 createdAt/headDeadline 超时升级、无 DEAD 告警出口;处理侧 head-deadline 判据亦不可达(见 §3.2 注) WP2
U15 job 与队头消息无统一全序(head-state 近似,非入队时间)——ACM2-12 口径:作业窗口执行,不追求与消息全序 作业窗口语义定稿(ACM2-12 Checks ④)
U16 /cminmsgs/send 无 @Consumes/字符集(实测 text/plain 415)、无错误路径契约(@ControllerAdvice WP2legacy 逐字对拍固化)
U17 Eureka 注册名仍取 micronaut.application.name=msgexchange-nextgen);msgx.service-name 无运行时消费方 → 影子/切流前注册名与文档契约脱节 micronaut.application.name=${msgx.service-name}application.yml
U18 DDL 缺口(随 PG 化收窄:REQ_TRACK.COUTMSGS_ID 已 BIGINT;时间列口径 U05 定) U05 批次
U19U21 请求状态机量纲/死分支、identity 绑定静默跳过、21 类静态接口/表对齐(REF_MASTER/SOURCE WP2
U22U24 载荷类型收敛、eventSeq 未接线、每报文全量读语义定案 WP3
U25/U28/U30 一致性哨兵实装、README 安全节/入口、索引与杂项 WP3/4
U27 专有材料(SIS md 703KB / XSD 版权头)治理决策 WP4ACL 核验先行)
ACM2-12 存储边界(自有 PG + 共享信箱 + Redis 动态/gen + 阶段 B 缓做):迁移 SQL/配置/接口注释已按定案调整(V1.0.0 PG);信箱适配层、gen Lua、作业窗口语义、影子重设计未实装 ACM2-12 Checks ①–⑥
ACM2-11 决策史(已被 ACM2-12 吸收):曾讨论 21 类静态独立 PG 参考库;定案为并入自有 PG REF_MASTER(见 ACM2-12),勿再按 datasources.reference 第二库规划 仅作决策脉络参考