docs: 收敛架构设计规范与用户故事实施清单 (ACM2-28)
This commit is contained in:
+192
-252
@@ -1,288 +1,228 @@
|
||||
# msgexchange-v2 设计文档
|
||||
|
||||
> **系统角色**:机场 OMMS **上游报文处理中间件**——消费 CIIMS/AODB 等上游经共享 MySQL 信箱
|
||||
> (`CMINMSGS`)投递的 XML 报文,解析处理后维护 Redis 动态并向 Kafka / 出站信箱投递;
|
||||
> **非**报文源系统。生产主路径 = **JDBC 轮询**发现新信;HTTP `POST /cminmsgs/send` = compat 写路径。
|
||||
> 本文对应仓库当前实现,给出模块级设计语义与依据;架构总览见
|
||||
> [architecture.md](architecture.md),架构基线为 Plane **ACM2-3**,其存储/事务边界由后续
|
||||
> **ACM2-12** 覆盖;实施计划与逐项验收为
|
||||
> **ACM2-10(U01–U30)**。文中标注「TODO/未实装」的条目均为已知开放项,不属文档遗漏。
|
||||
## 1. 阅读说明
|
||||
|
||||
## 0. 系统边界速览
|
||||
本文说明模块如何协作、状态如何流转,以及失败后如何恢复。系统范围、存储归属和部署约束见 [architecture.md](architecture.md),不在这里重复。
|
||||
|
||||
```
|
||||
上游(CIIMS/AODB…) ──外部写──▶ 共享 MySQL CMINMSGS(DATE_PROCESSED IS NULL)
|
||||
│ JDBC 轮询/重扫(InboxPoller,U05)
|
||||
▼
|
||||
自有 PG PROC_STATE(PENDING) ──▶ 主泵 FIFO 处理
|
||||
│
|
||||
┌───────────────┼───────────────┐
|
||||
▼ ▼ ▼
|
||||
Redis 动态 Kafka msg/schd COUTMSGS 出站
|
||||
(阶段 A 权威) (下游订阅) (他人读取发送)
|
||||
以下流程是阶段 A 的目标设计,不是实现完成清单。当前代码仍有占位和过渡实现,与设计的主要差异集中在第 10 节。阶段 B 的历史投影和清场暂不启用,也不改变 Redis 作为航班动态权威存储的定位。
|
||||
|
||||
(compat)POST /cminmsgs/send ──▶ insertRaw + PG 入队(手工/对拍,非主拓扑)
|
||||
```
|
||||
## 2. 数据与领域模型
|
||||
|
||||
- **中间件定位**:本系统位于 CIIMS 与下游消费者之间,负责**采集 → 解析 → 决策 → 投递**;
|
||||
报文原文由上游写入共享信箱,本系统只读(主路径)或 compat 写(辅助)。
|
||||
- **与 legacy 对齐**:legacy `MsgExchangeRunner` 同样以 1s 轮询 `CMINMSGS` 为处理入口;
|
||||
legacy HTTP 收报接口在 nextgen 中保留为 compat,不改变生产主拓扑。
|
||||
### 2.1 持久化记录
|
||||
|
||||
## 1. 领域模型
|
||||
所有内部表都属于自有 PostgreSQL;共享 MySQL 只保留约定的信箱读写边界。
|
||||
|
||||
### 1.1 状态机与错误分类
|
||||
|
||||
```
|
||||
ProcStatus(PROC_STATE.STATE,消息处理侧):
|
||||
PENDING ──处理成功──▶ SUCCEEDED(终态;回填共享库 CMINMSGS 为外部副作用,最终一致)
|
||||
│ ──同 identity 已绑定──▶ SKIPPED(终态,lastError=duplicate-of:<id>)
|
||||
└──失败──▶ FAILED(非终态,attempts+1 + nextAttemptAt 退避)
|
||||
│ attempts ≥ maxAttempts 或 队头滞留超 head-deadline
|
||||
▼
|
||||
DEAD(终态/DLQ,ERROR_CLASS=EXHAUSTED 规范化)
|
||||
|
||||
EventStatus(MSG_EVENT.STATE,投递侧):
|
||||
PENDING ──▶ SENT;失败退避回 PENDING;attempts 耗尽整批/单条 → 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=NULL,errorClass/lastError 保留审计);
|
||||
含 MALFORMED 的请求对该类静默忽略。运维接口(controller/runbook)属 U11 遗留。
|
||||
|
||||
### 1.2 报文模型(sealed 分派)
|
||||
|
||||
- `MetaFields(sndr, type, styp, seqn, dttm)`——实名沿用 legacy META.java。
|
||||
- `MsgKind` sealed:`Schd(RESP|DNLD|ADFT)` + `Flop(29 类 STYP)`;`typeTag` 产出
|
||||
`SCHD-XXX` / `FLOP-xxx`,与 `HandlerRegistry.keyOf` 同源(查表键=日志类型,禁止分叉)。
|
||||
- 解码失败二分:`DecodeResult.Err(MALFORMED)` → DEAD;`Err(CODEC_ERROR)` → FAILED 退避
|
||||
(T06/U11)。
|
||||
- Handler 为**纯函数**:`decide(flightView, msg) → Decision`(flightChanges / 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 | 泵作业调度(作业不插队,队头空闲/退避窗口执行) | kind:ARCHIVE/HISTORY_SWEEP/PROJECTION_REBUILD |
|
||||
| REQ_TRACK | 15 类请求状态机 | REGISTERED/SENT/WAITING/DONE/EXPIRED;`COUTMSGS_ID BIGINT`(U18 修正) |
|
||||
| REF_MASTER | 21 类静态主数据 | `(RTYPE,RKEY)` PK;SOURCE=ADMINAPI/AODB/PIPELINE;REFRESHED_AT |
|
||||
| `PROC_STATE` | 入站消息的处理状态、身份、重试次数和错误原因 | `CMINMSGS_ID` 主键防止重复入队;`IDENTITY_KEY` 唯一约束防止业务重复;按最小未完成消息 ID 取队头。 |
|
||||
| `MSG_EVENT` | 等待投递的事件(outbox) | `EVENT_ID` 决定投递顺序;`TARGET` 区分目标;`PARTITION_KEY` 在 `schd` 中为 `FLID`。 |
|
||||
| `PUMP_JOB` | 持久化维护作业 | 状态为 `QUEUED / RUNNING / DONE / FAILED`;不与业务消息共用排序序号。 |
|
||||
| `REQ_TRACK` | 上游请求及应答关联 | 保存请求类型、参数、出站信箱 ID、发送和完成时间;同类只允许一个开放请求。 |
|
||||
| `REF_MASTER` | 静态参考数据 | `(RTYPE, RKEY)` 唯一,`SOURCE` 记录数据来源。 |
|
||||
| `PROC_STATE_HST` | 终态处理记录的归档目标 | 属于目标设计,当前迁移尚未建表;不得改写为共享库历史表。 |
|
||||
|
||||
已知 DDL 缺口(U18,随本库 PG 化修正/收窄):REQ_TRACK.COUTMSGS_ID 已按 BIGINT;
|
||||
时间列须 DATETIME(6)/显式 UTC 口径在 U05 数据层实现时定;FLIGHT_STATE 因缓做不在本库。
|
||||
字段与索引定义以 `src/main/resources/db/migration/` 为准。报文原文仍从共享信箱读取,因此必须协调原文保留期,不能在消息尚需处理或重放时提前清理。
|
||||
|
||||
**存储边界(ACM2-12 定案)**:
|
||||
- **自有 PostgreSQL** = 上表全部(消息管道 + 调度 + 请求 + 21 类)。本地事务只在此库:
|
||||
处理侧「MSG_EVENT 插入 + PROC_STATE→SUCCEEDED」同事务;其余跨存储一律外部副作用。
|
||||
- **共享 MySQL(cdairport,他人系统)仅信箱 DML、不建表**:上游外部写 CMINMSGS;
|
||||
本系统 JDBC 轮询读 + 处理回填;出站写 COUTMSGS(他人读取发送)。见 §3.1/§3.2 的事务模型。
|
||||
- **Redis**:航班动态 flightInfo + 快照 **SCHD_GEN(gen)**——Lua 内原子「覆盖+按代差删+
|
||||
版本推进」;重放幂等由 Lua 承接(协议重设计属 U09),`RefDataRepository` 为目标实现的
|
||||
过渡占位接口。
|
||||
- **FLIGHT_STATE(阶段 B 权威):缓做不落表**(Redis 永续动态权威)。
|
||||
- 实现状态:迁移 SQL 已按 PG 落地(V1.0.0);Repository 接口归属注释已对正(自有 PG /
|
||||
信箱封装 / gen→Redis 占位 / FlightState 缓做);Micronaut Data 实装与信箱适配层
|
||||
(CminmsgMailbox/OutboxMailbox)属 U05 批次。
|
||||
Redis 保存 `flightInfo` 与快照代际元数据 `gen`。`gen` 记录快照所属日期、版本及该代航班集合,不属于静态参考数据。
|
||||
|
||||
## 3. 核心流程设计
|
||||
### 2.2 消息、身份与决策
|
||||
|
||||
### 3.1 流程 1:收报(JDBC 轮询 + HTTP compat)
|
||||
`XmlCodec` 将 XML 解码为 `DecodedMessage`,包含 `SNDR / TYPE / STYP / SEQN / DTTM` 元数据和业务载荷。`MsgKind` 区分 `SCHD` 与 `FLOP` 子类型;Handler 查找与日志类型标识使用同一套映射。
|
||||
|
||||
**主路径(生产/SIS 口径)**:上游经 CIIMS 等外部系统写入共享 MySQL `CMINMSGS`
|
||||
(`DATE_PROCESSED IS NULL`);本系统 `ingress` 经 **JDBC 轮询**发现新信(与 legacy
|
||||
`MsgExchangeRunner.getNewMsgsAfterId` 同语义,1s 节律),自有 PG 入队:
|
||||
|
||||
1. 目标态分两路:快路径按持久化 watermark 查询 `CMINMSGS_ID > watermark AND
|
||||
DATE_PROCESSED IS NULL`;补偿路径按受控周期重扫“未处理且 PG 无对应 PROC_STATE”的记录。
|
||||
watermark 仅在本批 PG 入队均已确认后推进。**当前初版** `InboxPoller` 固定 `afterId=0`,
|
||||
每轮全量扫描未处理记录并以 PG 判重,尚未实现上述水位与补偿频控;
|
||||
2. `procState.insert(id)`:自有 PG 建 PENDING 行入队;本步失败 → 下轮重扫补建;
|
||||
3. 不解析报文、接收层无业务 identity 唯一约束(I3);`PROC_STATE.CMINMSGS_ID` 主键负责
|
||||
轮询重扫幂等。当前无显式 `wakePump()`,主泵 1s 轮询兜底。
|
||||
|
||||
**compat 路径(现役 HTTP 写)**:`InboxService.accept`(`POST /cminmsgs/send`)=
|
||||
共享信箱 `insertRaw` + 自有 PG 入队(跨库,非同一事务);「已持久化」响应语义与现役
|
||||
对拍(U16)。用于手工注入/影子对拍,**非**上游报文到达的主拓扑。
|
||||
|
||||
### 3.2 流程 2:主泵 tick(`Pump.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 必有 now−updatedAt≤60s<10m,**滞留超时分支实际
|
||||
不可达**,仅 attempts 维度生效;锚点语义随 U05/U13 定案并补测试(Pump 未注入 Clock,主泵级
|
||||
门禁无单测锁定,见 §7);否则 sleep 至 nextAttemptAt。
|
||||
4. 正常队头 → `MessageProcessor.processOne`:
|
||||
入口守卫(FAILED 且已 exhausted → DEAD)→ `rawOf` 缺失 → DEAD(MALFORMED) →
|
||||
decode(MALFORMED→DEAD / CODEC_ERROR→FAILED)→ ignoreMsg 匹配(LDM/REGN/RSTA/EROR,
|
||||
命中→SKIPPED,回填遵循 US-09;当前未实装)→ identity 首绑
|
||||
(`tryBindIdentity` 失败 → SKIPPED,I3)→ Schd RESP/DNLD → `SnapshotFlow`(ACM2-16 定案:
|
||||
DNLD 与 RESP 均走 SnapshotFlow;RESP 成功后在同事务完成匹配开放 RQFD 的 `REQ_TRACK→DONE`;
|
||||
迟到或无匹配 RESP 严禁更新快照,直接转 SKIPPED 并审计)→
|
||||
其余 → `Handler.decide(redis.hgetAllFlightInfo(), msg)` →
|
||||
阶段 A:Redis 先写(I2 happens-before,TODO redisApply)→ 自有 PG 事务 2:
|
||||
MSG_EVENT 插入 + PROC_STATE→SUCCEEDED(同库原子,@Transactional);
|
||||
CMINMSGS 回填(DATE_PROCESSED/STATUS)为共享信箱**外部回填**(ACM2-19 定案):
|
||||
PG 事务提交后异步触发执行,持久化补偿、失败退避重试+告警、影子禁写;全部终态
|
||||
(SUCCEEDED / ignore SKIPPED / duplicate SKIPPED / DEAD)均必须回填,非终态禁止回填。
|
||||
5. 异常边界(U08):`processOne` 内 try/catch → `ProcFailure.fail(INFRA)`(attempts+1、
|
||||
退避、达上限 DEAD);`InterruptedException` 恢复中断位后**上抛**;loop 仅 catch
|
||||
`Exception` 作最后防线,`Error` 任其终止进程(异常必可见)。
|
||||
|
||||
幂等键(I3):`SNDR|TYPE|STYP|SEQN`(`Identity.of` 唯一入口);「含日边界」可配置且
|
||||
**默认关闭**(SEQN 重置作用域 CONFIRM 前,上线后不改幂等键)。
|
||||
|
||||
### 3.3 流程 3:投递(`Dispatcher`)
|
||||
|
||||
- 每 target 严格 FIFO(I1 双层同策略):`headUnsent` 取队头;队头退避未到期 → 等待不跳过。
|
||||
- 逐条循环显式排除 `KAFKA_SCHD`(D3/U06):schd 唯一出口 `flushSchd`——
|
||||
`claimBatch`(ORDER BY EVENT_ID)→ 按 FLID 分组取 max(EVENT_ID)(FIX:现役 buffer
|
||||
无去重会重发旧值)→ FLTR JSON 数组一次发出;失败整批 attempts+1 退避、队首未到期不
|
||||
claim、`lastFlush` 仅成功后推进;达上限整批 DEAD(DLQ)。周期/批上限取参数表
|
||||
(3s / 500)。
|
||||
- 轮询间隔取参数表(下限 50ms,N18);无 200ms 硬编码。
|
||||
- **Kafka 生产契约(ACM2-23)**:对接 Kafka 2.8+ / 3.x+,生产者强制 `acks=all`、`enable.idempotence=true` 与 `max.in.flight.requests.per.connection=1`;切流前须确认 Broker 支持 `InitProducerId(22)`,严禁非幂等降级;README 不提供生产降级 env。
|
||||
- 阶段 B(定案 2/D2,ACM2-12 缓做):ES 投递成功 → 同线程同步 `insertSync` 删除事件
|
||||
(`deleteOf`:Jackson 结构化序列化,refs 可空恒合法 JSON——U14)。当前不启用。
|
||||
|
||||
### 3.4 流程 4:日计划快照(`SnapshotFlow`,RESP/DNLD)
|
||||
业务身份统一由 `Identity.of` 生成:
|
||||
|
||||
```text
|
||||
SNDR | TYPE | STYP | SEQN
|
||||
```
|
||||
SCHD-RESP / SCHD-DNLD
|
||||
→ staging(流式解析+整包校验,TODO 阶段2;未实装→FAILED(UNSUPPORTED))
|
||||
→ 守卫判定:若为 SCHD-RESP,检查开放 RQFD(dttm < sentAt 或无匹配/已过期 → 严禁更新快照,转 SKIPPED 并审计)
|
||||
→ Redis Lua SNAPSHOT_REPLACE(同一 hash 原子「覆盖新代+按代差删」,删除集=旧代flids−新代)
|
||||
→ gen 版本推进(ACM2-12:gen 随 flightInfo 同在 Redis,Lua 内原子版本 CAS)
|
||||
→ 自有 PG 本地事务(原子性):PROC_STATE→SUCCEEDED + (RESP 匹配时)REQ_TRACK→DONE + MSG_EVENT 插入
|
||||
→ PG 提交后异步触发信箱回填(持久化补偿,ACM2-19)
|
||||
|
||||
数据流说明:SCHD-RESP 处理依赖 US-08 已登记的开放 REQ_TRACK;US-06 与 US-08 为单向数据流耦合(US-06 依赖 US-08 登记能力),不构成双向故事依赖(ACM2-24)。
|
||||
接收时只按信箱 ID 去重;解码后才首次绑定业务身份。重试保留原有绑定,不能把自己判为重复消息。身份被另一条记录占用时,当前消息转为 `SKIPPED`,记录 `duplicate-of:<id>`。是否加入日期边界取决于上游序号重置规则,默认关闭;上线后不能随意更换身份算法。
|
||||
|
||||
**已知缺口(U09,未定案,ACM2-12 后重设计为 Redis 内协议)**:gen 与 Lua/SUCCEEDED 不再
|
||||
分属两存储即可同原子(全部在 Redis Lua);真正跨存储的窗口收窄为「Lua 已完成、PG SUCCEEDED
|
||||
未写」——重放判据(版本不二次自增)与按代差删在 Lua 内以版本 CAS 承接,恢复协议待定案并补
|
||||
测试。现有代码的「CAS 重放二次自增」缺陷(版本 1→2)与实现注随协议重设计一并消除。
|
||||
Handler 是纯函数:
|
||||
|
||||
### 3.5 泵作业(`JobExecutor`;PUMP_JOB 自有 PG,作业窗口执行)
|
||||
```text
|
||||
Handler.decide(flightView, message) → Decision
|
||||
Decision = 航班变更 + msg 通知 + schd 状态 + 出站意图 + 静态数据变更
|
||||
```
|
||||
|
||||
- HISTORY_SWEEP(3:30 清场,I4 同步链):判史 → 同步写 ES → 仅删成功集。
|
||||
**占位门禁(U10/T07 修订)**:ES saveSync 接线前 `pickHistory` 恒空集、删除量恒 0,
|
||||
禁止「全量可删」fail-open 默认;现役五条判史规则 golden 通过后才允许接线。
|
||||
**阶段归属**:按 ACM2-12 阶段 B 缓做口径标为 DEFERRED,不作为阶段 A 切流门禁;启用前
|
||||
重新确认“历史链路不变”与阶段 B 投影范围、ES/OpenSearch 产品边界。
|
||||
- ARCHIVE(3:00):默认保留 **1 天**(接收时间早于 1 天)且仅终态(SUCCEEDED/SKIPPED/DEAD)可归档;保留期可配置为 1~7 天;
|
||||
**定案口径(ACM2-17)**:严禁向共享 MySQL 写入 `CMINMSGS_HST`(共享库严格保持两表 DML 契约);
|
||||
归档目标为自有 PG `PROC_STATE_HST`(及 `MSG_EVENT_HST`)。非终态(PENDING/FAILED)禁止归档。
|
||||
- PROJECTION_REBUILD(阶段 B 缓做,ACM2-12;重新评估后再启用)。
|
||||
Handler 不写 Redis、Kafka 或数据库。主泵负责应用决策;各类变更的持久化与重试边界必须明确,不能把“返回了 Decision”当成副作用已执行。
|
||||
|
||||
## 4. 失败与重试统一设计(U08)
|
||||
### 2.3 状态与错误分类
|
||||
|
||||
| 侧 | 组件 | 迁移语义 |
|
||||
|---|---|---|
|
||||
| ProcState(处理/快照) | `ProcFailure` + `FailureScheduler` | attempts+1 → exhausted ? DEAD(EXHAUSTED) : FAILED+nextAttemptAt |
|
||||
| MsgEvent 逐条(投递) | `Dispatcher.retryOrDead` | 同上(DLQ 保留行,attempts 审计) |
|
||||
| MsgEvent 批量(schd) | `flushSchd` 整批 | 队首未到期不 claim;整批退避;达上限整批 DEAD |
|
||||
```text
|
||||
处理:PENDING / FAILED → SUCCEEDED(成功)
|
||||
→ SKIPPED(忽略或重复)
|
||||
→ FAILED(等待重试)
|
||||
→ DEAD(非法报文或重试耗尽)
|
||||
|
||||
- 退避表 `[1s,2s,4s,8s,16s]`,单档封顶 `backoff-cap-ms=60s`;`attempt≤0` 兜底首档(N28)。
|
||||
- 时间一律经可注入 `java.time.Clock`(`TimeFactory`;测试用 MutableClock,无真实睡眠)。
|
||||
- loop 兜底 catch 不做状态迁移(迁移已在边界完成),仅防线程静默死亡。
|
||||
投递:PENDING → SENT
|
||||
→ PENDING(退避后重试)
|
||||
→ DEAD(重试耗尽)
|
||||
```
|
||||
|
||||
## 5. 不变量与实现落点
|
||||
`SUCCEEDED / SKIPPED / DEAD` 是处理终态,不再阻塞后续消息;`FAILED` 不是终态,仍占据队头。`DEAD` 表示需要处置,不等于业务成功。
|
||||
|
||||
| 不变量 | 语义 | 落点 | 状态 |
|
||||
|---|---|---|---|
|
||||
| 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 运行期保护未做) |
|
||||
| 错误类别 | 处理方式 |
|
||||
|---|---|
|
||||
| `MALFORMED` | 报文非法,直接 `DEAD`,不在原记录重放白名单内。 |
|
||||
| `CODEC_ERROR` | 解码能力问题,退避重试;修复后允许重放。 |
|
||||
| `UNSUPPORTED` | Handler 或快照能力未实现,按可恢复失败处理,不直接当作非法报文;仍受重试上限约束。 |
|
||||
| `INFRA` | 基础设施或执行异常,退避重试。 |
|
||||
| `EXHAUSTED` | 重试耗尽或滞留超时,转 `DEAD`,人工复核后允许重放。 |
|
||||
|
||||
## 6. 配置参数(`msgx.*`,ACMA-8 参数表初值)
|
||||
## 3. 收报与主泵
|
||||
|
||||
| 键 | 默认 | 说明 |
|
||||
|---|---|---|
|
||||
| `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) |
|
||||
### 3.1 收报
|
||||
|
||||
基础设施键位口径(Micronaut 5.1,U03):`datasources.default.*`(**自有 PostgreSQL**,
|
||||
ACM2-12)、`flyway.datasources.default.*`、`mailbox.shared-mysql.*`(共享信箱,仅 DML)、
|
||||
`kafka.producers.default.*`、`eureka.client.*`;logback 独立于本文件,
|
||||
环境变量前缀 `MSGX_LOGSTASH_*`。
|
||||
`InboxPoller` 默认每秒读取未处理信箱记录,在 PG 建立 `PENDING`,不解析业务载荷。PG 插入必须按信箱 ID 幂等,失败由后续扫描补建。
|
||||
|
||||
## 7. 测试策略
|
||||
水位优化分为两条路径:快路径读取水位之后的新记录,补偿路径重扫遗漏的未处理记录。只有本批 PG 入队全部确认后才能推进水位。**水位不是已处理标记,也不能单独证明较小 ID 已收齐**;迟提交和补扫场景的顺序保证需要在启用前验证。
|
||||
|
||||
- **接口驱动 + 假仓储**:管道语义全部离线单测(无 DB/Redis/Kafka),时间用 `MutableClock`。
|
||||
- **不变量测试**:FIFO/HOL、schd 批退避与 DLQ、重试上限、重放白名单、
|
||||
聚合最新态、配置绑定、DI 装配冒烟(PipelineSmokeTest:收报→FAILED(CODEC_ERROR,stub codec
|
||||
恒 CODEC_ERROR)→重放→schd 聚合发出)。
|
||||
**范围注**:以上覆盖的是边界级(processOne/flushSchd/仓储)语义;主泵 tick 级 HOL/毒丸/退避
|
||||
门禁因 Pump 未注入 Clock 而无单测锁定(§3.2 注),DispatcherTickTest 仅锁批退避与「队首未到期
|
||||
不推进」。
|
||||
- 最近生成的 JUnit XML 报告为 39 项测试全绿;本轮文档审查因沙箱不能写用户级 Gradle 缓存,
|
||||
未重新执行(wrapper 钉 9.6.1 + JDK 25;受限环境需将 `GRADLE_USER_HOME`/`TMPDIR` 指向可写目录)。
|
||||
- **U29 门禁(规划)**:不变量清单化入 CI 红即阻塞;「新增逻辑必伴生不变量测试」入贡献约定。
|
||||
兼容 HTTP 入口执行“写入共享信箱 → PG 入队”。两步不在同一事务中:信箱成功而 PG 失败时,原文不能丢失,由轮询补建;客户端失败重试可能再次写信箱,业务身份去重仍然必需。
|
||||
|
||||
## 8. 可观测性设计(U12 已落地部分)
|
||||
### 3.2 主泵调度
|
||||
|
||||
- logstash TCP 经 AsyncAppender(queueSize 4096 / neverBlock / discardingThreshold 0):
|
||||
logstash 不可达丢弃日志而非阻塞业务线程(N31b 顺序约束:先降级再补日志)。
|
||||
- MDC `traceId` = cminmsgsId/eventId(`TraceLog.withTrace`),覆盖处理/投递日志片段;
|
||||
收报→投递全链贯穿与 micrometer gauge(队列深度/投递延迟)未实装。
|
||||
- 生命周期结构化日志:收报 INFO、SUCCEEDED/SKIPPED INFO、FAILED WARN、DEAD/毒丸 ERROR、
|
||||
flush 批次 INFO、整批 DLQ ERROR——告警暂以 ERROR 日志为落点(DEAD 告警出口属 U13/U25)。
|
||||
每次 `Pump.tick`:
|
||||
|
||||
## 9. 已知缺口与定案待办(对照 ACM2-10)
|
||||
1. 读取最小未完成消息,必须包含 `FAILED`,不能只查当前可执行的记录。
|
||||
2. 无消息,或队头仍在退避窗口内时,允许执行一个维护作业;消息已可执行时优先处理消息。
|
||||
3. 队头达到重试或滞留上限时转 `DEAD(EXHAUSTED)`;未到重试时间则等待,不领取后续消息。
|
||||
4. 其余情况调用 `MessageProcessor.processOne`。
|
||||
|
||||
| 项 | 缺口 | 计划 |
|
||||
|---|---|---|
|
||||
| U05 | JDBC PG 仓储、CMINMSGS 适配器与 **InboxPoller** 已有初版;但实现仍是逐操作独立连接,`MSG_EVENT + SUCCEEDED` 无本地事务,回填仍同步夹在两者之间;OutboxMailbox、可靠补偿、ARCHIVE 目标与真实双库集成测试未完成。`JdbcRefDataRepository` 仍为进程内过渡态,`JdbcFlightStateRepository` 为阶段 B 空实现 | 完成本地事务边界、回填/入队补偿、COUTMSGS 适配器与 Testcontainers 双库验收;生产启用前 fail-fast |
|
||||
| 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) | WP2(legacy 逐字对拍固化) |
|
||||
| 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 批次 |
|
||||
| U19–U21 | 请求状态机量纲/死分支、identity 绑定静默跳过、21 类静态接口/表对齐(REF_MASTER/SOURCE) | WP2 |
|
||||
| U22–U24 | 载荷类型收敛、eventSeq 未接线、每报文全量读语义定案 | WP3 |
|
||||
| U25/U28/U30 | 一致性哨兵实装、README 安全节/入口、索引与杂项 | WP3/4 |
|
||||
| U27 | 专有材料(SIS md 703KB / XSD 版权头)治理决策 | WP4(ACL 核验先行) |
|
||||
| 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` 第二库规划 | 仅作决策脉络参考 |
|
||||
作业执行时长和饥饿边界需要限制,不能用长期作业阻塞已到期消息。滞留超时应基于稳定的起始时刻,不能用每次失败都会刷新的 `updatedAt` 代替。
|
||||
|
||||
## 10. 用户故事
|
||||
### 3.3 单条处理
|
||||
|
||||
面向需求优化的用户故事已集中到 [user-stories.md](user-stories.md)。该文档将目标能力、验收标准、
|
||||
依赖与待确认问题分开,覆盖阶段 A US-01~US-14、延后 US-15、上线 EPIC 及 legacy HTTP 去留,
|
||||
避免把当前实现、目标设计和遗留兼容行为混成同一项承诺。
|
||||
```text
|
||||
读取原文 → 解码 → 忽略规则 → 首次绑定身份
|
||||
├─ RESP / DNLD:快照流程
|
||||
└─ 其他:Handler 决策
|
||||
↓
|
||||
Redis 应用变更
|
||||
↓
|
||||
PG 本地事务:待发事件 + 处理终态
|
||||
↓
|
||||
提交后补偿回填信箱
|
||||
```
|
||||
|
||||
- 原文缺失当前归为 `MALFORMED`;读取异常不能伪装成“缺失”,应进入基础设施重试。
|
||||
- 忽略规则覆盖约定的 `LDM / REGN / RSTA / EROR`,转 `SKIPPED` 并审计;不能产生业务副作用。
|
||||
- Redis 更新必须先于对应事件提交。PG 事务只覆盖本库,不能靠事务注解把 Redis 或 MySQL 操作变成原子操作。
|
||||
- 所有终态都需要回填信箱,包括成功、忽略、重复和死信;非终态禁止回填。回填必须在 PG 提交后执行,并有持久化补偿、退避与告警;影子环境禁写。
|
||||
|
||||
## 4. 日计划快照与请求匹配
|
||||
|
||||
### 4.1 快照发布
|
||||
|
||||
`SCHD-RESP` 和 `SCHD-DNLD` 都进入 `SnapshotFlow`:
|
||||
|
||||
1. **暂存校验**:流式解析后完成整包校验和航班规范化;失败前不修改权威状态。暂存数据可在崩溃后从原文重建。
|
||||
2. **应答守卫**:`RESP` 必须匹配开放的 `RQFD` 请求;无匹配、已过期或报文时间早于发送时间时,转 `SKIPPED` 并审计,不更新快照。
|
||||
3. **原子替换**:Redis Lua 在同一次操作中校验版本、覆盖新代、删除旧代差集并推进 `gen`。删除集为“旧代航班集合 − 新代航班集合”,不是全部现存航班,不能误删快照集合外的增量航班。
|
||||
4. **提交结果**:在 PG 同一事务中保存 `MSG_EVENT`、将消息置为 `SUCCEEDED`,并将匹配 `RESP` 的请求置为 `DONE`;随后执行信箱回填。
|
||||
|
||||
关键恢复窗口是“Redis 已替换,PG 尚未提交”。重试同一快照必须识别已应用结果,不能再次增加版本,也不能用旧快照覆盖新状态。**版本校验与重放识别协议仍待实现验证**,单纯“读当前版本再加一”不满足要求。
|
||||
|
||||
### 4.2 上游请求与静态数据
|
||||
|
||||
`RequestCoordinator` 管理请求生命周期:
|
||||
|
||||
```text
|
||||
REGISTERED → SENT → WAITING → DONE
|
||||
└──→ EXPIRED
|
||||
```
|
||||
|
||||
注册同类新请求前使旧开放请求过期。只有 `COUTMSGS` 写入确认后才标记 `SENT` 并关联出站记录;写信箱成功但本地未确认的情况需要补偿与去重,不能无条件重新发送。
|
||||
|
||||
应答优先按已确认的回显字段精确匹配。回显契约未确认时,按同类开放请求和 `DTTM ≥ sentAt` 判断的降级方式存在跨代误配风险,必须明确接受并审计,不能宣称精确关联。比较前统一时区和时间单位。
|
||||
|
||||
静态应答写入自有 PG 的 `REF_MASTER`,日计划应答走快照流程。请求完成必须在相应数据处理成功之后;超时和迟到应答不能修改已关闭请求对应的状态。
|
||||
|
||||
## 5. 事件投递
|
||||
|
||||
### 5.1 普通事件
|
||||
|
||||
`Dispatcher` 按 `TARGET` 读取最小未发送 `EVENT_ID`。队头退避未到期时,该目标停止推进;发送确认后才标记 `SENT`,失败记录次数和下次执行时间。所有外部调用需要有界超时,避免阻塞整个投递线程。
|
||||
|
||||
投递是至少一次:下游已接收但本地未标记成功时可能重发。Kafka 生产约束沿用架构决策 D11,但生产者幂等不替代应用层事件去重;跨重启的事件身份和下游去重契约仍需落实。共享出站信箱也必须单独解决重复写入,不能假设 Kafka 的保证适用于 MySQL。
|
||||
|
||||
### 5.2 `schd` 聚合
|
||||
|
||||
`KAFKA_SCHD` 不进入逐条投递循环,只由 `flushSchd` 发送:
|
||||
|
||||
1. 队头可执行后,按 `EVENT_ID` 顺序领取有界批次。
|
||||
2. 按 `FLID` 分组,保留批次内最大 `EVENT_ID` 对应的状态,组成 FLTR JSON 数组发送。
|
||||
3. 成功后将本批被代表的事件一起标记完成,推进 `lastFlush`;失败则整批增加次数并退避,达到上限整批转 `DEAD`。
|
||||
|
||||
默认聚合周期 3 秒、批上限 500。它提供最新状态通知,不保留每次中间变化;批次不能绕过尚在退避的队头。
|
||||
|
||||
## 6. 失败恢复与维护作业
|
||||
|
||||
### 6.1 失败、重试与重放
|
||||
|
||||
`ProcFailure` 和 `FailureScheduler` 统一处理侧的失败落账;投递侧按单条或聚合批次执行同样的次数与退避规则。默认最多 5 次,退避档位为 1、2、4、8、16 秒,单档封顶 60 秒。
|
||||
|
||||
失败必须在持有具体消息或批次的位置记录,外层循环只做兜底日志和等待,不重复增加次数。线程中断应恢复中断标记并向上传递;不捕获 JVM `Error` 作为普通业务失败。所有时间判断通过注入的 `Clock` 完成。
|
||||
|
||||
`ReplayService` 只允许 `CODEC_ERROR / UNSUPPORTED / INFRA / EXHAUSTED` 从 `FAILED / DEAD` 回到 `PENDING`,重置次数和下次执行时间,保留身份与错误审计。旧消息进入终态后,后续消息可能已经执行;因此**重新入队不等于恢复历史顺序**,人工重放前必须评估状态覆盖和版本保护,不能直接批量重放到生产。
|
||||
|
||||
### 6.2 归档与阶段 B 作业
|
||||
|
||||
`ARCHIVE` 仅归档自有库的终态记录,默认按接收时间保留 1 天,保留期可配置为 1~7 天。归档表、接收时间依据和关联事件处理尚需落地;不得归档未完成记录,也不能因移走身份记录而意外失去业务去重能力。共享信箱保留策略由库所有方管理。
|
||||
|
||||
`HISTORY_SWEEP` 与 `PROJECTION_REBUILD` 暂缓。未来清场必须只删除已确认成功写入历史存储的集合;历史写入未接通时默认删除零条。历史写入与删除事件入队之间仍需恢复方案,顺序调用不构成原子提交。
|
||||
|
||||
## 7. 接口与运行配置
|
||||
|
||||
兼容入口为 `POST /cminmsgs/send`,请求体为原始 XML,当前成功响应为 HTTP 200、`text/plain` 格式的信箱记录 ID。这只表示接收结果,不表示业务处理成功。请求媒体类型、字符集和失败响应仍需与现役逐项对拍;查询与其他兼容端点不能因列入需求就视为已提供。
|
||||
|
||||
运行配置以 `application.yml`、`application-dev.yml` 和 `.env.example` 为准,设计上重点区分:
|
||||
|
||||
- `pipeline.autostart` 与 `msgx.stubs`:分别控制管道启动和内存适配器;生产禁止 stub,默认不自动启动。
|
||||
- `pipeline.poll-interval / max-attempts / head-deadline`:控制轮询、重试上限和队头滞留;不能改变 FIFO。
|
||||
- `schd.flush-period / flush-limit`:控制状态通知的聚合延迟与批量大小。
|
||||
- `identity.include-day-boundary`:影响去重语义,不能作为普通调优项切换。
|
||||
- `phase`:阶段 A 是当前范围,不应把切为 B 当成已具备历史投影能力。
|
||||
|
||||
日志关联消息 ID、事件 ID 和批次;失败记录错误分类、次数、下次执行时间。健康检查反映依赖实际可用性;队头滞留、积压、死信和补偿失败需要指标及告警。日志出口故障不得阻塞业务线程。
|
||||
|
||||
## 8. 验证要求
|
||||
|
||||
单元测试使用内存仓储和可推进的 `Clock`,不依赖睡眠或在线中间件。以下不变量必须有回归测试,接口级单测不能替代主泵调度测试:
|
||||
|
||||
| 场景 | 必须验证的结果 |
|
||||
|---|---|
|
||||
| 重复扫描、入队中断、较小 ID 迟到 | 不重复入队、不丢记录、不让后续消息越序。 |
|
||||
| 队头失败、退避及作业竞争 | 消息不越队;到期后恢复;作业不使消息无限饥饿。 |
|
||||
| 同身份多条记录、失败后重试、归档后重复 | 只产生一次有效业务处理,不把自身重试判为重复。 |
|
||||
| Redis 成功后 PG 失败、快照重复或迟到 | 不重复推进版本、不回退状态、不误删增量航班。 |
|
||||
| PG 提交失败、信箱回填失败 | 事件与处理结果一起回滚;已提交结果只补偿回填。 |
|
||||
| 投递确认丢失、批次失败、次数耗尽 | 允许可识别的重发、保持目标顺序、整批退避并保留死信。 |
|
||||
| 请求超时、无匹配 RESP、时间单位不一致 | 不误用迟到应答,不提前完成请求。 |
|
||||
| stub 误配置、重复实例、停机中断 | 生产拒绝不安全启动,工作线程能正确退出。 |
|
||||
|
||||
真实适配器还需 PG/MySQL 事务与补偿集成测试、Redis Lua 中断恢复测试、Kafka 故障投递验证;启动冒烟只证明装配可用,不证明这些一致性要求已满足。常规验证命令为 JDK 25 下执行 `./gradlew test`。
|
||||
|
||||
## 9. 实现入口
|
||||
|
||||
生产代码根目录为 `src/main/kotlin/com/gzzn/omms/msgexchange/`:
|
||||
|
||||
| 关注点 | 主要入口 |
|
||||
|---|---|
|
||||
| 收报与兼容接口 | `ingress/InboxPoller.kt`、`InboxService.kt`、`InboxController.kt` |
|
||||
| 调度与处理 | `processing/Pump.kt`(含 `MessageProcessor`)、`Handler.kt`、`Identity.kt` |
|
||||
| 快照与请求 | `processing/SnapshotFlow.kt`、`reference/RequestCoordinator.kt` |
|
||||
| 投递与作业 | `delivery/Dispatcher.kt`、`SchdAggregation.kt`、`jobs/JobExecutor.kt` |
|
||||
| 持久化与恢复 | `infra/persistence/`、`infra/redis/`、`infra/retry/` |
|
||||
| 启停与配置 | `PipelineLifecycle.kt`、`config/PipelineProps.kt` |
|
||||
|
||||
## 10. 当前实现差异
|
||||
|
||||
以下缺口直接影响上述设计是否成立,不能以类或接口已存在作为完成依据:
|
||||
|
||||
- **事务与外部副作用**:当前 JDBC 操作仍分散执行;普通处理按“插事件 → 同步回填信箱 → 更新状态”调用,尚未实现要求的 PG 原子提交与提交后补偿。Redis 普通变更、出站意图及真实出站信箱也未形成完整链路。
|
||||
- **收报与调度**:轮询仍从 `afterId=0` 扫描,持久水位与补扫策略未完成;作业以“队头为 FAILED”近似窗口,未区分是否已到期。主泵直取系统时间,滞留判据使用更新时刻,不能保证设计要求的超时升级。
|
||||
- **快照与业务能力**:快照暂存仍为占位,当前分流仅覆盖 DNLD,RESP 守卫与请求完成事务未接通;`gen` 仍走过渡仓储,未实现 Redis 内原子版本与重放协议。忽略规则及所需 Handler 还需补齐。
|
||||
- **请求、静态数据与归档**:请求在真实出站前就标记发送,时间匹配与审计仍需修正;静态数据存在进程内过渡实现,归档表与关联保留策略尚未落地。
|
||||
- **生产与运维**:缺少完整启动校验、运行期单写者保护、影子隔离和 Kafka 强制配置校验;投递滞留升级、死信告警、端到端追踪与指标未闭环。
|
||||
|
||||
进度与验收项见 Plane ACM2-10 实施计划(U01–U30),业务契约与待确认事项见 [user-stories.md](user-stories.md)。本文件不维护工单流水账、测试数量或历史方案全文。
|
||||
|
||||
Reference in New Issue
Block a user