Files
msgexchange-v2/docs/implementation.md
T

467 lines
41 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.
# 实现设计
本文件是实现设计的唯一出处,分两章:
- **处理管道**章:记录模型、状态机、收报与水位、主泵与事务边界、回填、快照与请求、投递、失败恢复与维护作业;
- **航班域**章:航班当前态的权威模型、合并与写入语义、删除与重建。
正文描述**目标设计**,不标注交付状态:可声明性见 [specification.md](specification.md)「声明边界」与「当前已知偏差」,进度在 Plane(ACM2)。契约与不变量只引稳定 ID;参数取值只引 `PARAM:<完整键>`(见 [reference.md](reference.md))。
## 1. 术语与持久化记录
### 1.1 管道术语
| 术语 | 语义 |
|---|---|
| `W`(水位) | 信箱 ID 的连续上界:`(min, W]` 已全部读入自有 PG;只随新 ID 成功入队推进(永久空洞放行是唯一例外)。 |
| `holeSince` | `W+1` 处空洞最早被观测到的时刻;无空洞时为 NULL,跨重启保留。 |
| 队头 | 最小的未完成消息(`PENDING``FAILED` 都占位)。 |
| 终态 | `SUCCEEDED` / `SKIPPED` / `DEAD`;到达后队列方可推进。 |
| 回填意图 | 「还欠一次信箱标记」的持久化事实,与终态同一条语句落库。 |
| `R``R_keep`、处理标记 | 定义见 [specification.md](specification.md)。 |
### 1.2 持久化记录
| 记录 | 用途 | 关键约束 |
|---|---|---|
| `PROC_STATE` | 入站消息的处理状态、身份、尝试次数、错误原因与回填事实 | `MSG_ID = CMINMSGS_ID` 主键防重复入队;`IDENTITY_KEY` 唯一约束防业务重复;按最小未完成 `MSG_ID` 取队头;`BACKFILL_NEXT_AT` 非空 = 还欠一次回填,`BACKFILL_AT` 非空 = 标记已确认,`BACKFILL_ABANDONED_AT/REASON` 非空 = 已停止自动重试(**不等于**标记已确认);`RECEIVED_AT` 复制自信箱接收时间、**可能为 NULL**、仅用于对账与展示;`ENQUEUED_AT` 是本地入队时间、非空、是超期判据的唯一依据;归档后的去重影子行置 `STATE='ARCHIVED'`、只保留 `IDENTITY_KEY``MSG_ID`,不占队头、不触发回填、不参与积压聚合(`G-PROC-HST`)。 |
| `MSG_EVENT` | 等待投递的事件(outbox) | `EVENT_ID``KAFKA:msg` 是稳定事件身份并决定投递顺序;对 `KAFKA:schd` 是每次接受 upsert 时替换的写代次。`TARGET` 区分 `KAFKA:msg` / `KAFKA:schd``PARTITION_KEY` 当前取 `FLID``Q4` 定案前为假定,见 `C-29`);`EVENT_TYPE` 区分 UPSERT 与 TOMBSTONE。`KAFKA:schd``FLID` 单行 upsert,只保留最新 `STATE_VERSION``SENT_AT` 在投递确认的同一条 UPDATE 内写入,是保留期判定的唯一基准。 |
| `REQ_TRACK` | 上游请求及应答关联 | 状态 `PENDING / SENT / DONE / EXPIRED`;保存请求类型、覆盖运营日、发送方、出站信箱 ID 与发送/完成时间;**「同类只允许一个开放请求」的唯一键 = `(请求类型, 覆盖运营日, 发送方)`,且仅对开放状态生效**。 |
| `REF_MASTER` | 静态参考数据(目标表) | `(RTYPE, RKEY)` 唯一;客户端与刷新流程见 [requirements.md](requirements.md) `US-13`/`US-14`。 |
| `FLIGHT_SCHD` | 航班标量及单值异常字段 | `FLID` 主键;`OPERATION_DAY` 一经确定不可变;版本与最近消息 ID 用于追踪。变长集合存于资源明细表与 `FLIGHT_ROUTE_POINT`,规则见「航班域」。 |
| `INBOX_CURSOR` | 消费水位 `W`、空洞计时 `holeSince`、播种事实 `SEEDED_AT` | 单行游标;`W` 只随新 ID 成功入队推进,遇空洞即停;`HOLE_SINCE` 持久化空洞观测时刻,进程重启不丢计时。`SEEDED_AT IS NULL` **不等于**从未消费(已有库新增列后同样为 NULL)。 |
| `SCHD_SNAP_LOG` | 日计划处理留痕 | 只追加、可重建,不参与状态决策;保留期见 [reference.md](reference.md)。 |
| `PROC_STATE_HST` | 终态处理记录的归档目标 | 只归档到自有 PG,不落共享库历史表(`G-PROC-HST``G-HST-RETENTION`)。 |
字段与索引以 `src/main/resources/db/migration/` 的迁移链为准(Oracle 11g 目录为占位,未接入 Flyway)。报文原文仍从共享信箱读取,原文保留期必须覆盖处理与重放窗口(`C-7`);清除前提、保留期下界与处理标记值集见 `C-5``C-12`
## 2. 消息、身份与决策
`XmlCodec`(实装 `JacksonXmlCodec`)把 XML 解码为 `DecodedMessage`,包含 `SNDR / TYPE / STYP / SEQN / DTTM` 元数据、`MsgKind` 与业务载荷。解码失败区分 `MALFORMED`(报文非法,不重试)与可随 codec 修复的 `CODEC_ERROR``MsgKind` 是一等分派键:`Schd(RESP/DNLD/ADFT)``Flop``Fdel``Unsupported`
业务身份统一由 `Identity.of` 生成:`SNDR | TYPE | STYP | SEQN`。接收时只按信箱 ID 去重,解码后才首次绑定业务身份;重试保留原有绑定,因此自身重试不会被判为重复。身份被另一条记录占用时,当前消息转 `SKIPPED`,记录 `duplicate-of:<id>`。是否加入日期边界取决于上游 `SEQN` 重置周期(见 `C-20`/`Q11``PARAM:msgx.identity.include-day-boundary`);上线后不能随意更换身份算法。
**身份绑定是独立的幂等单语句**`WHERE IDENTITY_KEY IS NULL`),不参与业务事务。它的前提是「报文不可变」(`PRE-7`):同一身份的重发不会被比对内容,若上游改发正文会被判为重复并跳过(`Q15`)。
分派与落库由 `MessageProcessor` 协调:按 `MsgKind` 把已绑定身份的队头消息交给对应事务协调器(DNLD/RESP → `ScheduleProcessor`ADFT → `AdftProcessor`FLOP → `FlopProcessor`FDEL → `FdelProcessor`,其余 → `FAILED(UNSUPPORTED)`)。这些处理器在 `PIPELINE_LOCK` 事务内读取当前完整态,调用纯领域决策逻辑得到下一完整态与待发事件,再统一落库并登记回填意图;它们不直接触碰 Kafka。领域决策逻辑不执行 I/O。
忽略规则:解码后、分派前按大小写不敏感的 `TYPE-STYP` / `TYPE-*` 匹配忽略清单(基线 `LDM` / `REGN` / `RSTA` / `EROR`,不混用 `ERROR`),命中转 `SKIPPED` 并记录 `ignored:<rule>`;忽略报文照常绑定身份,但不更新航班、不创建业务通知(`US-04`)。
## 3. 状态与错误分类
```text
处理:PENDING / FAILED → SUCCEEDED(成功)
→ SKIPPED(业务重复、忽略、无匹配)
→ FAILED(等待退避重试)
→ DEADMALFORMED / PROTOCOL / EXHAUSTED,均需人工处置)
投递:PENDING → SENT
→ PENDING(退避后重试)
→ DEAD(重试耗尽,记录保留作 DLQ)
```
`SUCCEEDED / SKIPPED / DEAD` 是处理终态,不再阻塞后续消息;`FAILED` 不是终态,仍占据队头。`DEAD` 表示需要处置,不等于业务成功。`ARCHIVED` 不是处理态:它是终态行归档后留在主表的去重影子行,不占队头、不触发回填、不参与积压聚合(见「生命周期与清除」)。
重试次数用尽时统一转 `DEAD(EXHAUSTED)``ERROR_CLASS` 被覆写为 `EXHAUSTED`,原始错误类别不再保留(`LAST_ERROR` 保留原因文本)。重放白名单与错误分类表见 [reference.md](reference.md)「错误分类与重放白名单」。
## 4. 收报与水位
### 4.1 收报流程
`InboxPoller` 按 ID 升序、有限批次读取水位之后的信箱记录(`ID > W`,**不以处理标记为谓词**),在自有 PG 建立 `PENDING` 并把水位推进到连续上界。每轮:
1. 读取游标 `(W, holeSince)`。信箱不可读时记日志、等下一轮,**不动水位**——这是基础设施失败,不能当成「没有新消息」。
2.`ID > W` 的升序前 `PARAM:msgx.pipeline.claim-batch` 行。
3.`W+1` 起逐 1 数,求连续上界;遇到第一个缺号即停止计数。
4. 空洞判定(仅当本批存在「缺号之后的行」时才可能成立):
- 缺号首次被观测到 → `holeSince = now`
- `now holeSince < PARAM:msgx.pipeline.max-commit-delay` → 水位停在缺号前,**本批缺号之后的行一律不入队**(否则晚提交的较小 ID 会排到它们后面,破坏 FIFO);
- `now holeSince ≥ PARAM:msgx.pipeline.max-commit-delay` → 判定为永久空洞,水位放行到「缺号后第一行 − 1」并清空 `holeSince`。放行**只跳过空洞本身,不越过任何已存在的行**。
5. 在同一个 PG 事务内:对水位以内的每一行 `insertIfAbsent(MSG_ID, RECEIVED_AT, ENQUEUED_AT)` 并写回 `(W, holeSince)`;主键冲突表示已入队(重复扫描与兼容入口并发都安全),不计入、不报错。
6. 提交。本批因空洞或批次上限未入队的行留待下一轮——**每轮最多解决一个空洞**。
`holeSince` 落在 `INBOX_CURSOR.HOLE_SINCE`,进程重启不丢计时。旧空洞补齐后出现的新空洞从新观测时刻重新计时,不继承旧等待时间。
**代价(必须接受并观测)**:水位遇空洞即停意味着空洞之后的所有消息最多要等一个老化窗口才能入队;自增回滚等会在 ID 序列留下永久空位,每出现一个永久空位就是一次等长的入队停摆,空位频繁时有效吞吐按比例下降。运行期必须观测永久空洞计数与水位滞后(指标见 [reference.md](reference.md))。
### 4.2 发现完整性依赖
水位的有效性依赖 `PRE-2``PRE-3``C-1`/`C-2`/`C-3`)。承诺缺失时 `W` 只是快路径提示,不足以证明该区间收齐。ID 分配 → 事务可见时延上界必须由库方直接给出,**不可由 SIS 报文 `Expiry` 推导**。
日常扫描只有一条快路径:`ID > W ORDER BY ID ASC LIMIT claim-batch`
### 4.3 切流播种
若信箱已有存量(典型情况是最老分区已被清除、`MIN(ID)` 远大于 1),从 `W=0` 启动会先把 `ID=1` 判成空洞、白等一个老化窗口。是否跳过存量属于**切流决策**,因此不做默认选择:只有显式配置 `PARAM:msgx.pipeline.cutover-watermark` 才播种,取值 `min`(读现存全部)/ `zero`(从 0 按空洞规则)/ `max`(跳过当前可见存量)/ 具体 ID。升级实例(已有水位或已有处理记录)拒绝重新播种;播种事实与水位同语句落库(`SEEDED_AT`)。
### 4.4 兼容 HTTP 入口
`POST /cminmsgs/send` 执行「写入共享信箱 → PG 入队」,两步不在同一事务:信箱成功而 PG 失败时原文仍在信箱中,由轮询补建;客户端失败重试可能再次写信箱,业务身份去重仍然必需。该入口直接写 `PROC_STATE`、不读不推水位,登记的行因此**超出水位**;主泵只领 `MSG_ID ≤ W`,所以在水位追平前不会被处理——顺序不受影响,代价是延迟到追平,且必须让收报轮询运行。响应语义见 `C-28`
### 4.5 单实例与水位
信箱读取不加锁,水位是单行覆盖写,本设计只在单活动实例下成立(`PRE-5`)。多实例并发收报会让水位互相覆盖(覆盖回退只会造成重复扫描,不会丢消息,但空洞计时会失真),必须先有实例级排他。
## 5. 主泵调度与单条处理
### 5.1 调度
每次 `Pump.tick` 只围绕最小未完成消息:
1. 无队头:按轮询间隔休眠。
2. 队头超出水位(`msgId > W`):**不领取**,休眠到下一轮。这类行只可能来自兼容入口的直接登记;允许领取会让它越过尚未入队的较小 ID。会打印一条限流 WARN(仅在水位值变化时打一次)。
3. 队头为 `FAILED` 且尝试次数达上限:转 `DEAD(EXHAUSTED)`,终态与回填意图同一条 UPDATE 落库,**不做跨库写**。该分支只写 `PROC_STATE`,不取 `PIPELINE_LOCK`、不在处理器事务内,也不在 `MessageLifecycleGate` 内。
4. 队头为 `FAILED` 且未到 `next_attempt_at`:休眠到可重试时刻,不处理后续消息。
5. 其余(新消息或退避到期的重试):调用 `MessageProcessor.processOne`,失败迁移在该边界内完成。
**终态判据只有尝试上限,没有按时间的毒丸**:处理卡死应由外部调用的有界超时兜底;用一个时间阈值把消息直接推入 `DEAD` 会绕过人工复核,并制造与人工重放并发的旁路写入者。队头年龄由 `msgx.pipeline.backlog.oldest_unprocessed_seconds` 观测,主泵不做这项判定。
**所有取时统一经注入 `Clock`**(收报空洞老化、主泵调度、处理器落库时间、回填重试、作业切日),不使用系统时钟。
### 5.2 processOne
```text
processOne(head)
1. 读原文:缺失 → DEAD(MALFORMED, raw-missing);读取异常 → FAILED(INFRA)
2. 解码:
报文非法(MALFORMED) → DEAD(MALFORMED),不重试
可修复解码错(CODEC_ERROR) → FAILED(CODEC_ERROR) 退避
3. 身份绑定(仅当 IDENTITY_KEY 为空):
已被本消息占用 → 继续
已被别的消息占用 → SKIPPED(duplicate-of:<id>),结束
空闲 → 写入 IDENTITY_KEY(独立单语句,不参与业务事务)
4. 忽略规则命中 → SKIPPED(ignored:<rule>)(非业务型终态)
5. 按 MsgKind 分派:
SCHD-DNLD / SCHD-RESP → ScheduleProcessor(快照事务)
SCHD-ADFT → AdftProcessor(单航班事务)
FLOP / FDEL → Flop / FdelProcessor(单航班事务)
Unsupported → FAILED(UNSUPPORTED) 退避
载荷缺失 → DEAD(MALFORMED)
整包协议拒绝 → DEAD(PROTOCOL),不落半包
6. 业务型成功:处理器在自己的事务内写航班变更 + 待发事件 + SUCCEEDED + 回填意图
7. 结束:主泵不做回填;回填意图已随终态落库,由扫描补写信箱标记
```
- 原文缺失归为 `MALFORMED`;读取异常按基础设施失败进入重试,与「原文缺失」区分。
- 类型未覆盖不等于报文非法:忽略类报文在规则实现前不按 `MALFORMED` 处理。
- **主泵不回填**:终态与回填意图由同一条 UPDATE 落库,回填一律由扫描驱动,不占用 FIFO 关键路径。回填只需消息 ID,缺 META 或解码失败的死信同样可补写。影子环境禁写。
### 5.3 事务边界
| 动作 | 显式事务 | 持 `PIPELINE_LOCK` | 触及航班表 / 事件 | 原子性来源 |
|---|---|---|---|---|
| 收报入队(`insertIfAbsent` + `cursor.save`) | 是 | 否 | 否 | 同库事务 |
| 身份首次绑定 | 否 | 否 | 否 | 单语句 + 唯一约束 |
| 业务型终态(处理器产出 `SUCCEEDED`) | 是 | 是 | 是 | 同库事务:航班变更 + 事件 + 终态 + 回填意图 |
| 非业务型终态(`MALFORMED` / `PROTOCOL` / `SKIPPED` / `EXHAUSTED`) | 否 | 否 | 否 | 单语句(终态与回填意图同一条 UPDATE) |
| 航班历史清理的物理删除 | 是 | 是(`INV-18`) | 是 | 同库事务:复查判据 + 归档成功后删除 |
| 回填(信箱标记 + `BACKFILL_AT`) | 否 | 否 | 否 | 跨库两次单写;幂等可重跑 |
| 人工重放(批量改回 `PENDING`) | 否 | 否 | 否 | 单语句批量;`MessageLifecycleGate` 与回填互斥 |
结论:「航班变更与处理终态同事务」只对业务型终态成立。`PIPELINE_LOCK` 的竞争写者是**航班历史清理**`INV-18`),不是别的处理器线程;没有第二写者时该锁不产生额外串行度。
### 5.4 历史积压
信箱中的成规模存量(上线前遗留、停机累积)**不是特殊模式**:它逐条走与日常完全相同的 FIFO 路径。
- 顺序由 `MSG_ID` 决定,不由执行方式决定。入队与处理由不同线程驱动、可以并发,「先入队后处理」只是可选的运维规程,不是正确性前提;系统不提供「只入队」模式。
- 不加速、不分流、不走旁路:不允许并行队头,也不允许实时消息跳过积压。
- 尝试上限与退避对积压同样生效,不因积压而放宽。
- 经确认不再处理的行置 `SKIPPED` 并记录原因,到达终态后走回填通道;不存在「整段 DELETE」的快速通道(授权与留痕见 `C-27``Q12`)。
- 消化期间的可观测项与完成时限口径见 [reference.md](reference.md) 与 `CLM-9`**扫描周期不是完成时限**。
## 6. 回填
### 6.1 事实与扫描谓词
终态落库时登记回填意图;标记回写由扫描驱动,跨库单写、幂等可重跑。扫描谓词(与实现一一对应):
```text
STATE ∈ {SUCCEEDED, SKIPPED, DEAD} -- 终态
AND BACKFILL_AT IS NULL -- 标记尚未确认
AND BACKFILL_ABANDONED_AT IS NULL -- 未放弃(放弃行可人工恢复)
AND ( BACKFILL_NEXT_AT IS NULL -- 异常兜底:终态行没有退避时间
OR BACKFILL_NEXT_AT ≤ NOW -- 退避到期
OR ENQUEUED_AT < NOW R ) -- 进入强补写窗口,覆盖退避
ORDER BY BACKFILL_ATTEMPTS ASC, MSG_ID ASC -- 公平轮转,永久失败行不占满批次
LIMIT PARAM:msgx.pipeline.backfill-batch
```
超期判据使用**本地入队时间**`ENQUEUED_AT`),不使用信箱的 `RECEIVED_AT`:后者来自外部时钟,前偏会在「打标即清除」语义下造成提前清除(`PRE-4`)。
### 6.2 四种结果与放弃
| 结果 | 判定 | 处置 |
|---|---|---|
| 写入成功 | 标记为空、写入 1 行 | 记 `BACKFILL_AT`,不再重试 |
| 早已有标记 | 写入 0 行且信箱行存在 | **视为成功**,不覆盖已有值,记 `BACKFILL_AT` |
| 信箱行不存在 | 写入 0 行且信箱行不存在 | **立即放弃自动重试**(原因 `MISSING_ROW`)并告警。终态行存在而信箱行不存在,只可能是该行在入队后被删除(永久空洞 ID 从不入队,不会进入本扫描) |
| 暂时性故障持续超期 | 超时 / 连接失败持续到 `R` 仍未打标 | **停止自动重试**(原因 `TRANSIENT_DEADLINE`)并告警;`R` 之前只退避重试,**不按尝试次数放弃**;保留人工恢复能力 |
**放弃 ≠ 标记已确认**:放弃行不写 `BACKFILL_AT`,因此不满足 `C-8` 的清除前提,库方不得据此清除;放弃清单需人工对账确认后才可用于清除判定。
### 6.3 `R` 的作用
`R`(取值见 [reference.md](reference.md))有两个作用:
1. **取消退避**:已终态但超期未打标的行,每轮扫描都被尝试,不再等退避到期;
2. **暂时性故障的放弃期限**:超时、连接失败这类暂时性故障**在 `R` 之前只退避重试、不放弃**;到 `R` 仍未打标才停止自动重试、记入放弃清单并告警。
放弃判据用**时间**而不是**尝试次数**:固定次数不能稳定表达允许的故障持续时间,因此按 `R` 判断放弃,`PARAM:msgx.pipeline.backfill-max-attempts` 只用于告警。
关于「最终一定打标」,准确表述是三段,缺一不可:
1. 退避重试(`R` 之前不放弃);
2.`R` 仍失败则停止自动重试、告警,进入放弃清单,保留人工恢复(`reopen`);
3. `C-8` 允许以「放弃清单 + 人工确认」作为清除判定,避免一行永久卡住整个分区。
两个边界要说清:`MISSING_ROW`(信箱行不存在)是**确定性结论**,立即放弃,不受 `R` 保护;`R` 仍然**不保护重放窗口**——`R``R_keep` 只要求 `R ≤ R_keep`,重放窗口的唯一保证来源是 `C-7`
重放窗口的保护只有两条路:约定保留期(`C-6` + `C-7`,目标前提),或另设原文保留通道(`G-REPLAY-CHANNEL`)。若库方清除语义是「打标即可清除」,则当天打标的原文当天即可被清除,增大 `R` 无效。
## 7. 日计划快照与请求匹配
### 7.1 快照发布
`SCHD-DNLD``SCHD-RESP` 共用 `ScheduleProcessor.applyScheduleRecords`
1. **重放判定**`PROC_STATE` 已存在成功终态 → 幂等成功,仅追加留痕,不重复写入。
2. **整包校验**:声明记录数、航班标识与运营日推导等校验失败 → 整包 `DEAD(PROTOCOL)`,不写半包,既有状态保持不变。
3. **事务写入**:锁内按 `FLID` 点查归属日,发现同一航班跨运营日即整包回滚并 `DEAD(PROTOCOL)`;通过后合并写主表与资源明细。报文未携带的航班不因本次日计划报文被删除。
4. **提交结果**:同一事务保存 `KAFKA:schd` / `KAFKA:msg` 事件、置消息 `SUCCEEDED` 并预登记回填意图;提交后信箱回填由扫描承接,留痕在事务外追加。
字段缺失与清空语义、运营日规则见「航班域」。
`RESP` 应匹配开放请求:无匹配、过期或报文早于发送时间时不更新快照(`G-RESP-GUARD`)。
### 7.2 上游请求与静态数据
请求状态机(`G-REQ-TRACK`):
```text
PENDING → SENT → DONE
└──→ EXPIRED
```
- 注册同类新请求前使旧开放请求过期;只有 `COUTMSGS` 写入确认后才标记 `SENT` 并关联出站记录;写信箱成功但本地未确认的情况需要补偿与去重,不能无条件重新发送。
- 应答优先按已确认的回显字段精确匹配;降级匹配的跨代误配风险必须明确接受并审计(`C-23`)。
- 时间比较统一时区与单位,并需定义时钟偏斜容忍;容忍判据未定(`Q5`),在定义前不得把降级匹配描述为精确关联。
- 参考应答写入自有 `REF_MASTER`,日计划应答走快照流程;请求完成必须在相应数据处理成功之后,超时和迟到应答不能修改已关闭请求对应的状态。
- 出站承诺只到落信(`C-24``CLM-8`);主 / 共享删除的 EROR 回报义务见 `C-25`
## 8. 事件投递
### 8.1 普通事件(`KAFKA:msg`
`Dispatcher``MSG_EVENT` 取待发事件。**保序边界是 `FLID`**(与分区键一致):同一 `FLID` 内按 `EVENT_ID` 保序,队头失败即暂停该 `FLID`;不同 `FLID` 之间不互相阻塞,也不承诺跨 `FLID` 顺序。消费者按 `(FLID, STATE_VERSION, UPDATED_AT)` 防旧覆盖新。
本批领取的 `EVENT_ID` 集合在**读取时刻冻结**:发送与标记只作用于这批事件,期间新提交的事件留待下一轮,不参与本批,也不被本批的「完成」带走。
`EVENT_ID` 由全局串行分配产生:事件生产者在事务内写 outbox,主泵单线程,历史清理与主泵互斥(`INV-18`),因此**分配顺序 = 提交顺序**,不存在「已提交的较大 ID 先于未提交的较小 ID 被投递」。
发送确认后才标记 `SENT`,失败记录次数并按退避推后,达到上限转 `DEAD`(记录保留作 DLQ)。所有外部调用需要有界超时,避免阻塞投递线程。
投递是至少一次:Broker 或其他目标已接受但本地未标记成功时可能重发;目标端接受不等于业务消费者已消费。Kafka 生产约束沿用 `D3`,生产者幂等不替代应用层事件去重。
### 8.2 `schd` 聚合
`KAFKA:schd` 只提供最新状态通知,不保留每次中间变化,因此 outbox 按 `FLID` 单行 upsert:同一 `FLID` 只保留最新 `STATE_VERSION` 的事件与投递状态。两条写规则:
`KAFKA:schd` 行的 `EVENT_ID` 不是跨代次稳定的事件句柄:每次接受更新都从全局序列取得新值并替换原主键,用作条件确认的写代次。重放和人工处置只能针对当前 `(TARGET, PARTITION_KEY, EVENT_ID)`;旧代次被替换后不再能按旧 ID 寻址。升级时若已有重复行,按 `STATE_VERSION DESC, EVENT_ID DESC` 保留一行,使迁移与运行时只进不退规则一致。
- **只进不退**:仅当新事件的 `STATE_VERSION ≥` 行内现有版本才覆盖,防止迟到的旧事件把新状态压回去。该合并规则以 `C-21``FLID` 在保留期内不复用)为前提。
- **条件标记**:发送成功后按**读取时刻的版本**做条件标记(`WHERE STATE_VERSION = <本批版本>`);该行若期间已被更新的版本覆盖,则不标记,留待下一轮重发。
发送时:
1. 到期领取批次:按 `FLID` 取未发送行,批次大小受 `PARAM:msgx.schd.flush-limit` 约束;
2. 逐条发送:UPSERT 发送该 `FLID` 的最新整态(key = `FLID`);TOMBSTONE 发送 null 值删除通知;
3. 成功后按上条规则标记完成并推进 `lastFlush`;失败按退避推后,达到上限转 `DEAD`
聚合周期与批上限见 [reference.md](reference.md)。`KAFKA:msg``KAFKA:schd` 之间不承诺顺序。`schd` 行与 `msg` 行共用 `MSG_EVENT`,靠 `TARGET` 区分。
### 8.3 事件清理
`SENT` 的事件行按 `PARAM:msgx.pipeline.event-retention` 由维护作业清理。`KAFKA:msg``DEAD` 行保留作 DLQ,人工处置后再清理;`KAFKA:schd``DEAD` 行只保留到同一 `FLID` 出现新的、可接受的状态代次,新代次会把单行投影重置为 `PENDING` 并清空旧错误。该取舍服从 schd 只保存最新状态的契约,因此被替换的 schd DEAD 代次不再由 `MSG_EVENT` 提供持久审计句柄。
## 9. 失败恢复与维护作业
### 9.1 失败、重试与重放
`ProcFailure``FailureScheduler` 统一处理侧失败落账,投递侧按同一套次数与退避规则迁移事件。启动自检强制退避档位数与尝试上限匹配,让「表里有档但永不触发」的配置无法通过。失败必须在持有具体消息、事件或批次的位置记录,外层循环只做兜底日志和等待,不重复增加次数。线程中断应恢复中断标记并向上传递;不把 JVM `Error` 当普通业务失败捕获。
`ReplayService` 只允许重放白名单内的错误类(见 [reference.md](reference.md))从 `FAILED / DEAD` 回到 `PENDING`,并重置尝试次数、下次重试时间与错误原因,**不重置 `IDENTITY_KEY`**(保留身份,避免重放时把自己判成重复消息)。重放与回填通过 `MessageLifecycleGate` 在同一实例内互斥;**旧消息进入终态后后续消息可能已执行,重新入队不等于恢复历史顺序**,重放前必须评估状态覆盖和版本保护(`CLM-3`)。
### 9.2 中断恢复
恢复的唯一依据是各存储中已持久化的记录,不依赖进程内存状态。在 PG 从备份恢复的场景下,已提交的终态与已发出的事件可能回退,后果是重复投递与重复回填(按至少一次与幂等接受),但不得据此重放业务;留痕(`SCHD_SNAP_LOG`)在业务事务外追加,崩溃会丢该条留痕,不影响状态。
| 中断位置 | 重启后的判定 | 恢复动作 |
|---|---|---|
| 已落信、未入队 | 信箱行位于应扫描的 ID 范围且 PG 无记录(不以处理标记为判据) | 重扫补建入队记录 |
| 水位卡在空洞 | `HOLE_SINCE` 有值且未超过 `PARAM:msgx.pipeline.max-commit-delay` | 等待;超期后放行空洞本身并继续推进 |
| 事务执行中 | PG 无该消息终态 | 事务整体回滚,按 `PENDING` 重新处理 |
| 事务已提交、标记未写 | 终态行仍持有回填意图 | 仅补写标记;业务处理结果保持不变 |
| 标记写入中途 | 标记仍为空 | 重新写入;重复写入同一值无副作用 |
| 兼容入口已入队、水位未追平 | PG 已有该 ID 的记录 | 轮询读到该行时主键幂等,水位照常推进 |
| 回填时信箱行已不存在 | 写入 0 行且信箱行不存在 | 立即放弃自动重试(`MISSING_ROW`)并告警;放弃不等于标记已确认,仍需人工对账 |
| `RECEIVED_AT` 为 NULL | 超期分支以本地 `ENQUEUED_AT` 判定,不受库方时钟与 NULL 影响 | 按 `R` 超期强补写;未超期则按退避重试 |
| 投递目标已接受、`SENT` 未置 | 事件仍 `PENDING` | 允许重发,消费方按事件身份去重 |
| PG 从备份恢复 | 终态与事件回退到备份点 | 按至少一次接受重复;不重放业务、不据此改写航班 |
### 9.3 生命周期与清除
`JobRunner` 用独立 daemon 线程按周期触发回填扫描、航班历史清理与留痕清理;作业不参与消息 FIFO,也不使到期消息饥饿。`INV-18` 要求历史清理的删除与主泵处理互斥。
**通则**(对本系统所有持久对象适用)
- **时间不构成清除依据**:到期只是必要条件,**终局证据才是充分条件**(共享库见 `C-8`,航班见 `D1`)。
- **归档不是终点**:归档目标是新的无界集合,必须有独立保留期与清除作业 `G-HST-RETENTION`,否则只是把容量问题从热表移到冷表。
- **证据不随清除消失**:清除所依赖的证据(如 `C-8` 引用的回填放弃清单,本期承诺见 `C-16`)在其覆盖的信箱边界被清除前必须保持可查。
- **证据缺失或结果不明时按最保守处置**:航班清理为删 0 条(`D1`)。
**逐对象生命周期**(保留期取值一律见 [reference.md](reference.md)
| 对象 | 终局判据 | 归档目标 | 清除证据 | 执行方 | 偏差 |
|---|---|---|---|---|---|
| 共享信箱 `CMINMSGS` 原文 | 处理标记 / 回填放弃清单 | 库方历史表(`C-9` 方案 B | `C-8` | 库方 | 契约未确认(`Q6`/`Q7`/`Q9` |
| `FLIGHT_SCHD` + 资源明细 | 判史规则 | 历史存储 | 归档确认 + 版本复查 | 我们 | — |
| 航班历史存储 | 保留期 | — | — | 我们 | `G-FLIGHT-HIST-RETENTION` |
| `SCHD_SNAP_LOG` | 保留期 | 无(本地可重建) | 无 | 我们 | — |
| `MSG_EVENT` 已发送行 | `SENT` | 无 | 无 | 我们 | — |
| `PROC_STATE` 终态行 | 见下 | `PROC_STATE_HST` | 回填了结 | 我们 | `G-PROC-HST` |
| `PROC_STATE_HST` | 保留期 | — | — | 我们 | `G-HST-RETENTION` |
| `REQ_TRACK` 关闭态行 | 保留期 | 无 | 无 | 我们 | `G-REQ-TRACK-RETENTION` |
原文副本是**条件对象**:仅当库方清除语义不满足 `C-6` 时才成立(`G-REPLAY-CHANNEL``CLM-5`),窗口与 `R_keep` 相同,落在自有 PG。
**处理终态归档**`PROC_STATE``PROC_STATE_HST`
候选 = 终态 **且** 回填已了结 **且** 终局后经过归档阈值(基准是 `UPDATED_AT`:终态与了结都推进它,了结后不再更新)。两处不可省:
- **回填已了结** = `BACKFILL_AT` 非空,或已放弃 **且经人工对账**。放弃行不写标记,是 `C-8` 的清除授权证据,未对账前不得归档。
- **写入前复查** = `DEAD` 可被人工重放改回 `PENDING`。人工重放走 `MessageLifecycleGate`、不取 `PIPELINE_LOCK`,因此该锁不构成复查依据:归档在同一事务内先 `INSERT INTO PROC_STATE_HST … SELECT`,再按候选时的 `STATE` 条件写主表;影响 0 行即整体回滚、该行跳过。重放先一步改回 `PENDING` 时谓词不匹配,天然互斥。批量归档不得持 `PIPELINE_LOCK`——那会阻塞主泵 FIFO,与「作业不使到期消息饥饿」冲突。
归档范围只含终态;归档后仍须保留业务去重能力(`INV-9`)——去重记忆期长于工作状态在线期,实现取「主表保留去重影子行」:主行置 `STATE='ARCHIVED'`、只留 `IDENTITY_KEY``MSG_ID``IDENTITY_KEY` 唯一约束留在主表不动。队头推进、`backlog()` 与回填扫描的谓词显式排除 `ARCHIVED`,不靠状态包含列表隐式过滤。
**时间常数排序**(取值见 [specification.md](specification.md)「契约数值」)
1. `R ≤ R_keep``C-7`)。
2. 去重记忆期 ≥ `R_keep`;否则「归档后重复」不成立(`INV-9`)。
3. 回填放弃清单可见期 ≥ `R_keep`;否则库方清除缺 `C-8` 依据(`C-16`)。
4. 归档阈值计的是**终局之后**的时间,不是入队之后:终态行未了结回填时不进入候选。
**其余清理**
- **航班历史清理**:按 [reference.md](reference.md) 的历史判据选候选(含 `DELETED`),先成功归档再删除;未经 FDEL 的生命周期清除需先补发删除事件。语义与红线见「航班域」。
- **留痕清理**`SCHD_SNAP_LOG` 按保留期与 `(SCOPE_END, RECV_AT)` 删除,不依赖历史存储开关。
- **出站事件清理**:见「事件清理」。
共享信箱保留策略由库方管理(`C-5``C-12`)。历史写入与删除事件入队之间仍需恢复方案;顺序调用不构成原子提交。
## 10. 容量假设与设计取舍
本设计按以下量级选型(参数默认值的依据列见 [reference.md](reference.md),可声明性见 `CLM-10`):
- 单机场、单活动实例、单维护者;入站日消息量千级到万级;单条报文量级 ≤ 10⁴ 字节。
- 处理延迟秒级可接受;航班可见性延迟不劣于现役(轮询间隔 + 聚合周期秒级)。
- 因此:不引入多实例并行、分布式锁、分区表;用单行锁与单线程换确定性。
容量假设变化时,需要重新评估的项:批次大小与轮询间隔、聚合周期与批上限、指标取数口径(`backlog()``PROC_STATE` 聚合;`G-PROC-HST`)、以及 `MSG_EVENT` 保留期。
## 11. 航班域:权威模型与合并写入语义
本章是航班状态的唯一现行设计规范。其他各章只描述管道机制,不重复定义航班域规则。
系统从共享 MySQL 信箱接收 SIS/AODB 报文,把结果合并到自有 PostgreSQL 中的航班当前态,再通过 outbox 投递 Kafka。共享信箱和 Kafka 都不是状态权威,也不在本地事务的提交范围内。
- `FLID` 是航班实例的唯一标识;不得由航班号、日期或资源号推断身份。
- `FLIGHT_SCHD` 及其明细表是唯一权威当前态;展示视图只读,不能作为写入或对账来源(`INV-11`)。
- 单活动主泵按信箱 FIFO 推进。事务内 `PIPELINE_LOCK` 只串行化本地状态提交,不替代选主或消息认领。
- 状态写入、outbox 事件、处理终态和回填意图在同一 PostgreSQL 事务中提交(`INV-17`);回填与 Kafka 投递在提交后独立重试。
现场目标库为 Oracle 11g;Oracle 适配必须通过方言与集成验证后才能作为可切换的运行时选项。
### 11.1 权威模型
| 对象 | 职责 |
|---|---|
| `FLIGHT_SCHD` | 一行一个 `FLID`,保存标量字段、`STATE``STATE_VERSION``OPERATION_DAY`、最近消息 ID 和审计时间。 |
| 资源明细表 | 保存登机门、值机柜台、转盘、计划机位、滑槽、延误、靠撤桥、轮挡等变长集合;`SRVT`/`VIPF` 专用明细见 `G-SRVT-VIPF`。主键为 `(FLID, ORDINAL)`。 |
| `FLIGHT_ROUTE_POINT` | ROUT 与 ERUT 两类路线点,使用 `ROUTE_KIND` 区分;主键应包含该列,避免两类路线的序号冲突。 |
| `PROC_STATE` | 信箱消息的处理终态、业务身份幂等记录,以及回填事实(`RECEIVED_AT` / `BACKFILL_*`)。 |
| `MSG_EVENT` | 事务 outbox,承载整态投影、变更通知和删除 tombstone。 |
| `INBOX_CURSOR` | 共享信箱消费水位(读取进度,与处理标记互不替代)。 |
| `SCHD_SNAP_LOG` | 日计划处理留痕,只追加、可重建,不参与状态决策。 |
### 11.2 航班身份与运营日
`FLID` 是主键。`OPERATION_DAY` 从 SCHD 记录的 `SODT` 按配置的机场时区和切日规则推导;它不是消息接收日或落库日。
一旦已写入非空 `OPERATION_DAY`,同一 `FLID` 不得改到另一个运营日(`INV-12`)。遇到冲突,整包日计划按协议错误拒绝,既有状态保持不变(`INV-19`)。尚未由日计划收录的航班可以为 `NULL`;这不表示该航班没有运营日,只表示当前模型无法为它确定归属日。
### 11.3 字段与集合
标量与异常对象前缀字段存于主表。协议中的 `SRVT``VIPF` 是无界集合,目标形态必须按集合完整保存到专用明细表示;专用明细、合并与投递见 `G-SRVT-VIPF``MAFL` 不是 SIS/XML 入站字段,而是由共享航班的 `MAID``FLID``FLNO` 生成的主航班派生投影(`G-MAFL`)。
- `ORDINAL` 是持久化顺序,从 1 开始;`SOURCE_SEQ` 是上游序号,允许为空或重复。
- 相同资源号不代表同一条分配,禁止按资源号去重。
- 每次持久化完整航班状态时,明细表按该 `FLID` 先删后插,以完整合并结果为准(`INV-14`)。
- ROUT 与 ERUT 是两类独立集合,不能因相同序号覆盖彼此。
- 主/共享关系以主表的 `MAID` 为事实来源:`MAID` 是共享航班指向主航班 `FLID` 的引用(非共享航班为 `NULL`);`MAFL` 只在读取和事件投影时从子航班事实派生,不按入站标量解析或保存。
### 11.4 主/共享投影(`MAFL`
`MAFL` 是主航班的派生集合,元素为子航班的 `FLID``FLNO`;内容与变更传播分别由 `INV-21``INV-22` 保证。
- 子航班集合 = `STATE = ACTIVE``MAID = 主航班 FLID``FLIGHT_SCHD` 行;已 FDEL 的子航班(`STATE = DELETED`)自然退出投影,不需要改写主航班行。
- 只有 `MAID` 为空的主航班携带 `MAFL`;共享航班只携带自身 `MAID``CSOP``CSFT`,不携带 `MAFL`,避免下游双向合并。
- 投影按 `FLID` 升序,与到达顺序及 `FLNO` 变更无关:同一 `STATE_VERSION` 的投影逐字节稳定,重发与消费端比对才有意义。
- `MAID = FLID` 的自引用行不进入任何 `MAFL``MAID` 指向不存在主航班的悬挂引用不阻断该子航班自身处理,只是不产生投影。
- 子航班集合变化(新增、删除、`MAID` 迁移)必须让涉及的主航班在同一事务内推进 `STATE_VERSION` 并登记主航班事件(`KAFKA:msg` + `KAFKA:schd`);否则整态投影的只进不退写入会丢弃它(见「`schd` 聚合」)。共享航班自身不单独发通知。
- 派生主航班投影与产生它的状态写入必须同一事务或一致读快照;按 `MAID` 取子航班要求该列有索引(`INV-17`)。
## 12. 航班域:合并、删除与生命周期
领域决策逻辑(如 `FlightStateEngine` 及各类 Handler 规则)保持纯粹:它根据当前完整态和已解码报文,返回下一完整态与待发事件,不执行数据库或 Kafka I/O。处理器是事务协调器,负责在统一事务边界内调用决策逻辑并持久化结果(`US-03``INV-17`)。
### 12.1 SCHD 日计划
SCHD DNLD/RESP 在整包校验通过后,逐条将报文携带的航班写入当前态。日计划只更新或创建其携带的 `FLID`,**不会因其他航班未出现在本次报文中而删除任何记录**(`INV-15`);SIS 同向(`SIS:3.16-note-1` 要求子系统自行保留前一日延误航班)。
日计划在重叠字段上可以覆盖当前动态值;未携带的字段按合并规则保留,显式清空才清除。每个成功写入的航班推进 `STATE_VERSION``INV-13`),并在同一事务登记 `KAFKA:schd``KAFKA:msg` 事件。
**字段缺失语义与外部规范冲突**:SIS 规定最新日计划中未发送的可选字段表示 AODB 已无该数据、子系统应删除本地已有值(`SIS:3.16-note-4`RESP 与 DNLD 同格式,见 `SIS:3.17`),并要求以 AODB 最新数据覆盖本地(`SIS:1.6.2`)。这与上面的「未携带字段保留」相反。确认前两条并存,按 `Q13` 跟踪,不得据本节推定已与上游对齐。
消息重放由 `PROC_STATE` 的消息 ID 与 `IDENTITY_KEY` 控制;已成功提交的消息不得再次写入或重复登记事件。整包校验失败或运营日冲突时,整包不落地(`INV-19`)。
### 12.2 动态运行事件
FLOP 事件只修改它表达的字段或资源集合,其余航班状态保持不变。每个动态子类型的语义都必须有明确 Handler 规则和回归测试,不能只因已被路由就推定其业务语义完整(`INV-20`)。
动态事件保留既有 `OPERATION_DAY`,也不基于接收时间重新推导它。未知或已删除航班的具体处理遵从对应 Handler 的幂等规则。
### 12.3 删除与重建
FDEL 是业务删除入口:仅在 `ACTIVE → DELETED` 时推进版本、保留明细并与 tombstone 同事务登记;重复 FDEL 或不存在的航班按幂等成功处理。
物理删除仅由独立历史清理在归档成功后执行(`D1`)。日计划报文不是删除依据(`INV-15`)。若未经 FDEL 而由生命周期清理,清理前需要登记一次 tombstone;已经 FDEL 的记录不重复发出。
ADFT 的字段缺失语义尚待上游确认。在确认前采用保守的 Set-only 规则:出现字段可更新,缺失字段不清空;不得把它当成日计划或动态全量替换。新建 ADFT 若带可解析的 `SODT`,按同一运营日规则计算 `OPERATION_DAY`;否则保留为 `NULL`
主/共享航班级联:删除共享航班时重算主航班 `MAFL`(见「主/共享投影」)并向主航班通知;删除主航班时级联删除其子共享关联并发出删除通知;主/共享关系必须一次原子变更,不出现主已删、子残留的半状态。共享航班增量通常只更新并通知主航班,不直接发共享通知。这些语义同样约束 FDEL 之外的生命周期清理。主/共享关联的增删按 `FLID` 做值比较,不使用引用比较。
SIS 规定删除主航班时必须先删子共享航班、再删主航班,顺序不符时 RMS 应向 AODB 回发 EROR`SIS:1.6.1-1.d`,事件定义见 `SIS:4.8`)。本章的原子级联不发该回报,两者取舍见 `C-25`
### 12.4 Kafka 与读取
`KAFKA:schd` 是按 `FLID` 的完整状态投影。Dispatcher 可以合并同一 `FLID` 尚未投递的中间版本,只发最新状态;消费端用 `(FLID, STATE_VERSION, UPDATED_AT)` 防止旧投影覆盖新状态。
`KAFKA:msg` 只通知变化,不承载权威状态;两个 topic 不承诺顺序。FDEL 和必要的生命周期清理使用 tombstone:键为 `FLID`,删除记录以 null 值投递,通知下游移除旧状态。
读取完整航班必须在明确的一致性读边界内批量加载主表和全部明细。
### 12.5 生命周期
运营日过去不等于航班结束。历史清理须同时满足配置保留期与终态证据或足够静默期,先成功写入历史存储,后物理删除当前态;历史存储失败时必须删除零行(`D1``G-FLIGHT-HIST-RETENTION`)。