docs: 同步取报-处理-回填链路的修复结论

- 事务边界表:区分业务型/非业务型终态;回填改为扫描驱动(主泵不再 inline 回填)
- 符号与配置:补 delivery-batch / backfill-max-attempts / cutover-watermark /
  late-detect-* / backlog-cache-ttl-ms / 邮箱超时与指标口径
- 闭环状态:G2/G3/G9 标记已修;G1 标注「阶段 0 只读迟到检测已交付」
- 修正 R/R_keep 论证:R 不保护重放窗口;原文保留归 R_keep 与 Q7/Q9,
  并删除「由 SIS Expiry 推导提交时延」的无效推论
- 回填语义补全:四种结果、放弃 ≠ 标记已确认、公平轮转与饥饿说明
- 兼容入口与水位:登记行超出水位、水位追平前不被领取,并写明运维含义

Plane: ACM2-35 ACM2-37 ACM2-38 ACM2-39 ACM2-41
This commit is contained in:
windyboy
2026-09-11 08:08:39 +08:00
parent 1e75a81107
commit fed0b5df26
2 changed files with 301 additions and 83 deletions
+107 -27
View File
@@ -9,6 +9,29 @@
本文描述处理机制与流程;尚未交付的能力在本文件中明确标注,并以 §10 差异为准。航班状态规则统一由
[运营航班状态设计](flight-state.md) 维护。
**文档标记约定**(全文沿用,[message-lifecycle.md](message-lifecycle.md) 同):
- 【目标】= 设计要求,是否已交付以「实现差异」节为准;
- 【现状】= 已按设计实现的行为;
- 【缺口】= 尚未实现,条目登记在本文 §10 与 [message-lifecycle.md](message-lifecycle.md) §12
- 【待确认】= 依赖开放问题(Q 编号),当前取值只是假定。
正文只描述目标设计。需要点明交付状态时,用上述标记写一句,不展开解释——缺口的唯一清单是 §10 与 [message-lifecycle.md](message-lifecycle.md) §12。
### 1.1 符号与术语
| 符号 / 术语 | 语义 | 存储与字段 | 配置键 | 当前取值 | 约束与唯一定义处 |
|---|---|---|---|---|---|
| `W`(水位) | 信箱 ID 的连续上界:`(min, W]` 已全部读入自有 PG | `INBOX_CURSOR.COMMITTED_UP_TO` | — | 初值 0 | 只随新 ID 成功入队推进(永久空洞放行是唯一例外);[message-lifecycle.md](message-lifecycle.md) §5.1 |
| `holeSince` | `W+1` 处空洞最早被观测到的时刻;无空洞时为 NULL | `INBOX_CURSOR.HOLE_SINCE` | — | NULL | 跨重启保留;旧空洞补齐后新空洞重新计时;[message-lifecycle.md](message-lifecycle.md) §5.1 |
| `R`(超期补写期限) | 回填长期失败时的强制补写上限(按接收时间计) | 判据用 `PROC_STATE.RECEIVED_AT` | `msgx.pipeline.overdue-backfill` | 30 天 | 仅 `R ≤ R_keep`;**不保护重放窗口**(回填由扫描驱动、不等 `R`,打标时刻与 `R` 解耦);[message-lifecycle.md](message-lifecycle.md) §5.2 |
| `R_keep`(清除保留期) | 信箱行可被物理清除前的最短保留时间 | 库方侧 | 库方策略 | 待定(Q9) | `R_keep ≥ max(人工重放期限 + 人工处置期限, 审计期限, 回填重试上限)`**仅在"标记 + 保留期"清除语义下成立**;[message-lifecycle.md](message-lifecycle.md) §6 |
| `head-deadline` | 队头滞留上限,超过即毒丸升级 | 计时起点 `PROC_STATE.PROCESSING_STARTED_AT` | `msgx.pipeline.head-deadline` | 10 分钟 | 为空时以 `UPDATED_AT` 兜底;§3.2 |
| 队头 | 最小的未完成消息(`PENDING``FAILED` 都占位) | `PROC_STATE``ORDER BY MSG_ID`) | — | — | 单线程串行处理,后续消息不得越过;§3.2 |
| 终态 | `SUCCEEDED` / `SKIPPED` / `DEAD` | `PROC_STATE.STATE` | — | — | 到达后队列方可推进;§2.3 |
| 回填意图 | 「还欠一次信箱标记」的持久化事实 | `PROC_STATE.BACKFILL_NEXT_AT` 非空 | — | 随终态写入 | 与终态同一条记录、同一条语句;§3.3 |
| 处理标记 | 信箱行上表示「本系统已处理」的约定字段 | 信箱 `DATE_PROCESSED` / `STATUS` | `mailbox.processed-value` | `PROCESSED` | 只写空标记,不回撤、不覆盖;[message-lifecycle.md](message-lifecycle.md) §5.2 |
## 2. 数据与领域模型
### 2.1 持久化记录
@@ -17,12 +40,12 @@
| 记录 | 用途 | 关键约束 |
|---|---|---|
| `PROC_STATE` | 入站消息的处理状态、身份、重试次数、错误原因与回填事实 | `MSG_ID = CMINMSGS_ID` 主键防止重复入队;`IDENTITY_KEY` 唯一约束防止业务重复;按最小未完成消息 ID 取队头;`PROCESSING_STARTED_AT` 是 HOL deadline 的稳定起点`RECEIVED_AT``BACKFILL_*` 承载 [message-lifecycle.md](message-lifecycle.md) §5.2 的超期判据与回填重试。 |
| `PROC_STATE` | 入站消息的处理状态、身份、重试次数、错误原因与回填事实 | `MSG_ID = CMINMSGS_ID` 主键防止重复入队;`IDENTITY_KEY` 唯一约束防止业务重复;按最小未完成消息 ID 取队头;`PROCESSING_STARTED_AT` 是 HOL deadline 的稳定起点(为空时以 `UPDATED_AT` 兜底,见 §3.2);`BACKFILL_NEXT_AT` 非空 = 还欠一次回填(回填意图),`BACKFILL_AT` 非空 = 标记已确认,`BACKFILL_ABANDONED_AT/REASON` 非空 = 已停止自动重试(**不等于**标记已确认);`RECEIVED_AT` 复制自信箱接收时间,**可能为 NULL**,为 NULL 时 [message-lifecycle.md](message-lifecycle.md) §5.2 的超期兜底不生效。 |
| `MSG_EVENT` | 等待投递的事件(outbox) | `EVENT_ID` 决定投递顺序;`TARGET` 区分 `KAFKA:msg` / `KAFKA:schd``PARTITION_KEY` 恒为 `FLID``EVENT_TYPE` 区分 UPSERT 与 TOMBSTONE。 |
| `REQ_TRACK` | 上游请求及应答关联 | 保存请求类型、覆盖运营日、发送方、出站信箱 ID 与发送/完成时间;同类只允许一个开放请求。登记、超时与应答匹配尚未实现(§10)。 |
| `REF_MASTER` | 静态参考数据(目标表) | `(RTYPE, RKEY)` 唯一;尚未建表,客户端与刷新流程见 user-stories.md US-13/US-14US-14 两类映射的存储落点未定)。 |
| `FLIGHT_SCHD` | 航班标量及单值异常字段 | `FLID` 主键;`OPERATION_DAY` 一经确定不可变;版本与最近消息 ID 用于追踪。变长资源集合存于 8 张资源明细表与 `FLIGHT_ROUTE_POINT`,规则见 [flight-state.md](flight-state.md) §2,不在此重复。 |
| `INBOX_CURSOR` | 共享信箱消费水位 `W` | 单行游标;只随新 ID 成功入队推进,遇空洞即停,空洞超期判定为永久([message-lifecycle.md](message-lifecycle.md) §5.1。 |
| `INBOX_CURSOR` | 共享信箱消费水位 `W` 与空洞计时 `holeSince` | 单行游标;`W` 只随新 ID 成功入队推进(永久空洞放行是唯一例外),遇空洞即停;`HOLE_SINCE` 持久化空洞观测时刻,**进程重启不丢失计时**;[message-lifecycle.md](message-lifecycle.md) §5.1。 |
| `PROC_STATE_HST` | 终态处理记录的归档目标 | 尚未建表;不得改写为共享库历史表。 |
字段与索引定义以 `src/main/resources/db/migration/V1__flight_state_baseline.sql` 为准;Oracle 11g 的迁移形态见 `src/main/resources/db/migration/oracle11g/`(占位,未接入任何 Flyway 配置)。报文原文仍从共享信箱读取,因此必须协调原文保留期,不能在消息尚需处理或重放时提前清理;生命周期、回填与清除契约的唯一定义见 [message-lifecycle.md](message-lifecycle.md)。
@@ -37,9 +60,9 @@
SNDR | TYPE | STYP | SEQN
```
接收时只按信箱 ID 去重;解码后才首次绑定业务身份。重试保留原有绑定,不能把自己判为重复消息。身份被另一条记录占用时,当前消息转为 `SKIPPED`,记录 `duplicate-of:<id>`。是否加入日期边界取决于上游序号重置周期,默认关闭(`SEQN` 的取值范围与回绕已由 `SIS_AODB_RMS-V0.1.md` §2.8.1 定义,重置周期见 Q11);上线后不能随意更换身份算法。
接收时只按信箱 ID 去重;解码后才首次绑定业务身份。重试保留原有绑定,不能把自己判为重复消息。身份被另一条记录占用时,当前消息转为 `SKIPPED`,记录 `duplicate-of:<id>`。是否加入日期边界取决于上游序号重置周期,默认关闭(`SEQN` 的取值范围与回绕已由 `SIS_AODB_RMS-V0.1.md` §2.8.1 定义,重置周期见 Q11);上线后不能随意更换身份算法。身份绑定是**独立的幂等单语句**`WHERE IDENTITY_KEY IS NULL`),不参与业务事务,见 §3.3 表 1。
分派与落库由 `MessageProcessor` 统一协调:按 `MsgKind` 把已绑定身份的队头消息交给对应事务协调器(DNLD/RESP → `ScheduleProcessor`ADFT → `AdftProcessor`FLOP → `FlopProcessor`FDEL → `FdelProcessor`,其余 → `FAILED(UNSUPPORTED)`)。这些 Processor 在 `PIPELINE_LOCK` 事务内读取当前完整态,调用纯领域决策逻辑得到下一完整态与待发事件,再统一落库并登记回填意图;它们不直接触碰 Kafka。处理终态与业务变更在同一事务边界提交(见 flight-state.md §4
分派与落库由 `MessageProcessor` 统一协调:按 `MsgKind` 把已绑定身份的队头消息交给对应事务协调器(DNLD/RESP → `ScheduleProcessor`ADFT → `AdftProcessor`FLOP → `FlopProcessor`FDEL → `FdelProcessor`,其余 → `FAILED(UNSUPPORTED)`)。这些 Processor 在 `PIPELINE_LOCK` 事务内读取当前完整态,调用纯领域决策逻辑得到下一完整态与待发事件,再统一落库并登记回填意图;它们不直接触碰 Kafka。处理终态与业务变更在同一事务边界提交,**该结论只对处理器产出的「业务型终态」成立**;`MALFORMED / PROTOCOL / SKIPPED(重复) / EXHAUSTED` 这类非业务型终态不涉及航班表,只需一条 `PROC_STATE` UPDATE(终态与回填意图同语句写入),不取 `PIPELINE_LOCK`。完整的事务边界见 §3.3 表 1
### 2.3 状态与错误分类
@@ -65,43 +88,84 @@ SNDR | TYPE | STYP | SEQN
| `INFRA` | 基础设施或执行异常,退避重试。 |
| `EXHAUSTED` | 重试耗尽或滞留超时,转 `DEAD`,人工复核后允许重放。 |
重试次数用尽时统一转 `DEAD(EXHAUSTED)``ERROR_CLASS` 被覆写为 `EXHAUSTED`**原始错误类别不再保留**`LAST_ERROR` 保留原因文本)。由于重放白名单包含 `EXHAUSTED`,这类记录仍可人工重放(§6.1)。
## 3. 收报与主泵
### 3.1 收报
`InboxPoller` 默认每秒按 ID 升序、有限批次(`claim-batch`,默认 50)读取水位之后的信箱记录(`ID > W`,**不以处理标记为谓词**),在自有 PG 建立 `PENDING` 并把水位推进到连续上界;入队与水位推进在同一 PG 事务内提交,重复扫描幂等、中断后重扫补建。收报层不解析业务载荷,也不回填已处理标记。
水位 W 与扫描谓词的完整口径(连续上界、遇空洞即停、空洞老化)唯一见 [message-lifecycle.md](message-lifecycle.md) §5.1;水位不是已处理标记,其有效性以 Q2 的 ID 单调承诺为前提。空洞老化阈值取 `msgx.pipeline.max-commit-delay`
**收报流程**(每轮 `pollOnce`;完整判据、代价与前提见 [message-lifecycle.md](message-lifecycle.md) §5.1,本节不重复):
1. 读游标 `(W, holeSince)`;信箱不可读时记日志、等下一轮,**不动水位**——这属于基础设施失败,不能当成「没有新消息」。
2.`ID > W` 的升序前 `claim-batch` 行。
3.`W+1` 起逐 1 数求连续上界;出现缺号时按「遇空洞即停 / 超期放行」处理,本批缺号之后的行本轮一律不入队(判据见 §5.1)。
4. 在同一个 PG 事务内:对水位以内的每一行 `insertIfAbsent(MSG_ID, RECEIVED_AT)`,并写回 `(W, holeSince)`;主键冲突表示已入队,不计入也不报错。
5. 提交。本批中因空洞或批次上限未入队的行留待下一轮——**每轮最多解决一个空洞**。
水位不是已处理标记,其有效性以 Q2 的 ID 单调承诺为前提;空洞老化阈值取 `msgx.pipeline.max-commit-delay`
**单实例前提**:信箱读取不加锁,水位是单行覆盖写。本设计只在单活动实例下成立([architecture.md](architecture.md) §5、D2);多实例并发收报会让水位互相覆盖,必须先有实例级排他。
兼容 HTTP 入口执行“写入共享信箱 → PG 入队”。两步不在同一事务中:信箱成功而 PG 失败时,原文不能丢失,由轮询补建;客户端失败重试可能再次写信箱,业务身份去重仍然必需。
兼容入口只写 `PROC_STATE`、**不参与水位**,因此它登记的行会**超出水位**;主泵在水位追平前不领取(§3.2 步骤 2),顺序因此不受影响——代价是这类消息要等收报把 `W` 推到它的 ID 之后才开始处理(最长约一个空洞老化窗口)。**运维含义**:只使用兼容入口而不运行收报轮询时,这些行不会被处理,必须让 `InboxPoller` 运行(`msgx.pipeline.autostart=true` 或显式触发)。详见 [message-lifecycle.md](message-lifecycle.md) §5.1/§12。
### 3.2 主泵调度
每次 `Pump.tick` 只围绕最小未完成消息(`PENDING``FAILED` 都占队头):
1. 无队头:按轮询间隔休眠。
2. 队头 `FAILED` 且未到 `next_attempt_at`:未超限则等到可重试时刻;已达重试上限或超过队头滞留时限(`head-deadline`,默认 10 分钟)则转 `DEAD(EXHAUSTED)`
3. 队头可执行:调用 `MessageProcessor.processOne`,失败迁移在该边界内完成
2. 队头超出水位(`msgId > W`):**不领取**,休眠到下一轮。这类行只可能来自兼容入口的直接登记;允许领取会让它越过尚未入队的较小 ID(G2)。会打印一条**限流 WARN**(仅在水位值变化时打一次)
3. 队头`FAILED` 且已毒丸(`attempts ≥ max-attempts`,或 `now 计时起点 ≥ head-deadline`):转 `DEAD(EXHAUSTED)`,并立即尝试一次回填。**该分支只写 `PROC_STATE`,不取 `PIPELINE_LOCK`、不在处理器事务内,也不在 `MessageLifecycleGate` 内**(见 §3.3 表 1 与 §6.1
4. 队头为 `FAILED` 且未到 `next_attempt_at`:休眠到可重试时刻,不处理后续消息。
5. 其余(新消息或退避到期的重试):记录 `PROCESSING_STARTED_AT`(仅首次)后调用 `MessageProcessor.processOne`,失败迁移在该边界内完成。
维护作业由独立 job 线程调度(§6.1),不占用消息循环;作业有界且不使到期消息无限饥饿。调度取时经可注入 `Clock`HOL deadline 以首次处理时写入的 `PROCESSING_STARTED_AT` 为稳定起点;人工重放会清空该值,由新一轮首次处理重新记录。
维护作业由独立 job 线程调度(§6.1),不占用消息循环;作业有界且不使到期消息无限饥饿。**所有取时统一经注入 `Clock`**(收报空洞老化、主泵调度、处理器落库时间、回填重试、作业切日),不使用系统时钟;HOL deadline 以首次处理时写入的 `PROCESSING_STARTED_AT` 为稳定起点,**该列为空时(V3 迁移之前的存量行)以 `UPDATED_AT` 兜底**;人工重放会清空 `PROCESSING_STARTED_AT`,由新一轮首次处理重新记录。
### 3.3 单条处理
```text
读取原文 → 解码(MALFORMED → DEAD;编码错误 → FAILED 退避)
→ 首次绑定身份(冲突 → SKIPPED,记 duplicate-of
→ 按 MsgKind 分派处理器
DNLD / RESP → ScheduleProcessor(快照事务)
ADFT / FLOP / FDEL → Adft / Flop / FdelProcessor(单航班事务)
Unsupported → FAILED(UNSUPPORTED)
PG 单事务:锁 + 航班变更 + 待发事件 + 终态 + 回填意图预登记
→ 提交后回填信箱;失败由补偿待办重试
processOne(head)
1. 入口守卫:head 已是 FAILED 且 attempts 达上限 → DEAD(EXHAUSTED),结束
2. 读原文:缺失 → DEAD(MALFORMED, raw-missing);读取异常 → FAILED(INFRA)
3. 解码:
报文非法(MALFORMED) → DEAD(MALFORMED),不重试
可修复解码错(CODEC_ERROR) → FAILED(CODEC_ERROR) 退避
4. 身份绑定(仅当 IDENTITY_KEY 为空):
已被本消息占用 → 继续
已被别的消息占用 → SKIPPED(duplicate-of:<id>),结束
空闲 → 写入 IDENTITY_KEY(独立单语句,不参与业务事务)
5. 按 MsgKind 分派:
SCHD-DNLD / SCHD-RESP → ScheduleProcessor(快照事务)
SCHD-ADFT → AdftProcessor(单航班事务)
FLOP / FDEL → Flop / FdelProcessor(单航班事务)
Unsupported → FAILED(UNSUPPORTED) 退避
载荷缺失 → DEAD(MALFORMED)
整包协议拒绝 → DEAD(PROTOCOL),不落半包
6. 业务型成功:处理器在自己的事务内写航班变更 + 待发事件 + SUCCEEDED + 回填意图
7. 返回终态标记:只有终态才调用 backfill.attempt(msgId) 试写一次信箱标记
```
**表 1 事务边界**(哪些动作在一个事务里、哪些不是):
| 动作 | 显式事务 | 持 `PIPELINE_LOCK` | 触及航班表/事件 | 原子性来源 |
|---|---|---|---|---|
| 收报入队(`insertIfAbsent` + `cursor.save`) | 是 | 否 | 否 | 同库事务 |
| 身份首次绑定 | 否 | 否 | 否 | 单语句 + `uk_proc_identity` |
| 业务型终态(处理器产出 `SUCCEEDED`) | 是 | 是 | 是 | 同库事务:航班变更 + 事件 + 终态 + 回填意图 |
| 非业务型终态(`MALFORMED` / `PROTOCOL` / `SKIPPED` / `EXHAUSTED`) | 否 | 否 | 否 | 单语句(终态与回填意图同一条 UPDATE) |
| 毒丸升级 `DEAD(EXHAUSTED)`(§3.2 步骤 2) | 否 | 否 | 否 | 单语句 |
| 回填(信箱标记 + `BACKFILL_AT`) | 否 | 否 | 否 | 跨库两次单写;幂等可重跑 |
| 人工重放(批量改回 `PENDING`) | 否 | 否 | 否 | 单语句批量;`MessageLifecycleGate` 与回填互斥 |
结论:**「航班变更与处理终态同事务」只对业务型终态成立**;非业务型终态不涉及跨表一致性,因此不需要 `PIPELINE_LOCK`,但终态与回填意图仍由同一条 UPDATE 保证不分离。
- 原文缺失归为 `MALFORMED`;读取异常不能伪装成“缺失”,应进入基础设施重试。
- 忽略规则(`LDM / REGN / RSTA / EROR``SKIPPED`)尚未实现(§10);不能因类型未覆盖就把合法忽略报文当非法报文处理。
- 航班变更、待发事件、处理终态与回填意图同一 PG 事务原子提交(终态由处理器在自己的事务内落库,`MessageProcessor` 不再单独补写终态);跨存储双写窗口已根除
- 终态提交后立即尝试一次回填,失败由回填扫描按退避重试,抵达超期期限 R 时按 [message-lifecycle.md](message-lifecycle.md) §5.2 强制补写(US-09/Q7)。回填只需消息 ID,因此缺 META 或解码失败的死信同样可补写。`PENDING / FAILED` 禁止回填;影子环境禁写
- **主泵不回填**终态与回填意图同一条 UPDATE 落库,回填**一律由扫描驱动**(跨库写不能占用 FIFO 关键路径)。失败按退避重试,抵达超期期限 R 时强制补写,确认行不存在或达尝试上限则停止自动重试([message-lifecycle.md](message-lifecycle.md) §5.2US-09/Q7)。回填只需消息 ID,因此缺 META 或解码失败的死信同样可补写。`PENDING / FAILED` 禁止回填;影子环境禁写
- 回填有三种结果:写入成功;**此前已被标记(视为成功,不覆盖已有值)**;信箱行不存在(视为失败,当前无终态,见 §10)
## 4. 日计划快照与请求匹配
@@ -157,13 +221,13 @@ REGISTERED → SENT → WAITING → DONE
### 6.1 失败、重试与重放
`ProcFailure``FailureScheduler` 统一处理侧失败落账,投递侧(`Dispatcher`)按同一套次数与退避规则迁移事件。默认最多 5 次(attempts ≥ 5 判耗尽)退避档位 1、2、4、8、16 秒、单档封顶 60 秒时间经可注入 `Clock` 判定。
`ProcFailure``FailureScheduler` 统一处理侧失败落账,投递侧(`Dispatcher`)按同一套次数与退避规则迁移事件。默认最多 5 次(attempts ≥ 5 判耗尽)退避表配置为 1、2、4、8、16 秒、单档封顶 60 秒,但**耗尽判定与退避取值同源**`attempts ≥ max-attempts` 即转 `DEAD`,不再计算下次重试),因此默认配置下实际只用 1、2、4、8 四档:**16 秒档与 60 秒封顶不会被触发**(要么调高 `max-attempts`,要么接受"5 次尝试 = 4 档退避")。时间经可注入 `Clock` 判定。
失败必须在持有具体消息、事件或批次的位置记录,外层循环只做兜底日志和等待,不重复增加次数。线程中断应恢复中断标记并向上传递;不把 JVM `Error` 当普通业务失败捕获。
`ReplayService` 只允许 `CODEC_ERROR / UNSUPPORTED / INFRA / EXHAUSTED``FAILED / DEAD` 回到 `PENDING`重置次数与下次执行时间,保留身份与错误审计。它按错误类整批重放,尚无按记录预检、操作审计与管理入口(US-10)。旧消息进入终态后后续消息可能已执行,**重新入队不等于恢复历史顺序**;人工重放前必须评估状态覆盖和版本保护,不能直接批量重放到生产。
`ReplayService` 只允许 `CODEC_ERROR / UNSUPPORTED / INFRA / EXHAUSTED``FAILED / DEAD` 回到 `PENDING`重置 `ATTEMPTS=0``NEXT_ATTEMPT_AT=NULL``PROCESSING_STARTED_AT=NULL`**不重置 `IDENTITY_KEY`**(保留身份,避免重放时把自己判成重复消息)。它按错误类**全局批量**重放,尚无按记录预检、操作审计与管理入口(US-10)。重放与回填通过 `MessageLifecycleGate` 在同一实例内互斥,避免「旧回填给已重新入队的消息写标记」;【缺口】**毒丸升级路径不在该 gate 内**(§3.2 步骤 2),该竞态窗口登记于 §10。旧消息进入终态后后续消息可能已执行,**重新入队不等于恢复历史顺序**;人工重放前必须评估状态覆盖和版本保护,不能直接批量重放到生产。
**维护作业**`JobRunner` 用独立 daemon 线程每 30 秒触发 `BackfillService.sweep`回填补写扫描,指数退避 30 秒起步、封顶 15 分钟),每天机场时区 03:30 后触发一次 `HistorySweepJob`(§6.2)。回填意图在终态事务内登记在 `PROC_STATE``BACKFILL_NEXT_AT/ATTEMPTS/ERROR`),不再有独立待办表。作业不再经 `PUMP_JOB` 队列插队,不参与消息 FIFO,也不使到期消息饥饿。作业与回填通道的生命周期口径见 [message-lifecycle.md](message-lifecycle.md) §3/§4。
**维护作业**`JobRunner` 用独立 daemon 线程每 30 秒触发 `BackfillService.sweep`**调度周期,不是回填完成时限**:批次积压、单行调用超时与历史作业耗时都会延长实际延迟;回填扫描自身的退避 30 秒起步、封顶 15 分钟),每天机场时区 03:30 后触发一次 `HistorySweepJob`(§6.2)。回填意图在终态的同一条语句里登记在 `PROC_STATE``BACKFILL_NEXT_AT/ATTEMPTS/ERROR`),不再有独立待办表。扫描谓词见 [message-lifecycle.md](message-lifecycle.md) §5.2;超过超期期限 `R` 后超期分支恒成立并覆盖退避,但**永久失败不会无限重试**:确认行不存在立即放弃、暂时性故障到 `backfill-max-attempts` 后停止自动重试(两类都告警且可人工恢复)。扫描按 `BACKFILL_ATTEMPTS, MSG_ID` **公平轮转**并排除已放弃行,最旧的一批永久失败行不会再占满批次饿死后续记录。作业不再经 `PUMP_JOB` 队列插队,不参与消息 FIFO,也不使到期消息饥饿。作业与回填通道的生命周期口径见 [message-lifecycle.md](message-lifecycle.md) §3/§4。
### 6.2 历史清理与归档
@@ -183,14 +247,18 @@ REGISTERED → SENT → WAITING → DONE
- `msgx.pipeline.autostart``msgx.stubs`:分别控制管道启动与内存适配器;生产禁止 stub,默认不自动启动。
- `msgx.pipeline.poll-interval / claim-batch / max-attempts / backoff-ms / backoff-cap-ms / head-deadline`:控制轮询节奏、批次、重试上限、退避与队头滞留;这些参数不能改变 FIFO。
- `msgx.pipeline.max-commit-delay / overdue-backfill / backfill-batch`:空洞老化阈值(Q2 最大提交时延)、超期补写期限 RQ6回填扫描批量R 与老化阈值都不能为提速而下调
- `msgx.pipeline.max-commit-delay / overdue-backfill / backfill-batch / backfill-max-attempts`:空洞老化阈值(Q2 的「ID 分配 → 事务可见时延上界」,同时决定目标补偿扫描的窗口宽度)、超期补写期限 RQ6回填扫描批量与回填自动重试上限(达上限停止自动重试并可人工恢复);R 与老化阈值都不能为提速而下调。默认 5 分钟**只是缺少依据的占位假定值**——库方尚未给出该可见性时延上界,且它**不能由 SIS 报文的 `Expiry`480 分钟量级)推导**[message-lifecycle.md](message-lifecycle.md) §5.1),Q2 确认前属于上线门槛。`overdue-backfill`R)的完整约束见 [message-lifecycle.md](message-lifecycle.md) §5.2
- `msgx.pipeline.delivery-batch / delivery-drain-rounds`:普通事件(`KAFKA:msg`)批量投递的批大小与每轮最多连取批数;把投递从"每条一次 DB 往返 + 一轮一次 sleep"提升到由下游决定,同时保留"队头失败即停止本轮"的目标内保序。
- `msgx.operation-day.zone / cutoff-hour`:运营日时区与切日边界,决定 `OPERATION_DAY` 推导(flight-state.md §2.1)。
- `mailbox.processed-value`:写回共享信箱的处理标记值,仅限库方认可的 legacy 值集(Q7)。
- `msgx.schd.flush-period / flush-limit`:控制状态通知的聚合延迟与批量大小。
- `msgx.identity.include-day-boundary`:影响去重语义,不能作为普通调优项切换;序号重置周期见 Q11。
- `mailbox.shared-mysql.enabled``msgx.history.history-store-enabled`:分别门控真实信箱与历史存储接线,默认关闭。
- `msgx.health.backlog-cache-ttl-ms`:积压快照缓存窗口(默认 30 秒,`/health``/metrics` 共用);设为 0 仅用于测试/排障,不作为实时性的替代。
- `msgx.pipeline.cutover-watermark`:**一次性、显式**的切流播种(默认不配置 = 不播种)。取值为 `min`(读当前全部现存行,`W=MIN(ID)1`)、`zero`(从 0 按空洞规则扫描)、`max`(跳过当前可见存量,`W=MAX(ID)`)或具体 ID。代码**不做默认选择**、也不会自动退化成 `max`;升级实例(已有水位或已有处理记录)**拒绝重新播种**,重新切流必须是显式操作;播种事实记在 `INBOX_CURSOR.SEEDED_AT`(与水位在同一条语句落库),而该列为 NULL **不等于**从未消费。非法取值由启动自检挡下。
- `msgx.pipeline.late-detect-period / late-detect-batch`**只读**迟到检测(ACM2-41 阶段 0)的周期与单轮复查量;周期设 0 即关闭。检测只计数与告警,不补入队、不改变处理语义。
日志关联消息 ID、事件 ID 和批次;失败记录错误分类、次数、下次执行时间。健康检查反映依赖实际可用性;队头滞留、积压、死信和补偿失败需要指标及告警。日志出口故障不得阻塞业务线程。
日志关联消息 ID、事件 ID 和批次;失败记录错误分类、次数、下次执行时间。健康检查反映依赖实际可用性;队头滞留、积压、死信和补偿失败需要指标及告警。指标经 Micrometer 暴露(`PipelineMetrics`,启动时急切注册):`msgx.pipeline.backlog.unfinished``msgx.pipeline.backlog.oldest_unprocessed_seconds``msgx.pipeline.backfill.unmarked_terminal``msgx.pipeline.backfill.abandoned``msgx.pipeline.backfill.oldest_unmarked_seconds``msgx.pipeline.watermark.lag``msgx.pipeline.hole.aged_out.total``msgx.pipeline.late_arrival.detected.total`(迟到检测命中数——**> 0 表示上游提交确实晚于水位推进,需要与库方对契约**)。取数统一走 `BacklogSnapshotProvider``msgx.health.backlog-cache-ttl-ms`,默认 30 秒),`/health``/metrics` 共用同一份快照——`backlog()``PROC_STATE` 的全表聚合,不能被高频抓取打穿;无法取数时上报 `NaN`,不伪造 0。日志出口故障不得阻塞业务线程。
## 8. 验证要求
@@ -198,12 +266,21 @@ REGISTERED → SENT → WAITING → DONE
| 场景 | 必须验证的结果 |
|---|---|
| 重复扫描、入队中断、较小 ID 迟到 | 不重复入队、不丢记录、不让后续消息越序;终态而未回填的行不得阻断后续消息发现。 |
| 重复扫描、入队中断 | 不重复入队、不丢记录、不让后续消息越序;终态而未回填的行不得阻断后续消息发现。 |
| 较小 ID 迟到(Q2 未确认) | 【缺口 · 已固定基线用例】快路径只读 `ID > W``InboxPollerTest` 已把"水位越过后到达的较小 ID 不被发现"钉成基线。补偿扫描(G1)交付前禁止任何"迟到不越序"的验收声明。 |
| 空洞老化与重置 | 阈值内不推进水位、不越过入队;超期后只放行空洞本身;旧空洞补齐后出现的新空洞获得完整等待窗口。 |
| 兼容入口与空洞并发 | 兼容入口登记的行超出水位,主泵不领取;水位追平后按序处理。端到端用例:`PipelineSmokeTest`「compat injected high id is not claimed until the watermark catches up」。 |
| 队头失败、退避及作业竞争 | 消息不越队;到期后恢复;作业不使消息无限饥饿。 |
| 同身份多条记录、失败后重试、归档后重复 | 只产生一次有效业务处理,不把自身重试判为重复。 |
| 非业务型终态 | 不触碰 `FLIGHT_SCHD` / `MSG_EVENT`,只写 `PROC_STATE`,且终态与回填意图同语句生效。 |
| PG 事务失败、快照重复或迟到 | 整体回滚重试、不重复推进版本、不回退状态、不误删增量航班。 |
| 整包协议拒绝(运营日冲突、声明数不符) | 整包不落地、整体回滚,既有状态与版本不变,消息终态为 `DEAD(PROTOCOL)`。 |
| PG 提交失败、信箱回填失败 | 事件、处理终态与回填意图一起回滚;已提交结果只补写标记,不重放业务;中间态永不补写。 |
| 回填四种结果 | 写入成功 / 早已标记(不覆盖、记成功)/ 信箱行不存在(立即放弃并告警,**不得**视为已标记)/ 暂时性故障达上限(停止自动重试,可人工恢复)。 |
| 回填公平性 | 最旧的一批记录永久失败时,后续待回填记录仍能被扫描到;已放弃行不再进入扫描。 |
| 「打标即清除」语义 | 【待确认】若库方清除由打标触发,必须验证存在约定的保留期或独立原文保留通道,且**不依赖 `R` 的取值**。 |
| `RECEIVED_AT` 为 NULL | 超期兜底不生效的行为被显式验证,且不会导致标记被提前写入。 |
| 毒丸升级与人工重放并发 | 【缺口】毒丸路径不在 gate 内;窗口必须被复现,或用 gate 覆盖后验证互斥。 |
| 投递确认丢失、批次失败、次数耗尽 | 允许可识别的重发、保持目标顺序、整批退避并保留死信。 |
| 请求超时、无匹配 RESP、时间单位不一致 | 不误用迟到应答,不提前完成请求。 |
| stub 误配置、重复实例、停机中断 | 生产拒绝不安全启动,工作线程能正确退出。 |
@@ -226,10 +303,13 @@ REGISTERED → SENT → WAITING → DONE
## 10. 当前实现差异
以下缺口直接影响上述设计是否成立,不能以类或接口已存在作为完成依据
以下缺口直接影响上述设计是否成立,不能以类或接口已存在作为完成依据。收报与处理链路的逐条缺口另见 [message-lifecycle.md](message-lifecycle.md) §12。
- **事务与外部副作用**状态、事件、处理终态与回填意图 PG 事务提交;回填意图落在 `PROC_STATE``BACKFILL_AT/NEXT_AT/ATTEMPTS/ERROR`),`BACKFILL_TODO` 已随 V2 迁移下线,[message-lifecycle.md](message-lifecycle.md) §4 的两个崩溃窗口(提交后回填前崩溃、回填意图二次落账失败)不再存在。回填本身仍是跨库单写,失败按退避重试并由超期期限 R 兜底;R 的取值待 Q6 确认。
- **收报与调度**:水位 W 落库并与入队同事务推进,遇空洞即停、空洞超过 `max-commit-delay` 判定为永久(Q2 未书面确认前该阈值是假定值);旧空洞补齐后出现的新空洞会重置老化起点。较小 ID 迟提交与空洞场景的端到端顺序保证仍待与库方联合验证。HOL deadline 已改用 `PROCESSING_STARTED_AT` 和可注入 `Clock`;Q6 仍需确认长期积压与人工重放的期限口径。
- **事务与外部副作用**业务型终态(处理器产出的 `SUCCEEDED`)与航班变更、待发事件、回填意图在同一 PG 事务提交;非业务型终态(`MALFORMED / PROTOCOL / SKIPPED / EXHAUSTED`)与毒丸升级只写 `PROC_STATE`,终态与回填意图同一条 UPDATE,不涉及跨表一致性(§3.3 表 1)。回填意图落在 `PROC_STATE``BACKFILL_AT/NEXT_AT/ATTEMPTS/ERROR`),`BACKFILL_TODO` 已随 V2 迁移下线,[message-lifecycle.md](message-lifecycle.md) §4 的两个崩溃窗口(提交后回填前崩溃、回填意图二次落账失败)不再存在。回填本身仍是跨库单写,失败按退避重试并由超期期限 R 兜底;R 的取值待 Q6 确认。
- **收报与调度**:水位 W 与空洞计时 `HOLE_SINCE` 落库并与入队同事务推进,遇空洞即停、空洞超过 `max-commit-delay` 判定为永久(Q2 未书面确认前该阈值只是**缺少依据的占位假定值**——库方尚未给出「ID 分配 → 事务可见」的时延上界,且不能由 SIS 报文 `Expiry` 推导);旧空洞补齐后出现的新空洞会重置老化起点。**较小 ID 迟提交目前没有任何发现机制**[message-lifecycle.md](message-lifecycle.md) §5.1 描述的窗口补偿扫描尚未实现,水位越过空洞后到达的较小 ID 不会被快路径(`ID > W`)发现,端到端顺序保证仍待与库方联合验证。HOL deadline 已改用 `PROCESSING_STARTED_AT`(为空时以 `UPDATED_AT` 兜底)和可注入 `Clock`;Q6 仍需确认长期积压与人工重放的期限口径。**阶段 0 只读迟到检测已实装**(监视被放行的空洞 ID,命中即计入 `msgx.pipeline.late_arrival.detected.total` 并告警,不补入队);补偿扫描的阶段 1(按类型安全补入队)与 Q2 契约仍待推进。
- **兼容入口与水位(已修)**`POST /cminmsgs/send` 直接写 `PROC_STATE`、不参与水位,因此登记行会超出水位;主泵只领 `msgId ≤ W`,在水位追平前不领取(`Pump` 已实现,端到端用例守住)。代价是这类消息延迟到水位追平,且**必须让收报轮询运行**([message-lifecycle.md](message-lifecycle.md) §5.1/§12)。
- **回填闭环(已修 · V4)**:三种结果已区分(写入成功 / 早已标记 / 信箱行不存在),并新增放弃语义:确认 `MISSING` 立即放弃、暂时性故障达 `backfill-max-attempts` 后停止自动重试,两者都告警且可由人工恢复;**放弃 ≠ 标记已确认**(`BACKFILL_AT` 仍为空,清除前提不成立)。扫描改为公平轮转,饥饿问题关闭。仍存在的相关风险是"打标即清除"语义下**`R` 无法保护重放窗口**,必须另行约定保留期或引入独立原文保留通道([message-lifecycle.md](message-lifecycle.md) §12 G10)。
- **重放互斥**`MessageLifecycleGate` 只覆盖回填与人工重放,毒丸升级路径不在其中,存在「先标 DEAD 并打标、再被重放拨回 PENDING」的窗口。
- **快照与业务能力**DNLD/RESP/ADFT 与 FLOP/FDEL 处理器、整包校验与跨运营日整包拒绝均已接入;但 RESP 应答守卫与出站请求未实现,忽略规则(US-04)未实现,29 类 FLOP 与参考应答的逐类矩阵未补全,ADFT 缺失字段与 `FLID` 重用语义待上游确认。
- **航班读写**:唯一写入口与权威读已落地;ROUT/ERUT 联合主键、空值/未知属性保真、事件在事务内只算一次、逐航班多次查询仍待修正(见 flight-state.md §6)。`/all/flights` 尚未实现。
- **请求、参考数据与归档**`REQ_TRACK` 表与仓储已建,但无运行时协调与 `COUTMSGS` 出站适配;`REF_MASTER` 未建表;`PROC_STATE_HST` 未建表。航班历史清理脚手架已实现,历史存储未接通时删 0 条。
+194 -56
View File
@@ -4,15 +4,16 @@
本文定义 msgexchange-v2(下称"本系统")对共享 MySQL 信箱中报文的完整生命周期:从上游写入、本系统处理、处理标记回信箱,到信箱数据最终清除。各阶段的输入输出、幂等方式、故障恢复方法,以及本系统与信箱库管理方之间的分工,都在本文界定。
- 模块协作、状态机与配置参数见 [design.md](design.md);系统边界与"不建表、不改结构"的红线见 [architecture.md](architecture.md);验收口径见 [user-stories.md](user-stories.md)。
- 读者包括本系统开发与运维人员、以及负责共享 MySQL 的库方接口人;§6、§7 中标注 Q 编号的条款即需要与库方书面确认的开放项。
- 模块协作、状态机与配置参数见 [design.md](design.md);系统边界与"不建表、不改结构"的红线见 [architecture.md](architecture.md);验收口径见 [user-stories.md](user-stories.md)。符号(`W` / `holeSince` / `R` / `R_keep` / `head-deadline` 等)与配置键的对应关系统一见 [design.md](design.md) §1.1。
- 读者包括本系统开发与运维人员、以及负责共享 MySQL 的库方接口人;**§5§7 中标注 Q 编号的条款**即需要与库方书面确认的开放项。
- 状态标记沿用 [design.md](design.md) §1:【目标】/【现状】/【缺口】/【待确认】。正文描述目标设计,本文范围内的缺口清单集中在 §12;跨文档引用一律带文档名。
**角色与术语约定**(全文沿用):
- **上游**:向信箱写入报文的源头系统(CIIMS、AODB 等)。
- **信箱**:共享 MySQL 中的入站信箱表 `CMINMSGS`;出站方向为 `COUTMSGS`(§7)。
- **库方**:共享 MySQL 的管理方;表结构变更与数据清除只能由库方执行或书面授权。
- **处理标记**:信箱行上表示"本系统已处理"的约定字段;逻辑名 `DATE_PROCESSED` / `STATUS`,实际列名以库方契约为准(legacy 为 `CMINMSGS_DATE_PROCESSED` 等真实列)。
- **处理标记**:信箱行上表示"本系统已处理"的约定字段;逻辑名 `DATE_PROCESSED` / `STATUS`,实际列名以库方契约为准(legacy 为 `CMINMSGS_DATE_PROCESSED` 等真实列)。正文统一用"处理标记",不再混用列名。
- **自有 PG**:本系统唯一的业务数据库 PostgreSQL;信箱与自有 PG 之间不存在跨库事务。
## 1. 五个独立事实
@@ -29,35 +30,76 @@
分环节的原因是存储边界:信箱与自有 PG 之间的写入不共享事务。“落信”以共享 MySQL 为准,“入队”和“处理完成”以自有 PG 为准;“已回填”与“投递确认”各以外部操作成功和本地确认事实共同判定。后两个环节各自独立重试,都可能单独失败。“投递确认”只表示 Kafka Broker 或其他投递目标已接受,不表示业务消费者已消费。
“入队”这一环还包含一个容易被忽略的前提:**发现是完整的**。它不由本系统单独保证,而依赖信箱 ID 单调与最大提交时延两条上游承诺(Q2);承诺不成立时,本系统只能看到"读过的区间",不能证明"读全了"(§5.1)。
## 2. 状态总纲
```text
信箱行 未标记 ────────────────→ 已处理标记 ────────→ 已归档 ────────→ 已清除
↑ 回填/超期补写 §6(库方执行)
PROC_STATE PENDING / FAILED ──→ SUCCEEDED(成功
├→ SKIPPED(业务重复
└→ DEADMALFORMED / PROTOCOL / EXHAUSTED,人工处置)
MSG_EVENT PENDING ──→ SENT
└──→ PENDING(退避重试)──→ DEADDLQ)
信箱行(共享 MySQL
未标记 ──[本系统写处理标记]──→ 已处理标记 ──[库方按保留期执行]──→ 已归档 ──→ 已清除
↑ §6(本系统不执行 DDL
└── 回填与超期补写(§5.2
PROC_STATE(自有 PG
PENDING ──处理成功────────→ SUCCEEDED
│ ↑ ─→ SKIPPED(业务重复,非错误)
│ └── 人工重放(白名单错误类)
└─ 处理失败 → FAILED ──退避到期──→ PENDING
└── 次数耗尽 / 队头滞留超限 ──→ DEAD(人工处置)
MSG_EVENT(自有 PG,投递侧,见 design.md §5
PENDING ──投递确认──→ SENT
└── 失败退避 ──→ DEAD(DLQ)
```
- `PENDING``FAILED` 是处理中的状态,`FAILED` 继续退避重试并占用 FIFO 队头;`SUCCEEDED / SKIPPED / DEAD` 是终态,到达后队列方可推进。
- `DEAD``FAILED` 不是不可逆:经人工批准,指定错误类别的记录可重新置回 `PENDING` 处理(放行范围见 §9 / design.md §6.1)。重放处理的是当前状态,不恢复历史处理顺序。
- 航班变更、待发事件、处理终态与回填意图四类写入位于同一个 PG 事务,一起提交或一起回滚;信箱处理标记写在该事务之后,由与终态同行的回填意图驱动重试与超期补写
- **写事务边界**:只有处理器产出的业务型终态才与航班变更、待发事件、回填意图同处一个 PG 事务;非业务型终态(`MALFORMED / PROTOCOL / SKIPPED / EXHAUSTED`)与毒丸升级只写 `PROC_STATE` 一条记录,终态与回填意图由同一条语句写入。完整边界见 [design.md](design.md) §3.3 表 1
- 信箱处理标记写在该事务**之后**,由与终态同行的回填意图驱动重试与超期补写。
## 3. 处理阶段
**表 1 取报与处理**(收报、主泵、解码与身份):
| 阶段 | 输入 | 输出 | 幂等依据 | 失败处理 |
|---|---|---|---|---|
| 收报 | ID 区间内尚未入队的信箱行(谓词见 §5.1 | `PROC_STATE(PENDING)` | `MSG_ID` 主键 | 下轮扫描补入队 |
| 取队头 | 最小未完成 `MSG_ID` | 本批处理消息 | 表状态即队列 | 无队头则等待 |
| 解码与身份 | 信箱原文 | `DecodedMessage``IDENTITY_KEY` | 身份唯一约束 | 报文非法→`DEAD`处理能力不足→`FAILED` 退避 |
| 事务处理 | 当前完整状态 + 报文载荷 | 航班变更、终态、事件、回填意图 | 消息 ID + 业务身份 | 事务整体回滚重试 |
| 提交后回填 | 终态消息 | 信箱处理标记 | 标记单调(§11) | 终态事务内登记回填意图;失败退避重试,超期强制补写 |
| 投递 | `MSG_EVENT` | 投递目标接受确认 | `EVENT_ID` | 退避至 `DEAD`;至少一次 |
| 补偿作业 | 到期或已达 `NOW R` 的终态记录 | 追平标记 | 消息 ID(仅空标生效) | 指数退避,独立线程执行;`R` 覆盖退避 |
| 收报(发现 + 入队) | `ID > W` 的升序信箱行(谓词见 §5.1 | `PROC_STATE(PENDING)` 与推进后的 `(W, holeSince)` | `MSG_ID` 主键(冲突即已入队) | 下轮扫描补入队;整轮失败不动水位 |
| 取队头 | 最小未完成 `MSG_ID``PENDING``FAILED` 都占位) | 本轮唯一处理对象 | 表状态即队列 | 无队头则等待;`FAILED` 未到期则等待 |
| 解码与身份 | 信箱原文 | `DecodedMessage``IDENTITY_KEY` | 身份唯一约束(仅首次绑定) | 报文非法→`DEAD`解码能力不足→`FAILED` 退避 |
| 事务处理 | 当前完整状态 + 报文载荷 | 航班变更、终态、事件、回填意图 | 消息 ID + 业务身份 | 业务型事务整体回滚重试 |
领域决策逻辑不执行 I/O;Processor 作为事务协调器,只在持有 `PIPELINE_LOCK` 的事务内写入自有 PG 状态,不操作信箱与 Kafka。全部外部副作用发生在事务提交之后。
**表 2 回填**(终态之后的信箱标记):
| 阶段 | 输入 | 输出 | 幂等依据 | 失败处理 |
|---|---|---|---|---|
| 回填扫描(每 30 秒,**调度周期 ≠ 完成时限**) | 终态 ∧ 未确认标记 ∧ 未放弃 ∧(退避到期 ∨ 已达 `NOW R`) | 追平标记,或判定放弃自动重试 | 消息 ID(仅空标生效) | 退避 30 秒起步、封顶 15 分钟;越过 `R` 后每轮都试;**确认行不存在**或达尝试上限则停止自动重试(可人工恢复) |
主泵**不做**回填:终态与回填意图由同一条 UPDATE 落库,回填一律由扫描驱动(跨库写不能占用 FIFO 关键路径)。
**处理链路流程**`├─` 括起的是同一个 PG 事务):
```text
① 取报(InboxPoller,默认每秒一轮)
读游标 (W, holeSince) → 取 ID > W 的升序前 N 行 → 求连续上界 → 空洞判定
├─ PG 事务:insertIfAbsent(MSG_ID, RECEIVED_AT) × k + 写回 (W, holeSince) → 提交
└─ 未入队的行留到下一轮(每轮最多解决一个空洞)
② 调度(Pump,单线程)
headUnfinished()(最小未完成 MSG_ID
├─ 队头为 FAILED 且已毒丸 → DEAD(EXHAUSTED)(单语句,不取锁)
├─ FAILED 未到期 → 等待
└─ 可执行 → 记 PROCESSING_STARTED_AT(仅首次)→ processOne()
③ 处理(MessageProcessor
读原文(共享 MySQL)→ 解码 → 身份绑定(独立单语句)
└─ PG 事务(持 PIPELINE_LOCK):航班变更 + 待发事件 + 终态 + 回填意图 → 提交
④ 回填(终态之后,跨库两次单写,无事务)
写信箱处理标记 → 记 BACKFILL_AT;失败则回填意图留在该行上
└─ 回填扫描(每 30 秒)按表 2 的谓词重试,越过 R 后强制补写
```
领域决策逻辑不执行 I/O。Processor 作为事务协调器:业务型终态在持有 `PIPELINE_LOCK` 的事务内写入自有 PG;非业务型终态与毒丸升级只写 `PROC_STATE` 单条记录。所有外部副作用(读信箱原文、写信箱标记)都发生在 PG 事务边界之外,跨库不共享事务。
## 4. 故障恢复
@@ -66,10 +108,14 @@ MSG_EVENT PENDING ──→ SENT
| 中断位置 | 重启后的判定 | 恢复动作 |
|---|---|---|
| 已落信、未入队 | 信箱行位于应扫描的 ID 范围且 PG 无记录(不以处理标记为判据) | 重扫补建入队记录 |
| 水位卡在空洞 | `INBOX_CURSOR.HOLE_SINCE` 有值且未超过 `max-commit-delay` | 等待;超期后放行空洞本身并继续推进(§5.1) |
| 事务执行中 | PG 无该消息终态 | 事务整体回滚,按 `PENDING` 重新处理 |
| 事务已提交、标记未写 | 终态行仍持有回填意图(`BACKFILL_NEXT_AT` 非空) | 仅补写标记;业务处理结果保持不变 |
| 回填意图登记失败(与业务同事务) | 事务未提交 | 同"事务执行中",不构成独立窗口 |
| 标记写入中途 | 标记仍为空 | 重新写入;重复写入同一值无副作用 |
| 兼容入口已入队、水位未追平 | PG 已有该 ID 的记录 | 轮询读到该行时主键幂等,水位照常推进;顺序风险见 §5.1/§12 |
| 回填时信箱行已不存在 | 写入 0 行且信箱行不存在 | **立即放弃自动重试**(原因 `MISSING_ROW`)并告警;**放弃 ≠ 标记已确认**,仍需人工对账(§5.2) |
| `RECEIVED_AT` 为 NULL | 终态行无标记且无超期判据 | 仅由退避重试保证;【目标】以入队时间兜底(§5.2/§12) |
| 投递目标已接受、`SENT` 未置 | 事件仍 `PENDING` | 允许重发,消费方按事件身份去重 |
两个窗口已随"终态与回填意图同体同行"消除(口径与 design.md §10「事务与外部副作用」一致),不再是缺口:
@@ -77,53 +123,119 @@ MSG_EVENT PENDING ──→ SENT
1. **提交后回填前崩溃**:回填意图与处理终态是同一条记录的同一次写入、同一个事务;重启后扫描按该意图继续补写。
2. **回填意图二次落账失败**:不存在第二处落账——意图就在终态行上。
仍属未闭环的回填本身在跨库单写期间的持续失败由退避重试与 §5.2 的超期期限 `R` 兜底(登记于 design.md §10 与 user-stories.md US-09
仍属未闭环的有两处:回填本身在跨库单写期间的持续失败由退避重试与 §5.2 的超期期限 `R` 兜底;以及 §5.2 的"信箱行已不存在"没有终态(登记于 §12)。两者均登记在 design.md §10 与 user-stories.md US-09。
## 5. 消费水位、超期补写与历史积压
### 5.1 消费水位 W
### 5.1 消费水位 W 与空洞
水位 `W` 定义为信箱 ID 的连续上界从最小 ID 到 `W` 的区间已全部读入自有 PG,无空洞;遇到空洞即停止推进。`W` 只随新 ID 的成功入队推进,不依赖处理完成或标记回写。
**定义**水位 `W` 信箱 ID 的连续上界——从最小 ID 到 `W` 的区间已全部读入自有 PG,无空洞;遇到空洞即停止推进。`W` 只随新 ID 的成功入队推进(永久空洞放行是唯一例外),不依赖处理完成或标记回写。
水位的有效性依赖两条需要库方确认的承诺(合计为 Q2):
**收报流程**(每轮,谓词不以处理标记为条件):
1. 读取游标 `(W, holeSince)`;信箱不可读时记日志、等下一轮,**不动水位**——这属于基础设施失败,不能当成"没有新消息"。
2.`CMINMSGS_ID > W` 的升序前 `claim-batch` 行(`ORDER BY ID ASC LIMIT N`)。
3.`W+1` 起逐 1 数,求连续上界;遇到第一个缺号即停止计数。
4. 空洞判定(仅当本批存在"缺号之后的行"时才可能成立):
- 缺号首次被观测到 → `holeSince = now`
- `now holeSince < max-commit-delay` → 水位停在缺号前,**本批缺号之后的行一律不入队**。没有这条规则,晚提交的较小 ID 会排到它们后面,破坏 FIFO;
- `now holeSince ≥ max-commit-delay` → 判定为永久空洞,水位放行到"缺号后第一行 − 1"`holeSince` 清空。放行**只跳过空洞本身,不越过任何已存在的行**。
5. 在同一个 PG 事务内:对水位以内的每一行 `insertIfAbsent(MSG_ID, RECEIVED_AT)`,并写回 `(W, holeSince)`;主键冲突表示已入队(重复扫描与兼容入口并发都安全),不计入也不报错。
6. 提交。本批中因空洞或批次上限未入队的行留待下一轮——**每轮最多解决一个空洞**。
**空洞计时跨重启保留**`holeSince` 落在 `INBOX_CURSOR.HOLE_SINCE`,进程重启不丢失。旧空洞补齐后出现的新空洞,从新观测时刻重新计时,不继承旧等待时间。
**代价(必须接受并观测)**:水位遇空洞即停意味着**该空洞之后的所有消息最多要等 `max-commit-delay` 才能入队**。自增回滚等会在 ID 序列中留下永久空位,因此每出现一个永久空位就是一次等长的入队停摆;空位频繁时有效吞吐按比例下降。运行期必须观测永久空洞计数与水位滞后(§5.3 验收口径、OPS-2)。
**水位的有效性依赖 Q2 的两条上游承诺**
1. **ID 单调**:信箱 ID 按提交顺序分配,晚提交的较小 ID 不参与。承诺缺失时,`W` 只能作为快路径的提示,不能证明该区间收齐。
2. **最大提交时延**:上游提交到 ID 可见的最长时间。补偿扫描窗口(`NOW 最大提交时延`)的宽度由此确定
2. **ID 分配 → 事务可见时延上界**:从"库方为该报文分配 ID"到"该 ID 对其他事务可见"的最长时间。它同时决定空洞老化阈值与补偿扫描窗口宽度
取值可参照 SIS:报文在 CIIMS 的 `Expiry` 480 分钟断连 120480 分钟按 Level 2 处理,CIIMS 保证报文在成功接收或过期前按序保存`SIS_AODB_RMS-V0.1.md` §2.4.2.3.3、§3.16、§3.17、§5.2)。
**这条时延必须由库方直接给出,不能由 SIS 的报文有效期推导。** SIS 的 `Expiry` = 480 分钟断连 120480 分钟按 Level 2 处理(`SIS_AODB_RMS-V0.1.md` §2.4.2.3.3、§3.16、§3.17、§5.2描述的是**报文保留与传输恢复**,与"共享 MySQL 里 ID 分配后多久对读事务可见"不是同一个量,二者之间没有推导关系。因此当前默认的 `max-commit-delay = 5 分钟`**缺少依据**,只是占位假定值,属于上线门槛([design.md](design.md) §7/§10)。报文有效期只能作为**原文保留期 / 重放窗口**(§6、§9)的参照,不能反过来给本阈值背书
日常执行方式:快路径从 `ID > W` 起按升序有限批次扫描;另按周期对窗口内可能迟到或空洞的行做补偿扫描。需要区分:水位表示"读取进度",与"已处理标记"是两个事实,不能互相替代。
**两条扫描路径(发现机制)**
**空洞老化**:自增回滚等会在 ID 序列中留下永久空位,而 Q2 只承诺 ID 单调、不承诺无空洞。W+1 处的空洞持续超过最大提交时延仍未被补齐时,即判定为永久空洞并放行水位(实现取 `msgx.pipeline.max-commit-delay`);放行只跳过空洞本身,不越过任何已存在的行。没有这条规则,水位会永久停摆于第一个空位,其后的行再也不会入队。
| 路径 | 目的 | 谓词 | 状态 |
|---|---|---|---|
| 快路径(日常) | 发现水位之后的新消息 | `ID > W ORDER BY ID ASC LIMIT claim-batch` | 【现状】已实现 |
| 迟到检测(只读,阶段 0) | 复查**被放行的空洞 ID** 是否后来真的出现在信箱 | 进程内监视队列 + `existingIds` 批量存在性检查(`late-detect-period`) | 【现状】已实现:只计数与告警,**不补入队** |
| 补偿扫描(阶段 1) | 发现"提交晚于水位推进"的迟到行并**安全补入队** | 按周期重扫 `W` 之前一个窗口(宽度由最大提交时延决定)内的 ID 区间 | 【缺口】未实现 |
两条扫描路径都按 **ID 区间**取行,不以处理标记为扫描谓词;标记只用于回填与库方清除,不参与消息发现。否则已入队但尚未回填的行会永久占据批次,这正是 US-01 条目 3 与 Q2 要求排除的场景
【缺口】在补偿扫描交付之前,"较小 ID 迟提交"**没有补入队机制**:快路径只读 `ID > W`,一旦水位越过某个 ID,该 ID 之后到达的消息永远不会被发现。阶段 0 的**只读迟到检测**已实现——监视"被放行的空洞 ID"是否后来真的出现,命中即计数(`msgx.pipeline.late_arrival.detected.total`)并告警,把原先的**静默丢失**变成可观测事实;但它**不补入队**,因此当前设计仍**不能对外声明"迟到不越序"**[design.md](design.md) §8/§10)。检测是进程内尽力而为(重启丢失监视),且只覆盖"曾被放行过的空洞 ID"——这已足够:水位只会越过连续存在的行与被判定为空洞的 ID
**水位与其余两条写路径的关系**
- **兼容 HTTP 入口**`POST /cminmsgs/send`):它直接写 `PROC_STATE`、不读不推水位,登记的行因此**超出水位**;主泵只领 `msgId ≤ W`,所以在水位追平前不会被处理——顺序不受影响,代价是延迟到追平,且**必须让收报轮询运行**(收报侧事实由 `InboxPollerTest` 固定,端到端由 `PipelineSmokeTest` 守住)。
- **回填**:不回写水位。水位不是处理标记,两者互不替代。
**单实例前提**:信箱读取不加锁,水位是单行覆盖写。本设计只在单活动实例下成立([architecture.md](architecture.md) §5、D2);多实例并发收报会让水位互相覆盖(覆盖回退只会造成重复扫描,不会丢消息,但空洞计时会失真),必须先有实例级排他。
**首次启动**:游标初值为 `(W=0, holeSince=NULL)`,第一轮会重读信箱中全部现存行;重复登记由主键幂等挡住,因此历史上"已入队但未回填"造成的漏读会一并补齐。
**切流播种(显式、一次性)**:若信箱已有存量(典型情况是最老分区已被清除、`MIN(ID)` 远大于 1),从 `W=0` 启动会先把 `ID=1` 判成空洞、白等一个老化窗口,同时也意味着要重新处理保留期内的全部存量。是否跳过存量属于**切流决策**,因此代码不做默认选择:只有显式配置 `msgx.pipeline.cutover-watermark` 才播种,取值 `min`(读现存全部)/ `zero`(从 0 按空洞规则)/ `max`(跳过可见存量)/ 具体 ID。升级实例(已有水位或已有处理记录)**拒绝重新播种**,重新切流必须是显式操作;播种事实与水位同语句落库(`INBOX_CURSOR.SEEDED_AT`),且该列为 NULL **不等于**"从未消费"(已有库新增列后同样为 NULL)。
### 5.2 超期标记补写
所有处理终态都只依赖消息 ID 执行回填;报文残缺或缺少 META 不妨碍正常回填。本节规定回填通道长期失败时的超期兜底,避免终态信箱行无限期保持空标记。处理规则为:
所有处理终态都只依赖消息 ID 执行回填;报文残缺或缺少 META 不妨碍正常回填。本节规定回填的三种结果与"长期失败"时的超期兜底
**已达 `PROC_STATE` 终态、且接收时间超期(信箱 `DATE_RECEIVED`,落库为 `PROC_STATE.RECEIVED_AT`;判据 `RECEIVED_AT < NOW R`)仍无标记的信箱行,由回填通道补写一个库方认可的"已处理"类标记。**
**处理规则**已达 `PROC_STATE` 终态、且接收时间超期(信箱 `DATE_RECEIVED`,落库为 `PROC_STATE.RECEIVED_AT`;判据 `RECEIVED_AT < NOW R`)仍无标记的信箱行,由回填通道补写一个库方认可的"已处理"类标记。
- 期限 `R` 必须不小于(人工重放期限 + 人工处置期限)之和;两项期限的取值口径统一见 Q6,确认前不得下调 `R`。提前补写会使仍可重放的消息先被库方清除、失去原件资格。
**扫描谓词**(与实现一一对应):
```text
STATE ∈ {SUCCEEDED, SKIPPED, DEAD} -- 终态
AND BACKFILL_AT IS NULL -- 标记尚未确认
AND ( BACKFILL_NEXT_AT IS NULL -- 异常兜底:终态行没有退避时间
OR BACKFILL_NEXT_AT ≤ NOW -- 退避到期
OR ( RECEIVED_AT IS NOT NULL -- 接收时间可用
AND RECEIVED_AT < NOW R ) ) -- 已达超期期限,覆盖退避
ORDER BY MSG_ID ASC LIMIT backfill-batch
```
**回填的可能结果**
| 结果 | 判定 | 处置 |
|---|---|---|
| 写入成功 | 标记为空、写入 1 行 | 记 `BACKFILL_AT`,不再重试 |
| 早已有标记 | 写入 0 行且信箱行存在 | **视为成功**,不覆盖已有值,记 `BACKFILL_AT` |
| 信箱行不存在 | 写入 0 行且信箱行不存在 | **立即放弃自动重试**(原因 `MISSING_ROW`)并告警。这是确定性结论,重试不会改变结果 |
| 暂时性故障达上限 | 超时/连接失败累计达到 `backfill-max-attempts` | **停止自动重试**(原因 `MAX_ATTEMPTS`)并告警;保留人工恢复能力 |
**放弃 ≠ 标记已确认**:放弃行不写 `BACKFILL_AT`,因此**不满足** §6"边界内全部行已打标"的清除前提,库方不得据此清除。放弃行可经人工恢复(清标记后重排一次回填)。
**约束**
- 期限 `R` 的语义是"回填长期失败时的强制补写上限"。**`R` 在任何清除语义下都不保护重放窗口**——因为回填不等 `R`:终态落库后由扫描尽快补写(§3 表 2),`R` 只在"回填持续失败"时把补写**提前**,从不推迟补写。
- 重放窗口的保护只能来自下面两条路之一,且 Q7/Q9 的确认是**阻塞性前提**,不是"调大 `R`"就能覆盖的参数问题:
- **约定保留期(目标前提)**:库方按"标记 + 保留期 `R_keep`"清除(打标本身不触发清除,§6),且 `R_keep ≥ 人工重放期限 + 人工处置期限`。此时 `R` 只需满足 `R ≤ R_keep`,与重放窗口无关。
- **另设原文保留机制**:若库方的清除语义是"一旦打标即可清除"(未确认),则**增大 `R` 无效**——当天进入终态、当天打标的消息会被当天清除,即使 `R = 30 天`。此时必须另行约定保留期,或引入独立的原文保留通道(例如由库方把原文归档到 `CMINMSGS_HST` 后供重放读取);该通道**尚未设计**(§12 G10)。
- 撤回原先"打标即清除语义下需 `R ≥ 重放期限 + 处置期限`"的表述:它是错误推论,因为**回填由扫描驱动、不等 `R`**——打标时刻与 `R` 本来就是解耦的,`R` 无法推迟打标。
- 中间态(`PENDING` / `FAILED`)不适用本规则:处理未完成时不打标,也不允许被清除。
- 补写值仅限于库方认可的 legacy 值集(Q7);死信、业务重复等内部原因记录在 `PROC_STATE` 与审计日志,不在信箱新增枚举。
- 补写只针对空标记;已有值不回撤、不覆盖,重复执行无副作用。
- 判据不依赖独立待办表:回填意图与处理终态同行(§2/§4),扫描条件为「终态 + 未确认标记 + (已到期 或 接收时间早于 `NOW R`)」——`R` 是覆盖退避的硬期限,保证该条件在有限时间内必然被处理
- `RECEIVED_AT` 为 NULL(上游未写 `DATE_RECEIVED`)时,超期分支不成立,`R` 兜底**不生效**;该行只能靠退避(封顶 15 分钟)反复重试。【目标】接收时间缺失时以入队时间兜底
- 判据不依赖独立待办表:回填意图与处理终态同行(§2/§4)。
- **越过 `R` 之后**:超期分支一旦成立便恒成立,该行**每轮扫描都会被重试**,退避被覆盖。不过永久失败不会无限重试:**确认行不存在**立即放弃,暂时性故障到 `backfill-max-attempts` 后也停止自动重试(两类都告警,且可人工恢复)。
- **扫描周期 ≠ 完成时限**:30 秒只是**调度周期**。`JobRunner` 串行执行 `backfill.sweep` 与历史作业后才等待 30 秒,批次积压、单行调用超时与历史作业耗时长都会延长实际回填延迟;只有在"正常无积压"时实际延迟才近似等于扫描周期。因此任何"标记延迟 ≤ 30 秒"的表述都不成立,**回填完成目标是独立指标**(§5.3 验收口径、OPS-2),需要硬时限时必须另行补齐容量、作业隔离与故障条件。
- **回填饥饿(已修)**:扫描改为**公平轮转**(`ORDER BY BACKFILL_ATTEMPTS ASC, MSG_ID ASC`)并排除已放弃的行,因此"最旧的一批永久失败行占满批次、后面的记录永远轮不到"不再成立。历史缺口见 §12 G9(已关闭)。
该规则同时保证件事成立:收报表不被永不回填的记录占满(Q2 中"积压挡批"场景)、§6 的清除边界可以达到、重放期限与信箱保留期具备可核验的下限关系(§9)
该规则保证件事成立:收报表不被永不回填的记录占满(Q2 中"积压挡批"场景)、`R``R_keep` 的约束关系可核验(§5.2/§6/§9)。**它仍不保证"§6 的清除边界能在有限时间内达到"**:放弃行永远不满足"边界内全部行已打标",而"打标即清除"语义(§12 G10)会让该前提根本不成立
### 5.3 历史积压消息的处理方案
"历史积压"指信箱中成规模的未处理存量:上线前遗留、停机期间累积、或消量未完成的批次。处理方案分三步,并约束四条行为红线。
**第一步:摸底。** 处理动作开始前,先确定积压的范围与构成:ID 区间、条数、时间跨度、报文类型分布(`SCHD-DNLD/RESP/ADFT``FLOP-*``FDEL` 等),并与库方确认其中哪些仍需业务处理、哪些按约定跳过(跳过基调只在该批确认时成立,授权与留痕见 Q12,标记值集见 Q7)。
**第一步:摸底。** 处理动作开始前,先确定积压的范围与构成:ID 区间、条数、时间跨度、报文类型分布(`SCHD-DNLD/RESP/ADFT``FLOP-*``FDEL`,逐类清单依据见 Q8),并与库方确认其中哪些仍需业务处理、哪些按约定跳过(跳过基调只在该批确认时成立,授权与留痕见 Q12,标记值集见 Q7)。
**第二步:入队,后处理。** 两阶段执行,禁止边收边处理一辆长队
**第二步:入队与处理的顺序关系。** 顺序由 `MSG_ID` 决定,**不由"先入队后处理"这种执行方式决定**
1. 消化阶段只做入队:`InboxPoller` 按升序有限批次将历史行全部建为 `PROC_STATE(PENDING)`,水位随之推到积压末端(连续推进以 §5.1 的 Q2 承诺为前提)。此阶段只写自有 PG,不触碰信箱标记。
2. 入队完成后交给主泵按最小未完成 ID 顺序消化。顺序阈限与日常完全相同:不加速、不分流、不走旁路。
1. 入队:`InboxPoller` 按升序有限批次将历史行全部建为 `PROC_STATE(PENDING)`,水位随之推到积压末端(连续推进以 §5.1 的 Q2 承诺为前提)。此阶段只写自有 PG,不触碰信箱标记。
2. 处理:主泵按最小未完成 ID 顺序消化。顺序阈限与日常完全相同:不加速、不分流、不走旁路。积压期间到达的实时消息 ID 更大,自然排在积压之后。
这样做的理由:入队廉价、处理昂贵;分开后入队阶段的中断恢复只涉及重扫(§4 第一行),水位也能尽早到达完整上界,快路径与补扫窗口立即生效
**入队与处理可以并发**:两者由不同线程驱动,FIFO 由"队头取最小未完成 `MSG_ID`"保证,并发不会破坏顺序。因此"先入队后处理"是**可选的运维规程**(例如先手动跑一段只收报的窗口,便于提前把水位推到完整上界),**不是正确性前提**,系统也不需要提供"只入队"模式
这样安排的收益:入队廉价、处理昂贵;即使并发,入队阶段的中断恢复也只涉及重扫(§4 第一行),水位能尽早到达完整上界,快路径与补偿扫描窗口立即生效。
**第三步:等待过程中的四个边界。**
@@ -132,7 +244,7 @@ MSG_EVENT PENDING ──→ SENT
- 报文类型的处理方式不变:旧的全量日计划(`SCHD-DNLD`)按最新一份**合并**即可收敛(缺席不删、缺失字段保留;字段级冲突见 flight-state.md §3.1 与 Q13),但仍逐条执行;动态增量(`FLOP-*` / `FDEL` / `ADFT`)持有时序语义,必须逐条。
- 可放弃但必须留痕:摸底批中经库方确认"不再处理"的行,处置方式为——`PROC_STATE``SKIPPED` 并记录跳过原因,到达终态后走 §5.2 通道补写标记;本方案不支持任何"整段 DELETE"的快速通道。
**验收口径**:积压消化期间持续输出三项指标——剩余积压条数、最老未处理信龄、预计消化时长;期间不允许出现 FIFO 越序、身份去重失效或头行滞留超时未告警。红线依据:architecture.md §5(消息严格 FIFO)、单写者约束。三项指标并入 user-stories.md OPS-2 的可观测性验收。
**验收口径**:积压消化期间持续输出三项指标——剩余积压条数`unfinished`)、最老未处理信龄(`oldestUnprocessedSeconds`)、**预计消化时长**(定义:`剩余积压条数 ÷ 近 N 分钟实际处理速率`,速率样本窗口 N 需在实现时固定;若实现不提供该估算,本条降级为"由运维按前两项指标自行推算")。另需观测水位滞后与永久空洞计数(§5.1)。期间不允许出现 FIFO 越序、身份去重失效或头行滞留超时未告警。红线依据:architecture.md §5(消息严格 FIFO)、单写者约束。以上指标并入 user-stories.md OPS-2 的可观测性验收。
## 6. 信箱数据清除
@@ -143,7 +255,7 @@ MSG_EVENT PENDING ──→ SENT
**方案 A:按日分区(首选)**,适用于库方可以为表增加分区的场合。
1. `CMINMSGS``DATE_RECEIVED` 建立日粒度 RANGE 分区;
2. 某分区到达保留期时,确认该分区全部行已持有处理标记(§5.2 保证该条件在有限时间内满足);
2. 某分区到达保留期时,确认该分区全部行已持有处理标记(§5.2 只在回填可持续成功时保证该条件;存在回填饥饿或"打标即清除"语义时见 §12 G9/G10);
3. 将该分区复制入历史表:`INSERT INTO CMINMSGS_HST SELECT``NOT EXISTS` 判重);
4. `TRUNCATE / DROP PARTITION` 执行清除:DDL 级操作,无行锁竞争,页外大字段(`CMINMSGS_CLOB_MSG` 列)整块释放;碎片与 binlog 影响待现场验证。
@@ -162,9 +274,11 @@ MSG_EVENT PENDING ──→ SENT
方案共同前提:
- 执行清除时,边界内不存在未打标记的行;未达终态的行顺延至处理完成后清除(§5.2 仅对终态行补标);
- 时间比较与换算统一采用机场时区 Asia/Shanghai 及明确的类型转换(口径同 user-stories.md §6);
- 清除保留期 `R_keep` 不小于 max(回填重试上限、重放期限、审计期限);且 §5.2 的 `R` 不大于 `R_keep`,否则尚在重放窗口内的消息会先于重放被清除;
- **方案要求库方按"标记 + 保留期"清除**(打标本身不触发清除):处理标记只是可清除的**必要条件**,触发条件是"到达保留期 `R_keep`" ∧ "边界内全部行已持有处理标记"。该语义属 Q7/Q9 确认范围;若确认结果是"打标即可清除",则**本节前提不成立,且增大 `R` 无法补救**(原因见 §5.2:回填不等 `R`,打标时刻与 `R` 解耦),必须另行约定保留期或引入独立原文保留通道(§12 G10)。
- 执行清除时,边界内不存在未打标记的行;未达终态的行顺延至处理完成后清除(§5.2 仅对终态行补标)。
- 时间比较与换算统一采用机场时区 Asia/Shanghai 及明确的类型转换(口径同 user-stories.md §6)。
- 保留期 `R_keep` 的下界(**本节是唯一表述处**):`R_keep ≥ max(人工重放期限 + 人工处置期限, 审计期限, 回填重试上限)`——**仅在"标记 + 保留期"清除语义下**,这是"重放窗口内原文仍在"的**唯一保证来源**。报文在 CIIMS 的 `Expiry`480 分钟量级,`SIS_AODB_RMS-V0.1.md` §3.16)可作为原文保留期的参照,但它是报文有效期,不等于本处所需的保留期,也不能用于推导 §5.1 的可见性时延。
- 超期补写期限 `R` 的约束(仅 `R ≤ R_keep`**唯一见 §5.2**`R` 不参与重放窗口保护,本节不重复。
- legacy 现役按接收时间超过 1 天即归档并删除 `CMINMSGS`(legacy 行为基线 §3.6);若沿用该窗口,则与上一行的保留期下限冲突,须在 Q9 中与库方一并确认;
- 信箱 ID 全程不断链:方案 A 天然满足;方案 B 依赖 `AUTO_INCREMENT` 种子,种子缺失时新 ID 与旧记录主键冲突,水位随之失效。
@@ -194,32 +308,56 @@ MSG_EVENT PENDING ──→ SENT
## 9. 重放与原文可用性
- 可重放的错误类别为 `CODEC_ERROR / UNSUPPORTED / INFRA / EXHAUSTED` 四类(以 design.md §6.1 为准);重放处理当前状态,不恢复历史顺序。
- 重放的前提是信箱原文仍可读取:`R_keep` 必须覆盖重放期限(§6),属于硬约束
- 重放的前提是信箱原文仍可读取。**该前提只在"标记 + 保留期"清除语义下才能由配置满足**`R_keep` 必须覆盖(人工重放期限 + 人工处置期限)(§6)。若库方是"打标即清除",则需独立的原文保留通道(§5.2、§12 G10)——**增大 `R` 无效**
- 确认信箱行或原文已缺失时,按 `MALFORMED → DEAD` 处理;共享库超时、连接失败等暂时性读取异常按 `INFRA` 退避重试,不得伪装成原文缺失。若原文缺失源于库方违反保留契约提前清除,按契约违例走运维追责通道(契约条款见 Q7/Q9);该消息的死信处置本身不变。
- **回填阶段的"信箱行缺失"是另一回事**(§5.2):它表示"本系统已经处理完、却找不到要标记的那一行"。此时不改变处理终态,也不重新处理业务;需要区分"该 ID 本就属于被判定的永久空洞"与"行被提前清除"。【缺口】当前无终态与告警,见 §12。
## 10. 开放问题索引
| 编号 | 待确认事项 | 本文相关章节 |
|---|---|---|
| Q2 | 信箱 ID 单调承诺;最大提交时延;空洞与迟到处理 | §5.1、§6 |
| Q6 | 重放 deadline 与人工处置期限的取值(决定 §5.2 的 `R` | §5.2、§9 |
| Q7 | 处理标记值集与写权限;原文保留期;处理时间语义 | §5.2、§6 |
| Q8 | 逐类覆盖清单(积压摸底的类型分布依据) | §5.3 |
| Q9 | 清除执行方与 DDL 授权;方案 A / B 选型 | §6 |
| Q10 | 出站消费方、ACK 列语义、出站清理与去重契约 | §7 |
| Q11 | 上游 `SEQN` 重置周期与业务身份的日期边界 | design.md §2.2 |
| Q12 | 历史积压批次“不再处理”的确认主体、审批留痕与跳过值集 | §5.3 |
| Q13 | 日计划缺失可选字段的删除语义(与 SIS §3.16 注释 4 冲突) | §5.3、flight-state.md §3.1 |
| Q14 | 主/共享删除顺序与 EROR 回报义务 | §7、flight-state.md §3.3 |
| 编号 | 待确认事项 | 当前假定值 | 阻塞谁 | 风险 | 本文相关章节 |
|---|---|---|---|---|---|
| Q2 | 信箱 ID 单调承诺;**ID 分配 → 事务可见时延上界**;空洞与迟到处理 | 时延按 5 分钟(`max-commit-delay`,**缺少依据的占位值**;不可由 SIS 报文 `Expiry` 推导) | 快路径的发现完整性声明、空洞老化阈值、补偿扫描窗口 | 迟到小 ID 永久不被发现;空洞被误判为永久 | §5.1、§6 |
| Q6 | 重放 deadline 与人工处置期限的取值(唯一作用是决定 `R_keep` 下界) | R = 30 天;重放/处置期限未定 | `R_keep` 的最终取值 | `R_keep` 不足导致重放窗口内原文被清除;**若清除语义为"打标即清除",则任何取值都无效**(§5.2 | §5.2、§6、§9 |
| Q7 | 处理标记值集与写权限;原文保留期;处理时间语义 | 写入 `PROCESSED` | 回填值集、原文保留期下限 | 值集不被库方接受;保留期短于重放窗口 | §5.2、§6 |
| Q8 | 逐类覆盖清单(积压摸底的类型分布依据) | — | 积压摸底与逐类矩阵 | 类型漏项导致积压处置误判 | §5.3 |
| Q9 | 清除执行方与 DDL 授权;方案 A / B 选型 | 首选方案 A | `R_keep` 与清除边界 | 方案不可执行;ID 断链 | §6 |
| Q10 | 出站消费方、ACK 列语义、出站清理与去重契约 | — | 出站信箱 | 重复写入、无法清理 | §7 |
| Q11 | 上游 `SEQN` 重置周期与业务身份的日期边界 | 不含日期边界 | 身份算法 | 跨周期误判重复 | design.md §2.2 |
| Q12 | 历史积压批次“不再处理”的确认主体、审批留痕与跳过值集 | — | 积压跳过处置 | 无授权跳过或删除 | §5.3 |
| Q13 | 日计划缺失可选字段的删除语义(与 SIS §3.16 注释 4 冲突) | 保留未携带字段 | 快照合并 | 与 SIS 语义冲突 | §5.3、flight-state.md §3.1 |
| Q14 | 主/共享删除顺序与 EROR 回报义务 | 自动幂等级联 | 出站事件类型 | 上游状态不一致 | §7、flight-state.md §3.3 |
Q2、Q6Q14 的完整登记见 [user-stories.md](user-stories.md) §6 的 Q 表。
## 11. 不变量
- 五个事实互不替代:入队不引用信箱标记,回填不引用投递,投递不引用回填。
- 发现与处理互不阻塞:收报只看 `ID > W`,终态而未回填的行不阻断后续消息的发现;处理状态不改变扫描谓词。
- 队头唯一:任一时刻只有一个可执行队头,`FAILED` 未退避到期时后续消息不得越过。
- 只领取已发现的行:主泵只领 `MSG_ID ≤ W`;水位之外的行只可能来自兼容入口,必须等水位追平后按序处理(否则会越过尚未入队的较小 ID)。
- 处理终态不可逆:已提交的 `SUCCEEDED` 不因回填或投递失败回改。
- 处理标记单调:任何路径只将空标写为已处理,不回撤、不覆盖。
- 回填只针对终态:`PENDING / FAILED` 永不写信箱标记。
- 一信一行:每条信箱行在 `PROC_STATE` 至多一条记录(`MSG_ID` 主键);同一业务身份至多绑定一条有效处理记录。
- 水位不越过任何已存在的行;遇空洞即停,只有 §5.1 判定为永久空洞时才放行。
- 写入水位与入队同事务:不允许出现"水位已推进、消息未入队"的持久化状态。
- 执行清除前,边界内全部行已持有处理标记;物理删除仅发生在归档成功之后(追加写入与分区留档均构成归档成功)。
- 信箱 ID 全局单调、不断链:以 Q2 的 ID 单调承诺与 §6 各方案前提成立为条件(方案 B 依赖 `AUTO_INCREMENT` 种子)。
- 对外投递按至少一次设计;端到端恰好一次不在交付范围内。
## 12. 与当前实现的差异
本节是本文范围内的缺口清单,与 [design.md](design.md) §10 对应。正文描述目标设计;下列条目尚未交付,不得作为已具备能力引用。
| 编号 | 缺口 | 影响 | 相关章节 |
|---|---|---|---|
| G1 | **窗口补偿扫描未实现**(阶段 0 已交付) | 快路径只读 `ID > W`;水位越过某个 ID 后,该 ID 之后到达的迟到消息永远不会被补入队。**已实现只读迟到检测**(`msgx.pipeline.late_arrival.detected.total` + 告警),把静默丢失变为可观测;**补入队(阶段 1)仍未实现**,因此不能声明"迟到不越序" | §5.1、§4 |
| G2 | ~~**兼容入口不参与水位**~~(**已修**) | 兼容入口登记的行超出水位,主泵只领 `msgId ≤ W`,因此不会越过尚未入队的较小 ID;代价是需等水位追平(且收报轮询必须运行)。收报侧事实与端到端验收均有用例 | §5.1、§4 |
| G3 | ~~**回填"信箱行缺失"无终态**~~(**已修 · V4**) | 运行时确认 `MISSING` 即放弃自动重试(原因 `MISSING_ROW`)并告警;暂时性故障达 `backfill-max-attempts` 后也停止自动重试;两者均可人工恢复,且**不等于**标记已确认 | §5.2、§4、§9 |
| G4 | **`RECEIVED_AT` 为 NULL 时 R 兜底失效** | 该行只能靠退避重试;V2 只对存量行用 `UPDATED_AT` 兜底,新行没有 | §5.2、§4 |
| G5 | **毒丸升级不在 `MessageLifecycleGate` 内** | 与人工重放并发时,存在"先标 DEAD 并写标记、再被重放拨回 `PENDING`"的窗口。**注意**:回填与重放**共用同一个 gate 单例**这一点已由装配断言守住(`PipelineSmokeTest`),但毒丸路径仍在 gate 之外 | §2、§5.3design.md §3.2/§6.1 |
| G6 | **"预计消化时长"指标无实现定义** | §5.3 验收口径的第三项指标尚无算法与实现(已降级为可推算项) | §5.3 |
| G7 | **`max-commit-delay` 默认值缺少依据** | 库方尚未给出"ID 分配 → 事务可见"的时延上界;5 分钟只是占位假定值,**不能由 SIS 报文 `Expiry` 推导**(§5.1)。取值无依据时可能大量误判永久空洞 | §5.1、§6 |
| G8 | **Q 编号条款尚未与库方书面确认** | 当前取值均为假定,见 §10 | §10 |
| G9 | ~~**回填扫描无公平轮转、无永久失败隔离**~~**已修 · V4** | 改用 `ORDER BY BACKFILL_ATTEMPTS, MSG_ID` 公平轮转并排除已放弃行,新行不会再被最旧的一批永久失败行饿死 | §5.2、§3、§6 |
| G10 | **"打标即清除"语义下没有原文保留通道** | 若库方清除由打标触发,则**回填成功后当天打标**会让原文当天即可被清除,增大 `R` 无效 → 重放窗口失去保护。需另行约定保留期,或引入独立原文保留通道(如库方归档 `CMINMSGS_HST` 供重放读取);该通道尚未设计 | §5.2、§6、§9 |