diff --git a/docs/design.md b/docs/design.md index 1e396a8..b323e4c 100644 --- a/docs/design.md +++ b/docs/design.md @@ -24,7 +24,7 @@ |---|---|---|---|---|---| | `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`(超期补写期限) | 回填长期失败时的强制补写上限(按接收时间计) | 判据用 `PROC_STATE.RECEIVED_AT` | `msgx.pipeline.overdue-backfill` | 30 天 | 仅 `R ≤ R_keep`;**不保护重放窗口**(完整论证唯一见 [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 | @@ -40,7 +40,7 @@ | 记录 | 用途 | 关键约束 | |---|---|---| -| `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 的超期兜底不生效。 | +| `PROC_STATE` | 入站消息的处理状态、身份、重试次数、错误原因与回填事实 | `MSG_ID = CMINMSGS_ID` 主键防止重复入队;`IDENTITY_KEY` 唯一约束防止业务重复;按最小未完成消息 ID 取队头;`PROCESSING_STARTED_AT` 为 HOL 计时起点(口径见 §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-14(US-14 两类映射的存储落点未定)。 | @@ -62,7 +62,7 @@ SNDR | TYPE | STYP | SEQN 接收时只按信箱 ID 去重;解码后才首次绑定业务身份。重试保留原有绑定,不能把自己判为重复消息。身份被另一条记录占用时,当前消息转为 `SKIPPED`,记录 `duplicate-of:`。是否加入日期边界取决于上游序号重置周期,默认关闭(`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。处理终态与业务变更在同一事务边界提交,**该结论只对处理器产出的「业务型终态」成立**;`MALFORMED / PROTOCOL / SKIPPED(重复) / EXHAUSTED` 这类非业务型终态不涉及航班表,只需一条 `PROC_STATE` UPDATE(终态与回填意图同语句写入),不取 `PIPELINE_LOCK`。完整的事务边界见 §3.3 表 1。 +分派与落库由 `MessageProcessor` 统一协调:按 `MsgKind` 把已绑定身份的队头消息交给对应事务协调器(DNLD/RESP → `ScheduleProcessor`,ADFT → `AdftProcessor`,FLOP → `FlopProcessor`,FDEL → `FdelProcessor`,其余 → `FAILED(UNSUPPORTED)`)。这些 Processor 在 `PIPELINE_LOCK` 事务内读取当前完整态,调用纯领域决策逻辑得到下一完整态与待发事件,再统一落库并登记回填意图;它们不直接触碰 Kafka。处理终态与业务变更在同一事务边界提交,**该结论只对处理器产出的「业务型终态」成立**——非业务型终态不涉及航班表,不取 `PIPELINE_LOCK`。完整边界唯一见 §3.3 表 1。 ### 2.3 状态与错误分类 @@ -96,56 +96,45 @@ SNDR | TYPE | STYP | SEQN `InboxPoller` 默认每秒按 ID 升序、有限批次(`claim-batch`,默认 50)读取水位之后的信箱记录(`ID > W`,**不以处理标记为谓词**),在自有 PG 建立 `PENDING` 并把水位推进到连续上界;入队与水位推进在同一 PG 事务内提交,重复扫描幂等、中断后重扫补建。收报层不解析业务载荷,也不回填已处理标记。 -**收报流程**(每轮 `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);多实例并发收报会让水位互相覆盖,必须先有实例级排他。 +收报流程、空洞判定(遇空洞即停 / 超期放行)及其代价与前提(Q2 承诺、`max-commit-delay` 老化阈值、单实例排他)**唯一定义于** [message-lifecycle.md](message-lifecycle.md) §5.1,本节不重复。水位不是已处理标记。 兼容 HTTP 入口执行“写入共享信箱 → PG 入队”。两步不在同一事务中:信箱成功而 PG 失败时,原文不能丢失,由轮询补建;客户端失败重试可能再次写信箱,业务身份去重仍然必需。 -兼容入口只写 `PROC_STATE`、**不参与水位**,因此它登记的行会**超出水位**;主泵在水位追平前不领取(§3.2 步骤 2),顺序因此不受影响——代价是这类消息要等收报把 `W` 推到它的 ID 之后才开始处理(最长约一个空洞老化窗口)。**运维含义**:只使用兼容入口而不运行收报轮询时,这些行不会被处理,必须让 `InboxPoller` 运行(`msgx.pipeline.autostart=true` 或显式触发)。详见 [message-lifecycle.md](message-lifecycle.md) §5.1/§12。 +兼容入口只写 `PROC_STATE`、不参与水位,登记的行因此超出水位,主泵在水位追平前不领取(§3.2 步骤 2);完整语义与运维含义唯一见 [message-lifecycle.md](message-lifecycle.md) §5.1。 ### 3.2 主泵调度 每次 `Pump.tick` 只围绕最小未完成消息(`PENDING` 与 `FAILED` 都占队头): 1. 无队头:按轮询间隔休眠。 -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)。 +2. 队头超出水位(`msgId > W`):**不领取**,休眠到下一轮。这类行只可能来自兼容入口的直接登记;允许领取会让它越过尚未入队的较小 ID。会打印一条**限流 WARN**(仅在水位值变化时打一次)。 +3. 队头为 `FAILED` 且已毒丸(`attempts ≥ max-attempts`,或 `now − 计时起点 ≥ head-deadline`):转 `DEAD(EXHAUSTED)`,终态与回填意图同一条 UPDATE 落库,**不做跨库写**。**该分支只写 `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` 为稳定起点,**该列为空时(V3 迁移之前的存量行)以 `UPDATED_AT` 兜底**;人工重放会清空 `PROCESSING_STARTED_AT`,由新一轮首次处理重新记录。 +维护作业由独立 job 线程调度、不占用消息循环(§6.1)。**所有取时统一经注入 `Clock`**(收报空洞老化、主泵调度、处理器落库时间、回填重试、作业切日),不使用系统时钟;HOL deadline 以首次处理时写入的 `PROCESSING_STARTED_AT` 为稳定起点,**该列为空时(V3 迁移之前的存量行)以 `UPDATED_AT` 兜底**;人工重放会清空 `PROCESSING_STARTED_AT`,由新一轮首次处理重新记录。 ### 3.3 单条处理 ```text processOne(head): - 1. 入口守卫:head 已是 FAILED 且 attempts 达上限 → DEAD(EXHAUSTED),结束 - 2. 读原文:缺失 → DEAD(MALFORMED, raw-missing);读取异常 → FAILED(INFRA) - 3. 解码: + 1. 读原文:缺失 → DEAD(MALFORMED, raw-missing);读取异常 → FAILED(INFRA) + 2. 解码: 报文非法(MALFORMED) → DEAD(MALFORMED),不重试 可修复解码错(CODEC_ERROR) → FAILED(CODEC_ERROR) 退避 - 4. 身份绑定(仅当 IDENTITY_KEY 为空): + 3. 身份绑定(仅当 IDENTITY_KEY 为空): 已被本消息占用 → 继续 已被别的消息占用 → SKIPPED(duplicate-of:),结束 空闲 → 写入 IDENTITY_KEY(独立单语句,不参与业务事务) - 5. 按 MsgKind 分派: + 4. 按 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) 试写一次信箱标记 + 5. 业务型成功:处理器在自己的事务内写航班变更 + 待发事件 + SUCCEEDED + 回填意图 + 6. 结束:主泵不做回填;回填意图已随终态落库,由扫描补写信箱标记 ``` **表 1 事务边界**(哪些动作在一个事务里、哪些不是): @@ -156,7 +145,7 @@ processOne(head): | 身份首次绑定 | 否 | 否 | 否 | 单语句 + `uk_proc_identity` | | 业务型终态(处理器产出 `SUCCEEDED`) | 是 | 是 | 是 | 同库事务:航班变更 + 事件 + 终态 + 回填意图 | | 非业务型终态(`MALFORMED` / `PROTOCOL` / `SKIPPED` / `EXHAUSTED`) | 否 | 否 | 否 | 单语句(终态与回填意图同一条 UPDATE) | -| 毒丸升级 `DEAD(EXHAUSTED)`(§3.2 步骤 2) | 否 | 否 | 否 | 单语句 | +| 毒丸升级 `DEAD(EXHAUSTED)`(§3.2 步骤 3) | 否 | 否 | 否 | 单语句 | | 回填(信箱标记 + `BACKFILL_AT`) | 否 | 否 | 否 | 跨库两次单写;幂等可重跑 | | 人工重放(批量改回 `PENDING`) | 否 | 否 | 否 | 单语句批量;`MessageLifecycleGate` 与回填互斥 | @@ -164,8 +153,7 @@ processOne(head): - 原文缺失归为 `MALFORMED`;读取异常不能伪装成“缺失”,应进入基础设施重试。 - 忽略规则(`LDM / REGN / RSTA / EROR` → `SKIPPED`)尚未实现(§10);不能因类型未覆盖就把合法忽略报文当非法报文处理。 -- **主泵不回填**:终态与回填意图由同一条 UPDATE 落库,回填**一律由扫描驱动**(跨库写不能占用 FIFO 关键路径)。失败按退避重试,抵达超期期限 R 时强制补写,确认行不存在或达尝试上限则停止自动重试([message-lifecycle.md](message-lifecycle.md) §5.2,US-09/Q7)。回填只需消息 ID,因此缺 META 或解码失败的死信同样可补写。`PENDING / FAILED` 禁止回填;影子环境禁写。 -- 回填有三种结果:写入成功;**此前已被标记(视为成功,不覆盖已有值)**;信箱行不存在(视为失败,当前无终态,见 §10)。 +- **主泵不回填**:终态与回填意图由同一条 UPDATE 落库,回填**一律由扫描驱动**(跨库写不能占用 FIFO 关键路径);四种结果、退避、超期强制补写与放弃恢复、调度周期口径全在 [message-lifecycle.md](message-lifecycle.md) §5.2(US-09/Q7)。回填只需消息 ID,缺 META 或解码失败的死信同样可补写;【缺口】影子环境禁写(§10)。 ## 4. 日计划快照与请求匹配 @@ -176,7 +164,7 @@ processOne(head): 1. **重放判定**:`PROC_STATE` 已存在成功终态 → 幂等成功,仅追加留痕,不重复写入。 2. **整包校验**:声明记录数、航班标识与运营日推导等校验失败 → 整包 `DEAD(PROTOCOL)`,不写半包,既有状态保持不变。 3. **事务写入**:锁内按 `FLID` 点查归属日,发现同一航班跨运营日即整包回滚并 `DEAD(PROTOCOL)`;通过后合并写主表与资源明细。报文未携带的航班不因本次日计划报文被删除。 -4. **提交结果**:同一事务保存 `KAFKA:schd` / `KAFKA:msg` 事件、置消息 `SUCCEEDED` 并预登记回填意图;事务提交后执行信箱回填,留痕在事务外追加。 +4. **提交结果**:同一事务保存 `KAFKA:schd` / `KAFKA:msg` 事件、置消息 `SUCCEEDED` 并预登记回填意图;提交后信箱回填由扫描承接(§3.3),留痕在事务外追加。 单事务保证未提交变更整体回滚。消息重放由 `PROC_STATE` 的消息 ID 与业务身份控制;版本号不能单独证明消息身份。 @@ -221,19 +209,19 @@ REGISTERED → SENT → WAITING → DONE ### 6.1 失败、重试与重放 -`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` 判定。 +`ProcFailure` 与 `FailureScheduler` 统一处理侧失败落账,投递侧(`Dispatcher`)按同一套次数与退避规则迁移事件。默认最多 5 次(attempts ≥ 5 转 `DEAD`,不再计算下次重试);退避表默认 1、2、4、8 秒,**启动自检强制档位数 = max-attempts − 1**,"表里有档但永不触发"的配置不可能出现;单档封顶 `backoff-cap-ms`(默认 60 秒)仅对超过封顶的档位生效。时间经可注入 `Clock` 判定。 失败必须在持有具体消息、事件或批次的位置记录,外层循环只做兜底日志和等待,不重复增加次数。线程中断应恢复中断标记并向上传递;不把 JVM `Error` 当普通业务失败捕获。 -`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。旧消息进入终态后后续消息可能已执行,**重新入队不等于恢复历史顺序**;人工重放前必须评估状态覆盖和版本保护,不能直接批量重放到生产。 +`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 步骤 3),该竞态窗口登记于 §10。旧消息进入终态后后续消息可能已执行,**重新入队不等于恢复历史顺序**;人工重放前必须评估状态覆盖和版本保护,不能直接批量重放到生产。 -**维护作业**:`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。 +**维护作业**:`JobRunner` 用独立 daemon 线程每 30 秒触发 `BackfillService.sweep`(调度周期;扫描谓词、超期、放弃与公平轮转口径唯一见 [message-lifecycle.md](message-lifecycle.md) §5.2),每天机场时区 03:30 后触发一次 `HistorySweepJob`(§6.2)。回填意图在写终态的同一条语句里登记在 `PROC_STATE`(`BACKFILL_NEXT_AT/ATTEMPTS/ERROR`)。作业不参与消息 FIFO,也不使到期消息饥饿。 ### 6.2 历史清理与归档 **航班历史清理**(`HistorySweepJob`,每天 03:30 触发):按 `HistoryProps` 的保留期与终态/静默判据选出候选(含 `DELETED`),先写历史存储,成功后物理删除主行与明细;历史存储未接通或 `msgx.history.history-store-enabled=false` 时删除 0 条。未经 FDEL、由生命周期直接清除的航班,清除前补发一次删除事件。语义与红线见 flight-state.md §6,不在此重复。 -**留痕清理**:`SCHD_SNAP_LOG` 保留 90 天,在历史清理窗口内按 `(SCOPE_END, RECV_AT)` 删除。 +**留痕清理**:`SCHD_SNAP_LOG` 保留 90 天,在历史清理窗口内按 `(SCOPE_END, RECV_AT)` 删除;该清理不依赖历史存储开关,随作业每日执行。 **处理终态归档**:`PROC_STATE_HST` 仍是目标表(user-stories.md US-11),尚未建表;不得归档 `PENDING / FAILED`,也不能因移走身份记录而失去业务去重能力。ES 历史投影(阶段 B)不启用。自有记录归档与信箱原文保留的关系见 [message-lifecycle.md](message-lifecycle.md) §8/§9。 @@ -247,18 +235,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 / backfill-max-attempts`:空洞老化阈值(Q2 的「ID 分配 → 事务可见时延上界」,同时决定目标补偿扫描的窗口宽度)、超期补写期限 R(Q6)、回填扫描批量与回填自动重试上限(达上限停止自动重试并可人工恢复);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.pipeline.max-commit-delay / overdue-backfill / backfill-batch / backfill-max-attempts`:空洞老化阈值(Q2 的「ID 分配 → 事务可见时延上界」;默认值依据缺口唯一见 [message-lifecycle.md](message-lifecycle.md) §12 G7)、超期补写期限 R(Q6;完整约束唯一见 [message-lifecycle.md](message-lifecycle.md) §5.2)、回填扫描批量与回填自动重试上限(达上限停止自动重试并可人工恢复);R 与老化阈值都不能为提速而下调。 +- `msgx.pipeline.delivery-batch / delivery-drain-rounds`:普通事件(`KAFKA:msg`)单目标每轮领取条数上限,与每轮最多连取批数(连取后让出一次循环跑 `schd` flush,防长积压饿死状态通知);保留"队头失败即停止本轮"的目标内保序。 - `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 即关闭。检测只计数与告警,不补入队、不改变处理语义。 +- `msgx.pipeline.cutover-watermark`:**一次性、显式**的切流播种(默认不配置 = 不播种),取值 `min` / `zero` / `max` 或具体 ID;播种的动机、升级实例拒绝重播与 `SEEDED_AT` 语义唯一见 [message-lifecycle.md](message-lifecycle.md) §5.1。非法取值由启动自检挡下。 +- `msgx.pipeline.late-detect-period / late-detect-batch`:**只读**迟到检测的周期与单轮复查量;周期设 0 即关闭。检测只计数与告警,不补入队、不改变处理语义。 -日志关联消息 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。日志出口故障不得阻塞业务线程。 +日志关联消息 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`,无可比记录的年龄/滞后类仪表上报 `-1`,两者都不伪造 0。日志出口故障不得阻塞业务线程。 ## 8. 验证要求 @@ -305,11 +293,10 @@ REGISTERED → SENT → WAITING → DONE 以下缺口直接影响上述设计是否成立,不能以类或接口已存在作为完成依据。收报与处理链路的逐条缺口另见 [message-lifecycle.md](message-lifecycle.md) §12。 -- **事务与外部副作用**:业务型终态(处理器产出的 `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」的窗口。 +- **事务与外部副作用**:事务边界以 §3.3 表 1 为准;[message-lifecycle.md](message-lifecycle.md) §4 的两个崩溃窗口已随 V2 迁移消除。回填本身仍是跨库单写,失败按退避重试并由超期期限 R 兜底;R 的取值待 Q6 确认。 +- **收报与调度**:水位与空洞计时已与入队同事务推进([message-lifecycle.md](message-lifecycle.md) §5.1 口径)。**较小 ID 迟提交尚无补入队机制**:只读迟到检测(阶段 0)已实装,仅计数与告警;窗口补偿扫描(阶段 1,G1)未实现,端到端顺序保证仍待与库方联合验证。空洞老化阈值 `max-commit-delay` 的取值缺依据(§12 G7)。HOL 计时已改用 `PROCESSING_STARTED_AT` + 可注入 `Clock`(§3.2);Q6 仍需确认长期积压与人工重放的期限口径。 +- **重放互斥**:`MessageLifecycleGate` 只覆盖回填与人工重放,毒丸升级路径不在其中,存在「先标 DEAD 并打标、再被重放拨回 PENDING」的窗口(§12 G5)。 +- **原文保留通道**:若库方清除语义为"打标即清除",回填成功后原文即可被清除,重放窗口失去保护;独立原文保留通道尚未设计(§12 G10)。 - **快照与业务能力**: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 条。 diff --git a/docs/message-lifecycle.md b/docs/message-lifecycle.md index b6180b3..5f5b809 100644 --- a/docs/message-lifecycle.md +++ b/docs/message-lifecycle.md @@ -30,7 +30,7 @@ 分环节的原因是存储边界:信箱与自有 PG 之间的写入不共享事务。“落信”以共享 MySQL 为准,“入队”和“处理完成”以自有 PG 为准;“已回填”与“投递确认”各以外部操作成功和本地确认事实共同判定。后两个环节各自独立重试,都可能单独失败。“投递确认”只表示 Kafka Broker 或其他投递目标已接受,不表示业务消费者已消费。 -“入队”这一环还包含一个容易被忽略的前提:**发现是完整的**。它不由本系统单独保证,而依赖信箱 ID 单调与最大提交时延两条上游承诺(Q2);承诺不成立时,本系统只能看到"读过的区间",不能证明"读全了"(§5.1)。 +"入队"这一环还有一个不由本系统单独保证的前提:**发现是完整的**——依赖 Q2 的两条上游承诺(§5.1)。 ## 2. 状态总纲 @@ -54,7 +54,7 @@ MSG_EVENT(自有 PG,投递侧,见 design.md §5) - `PENDING` 与 `FAILED` 是处理中的状态,`FAILED` 继续退避重试并占用 FIFO 队头;`SUCCEEDED / SKIPPED / DEAD` 是终态,到达后队列方可推进。 - `DEAD` 与 `FAILED` 不是不可逆:经人工批准,指定错误类别的记录可重新置回 `PENDING` 处理(放行范围见 §9 / design.md §6.1)。重放处理的是当前状态,不恢复历史处理顺序。 -- **写事务边界**:只有处理器产出的业务型终态才与航班变更、待发事件、回填意图同处一个 PG 事务;非业务型终态(`MALFORMED / PROTOCOL / SKIPPED / EXHAUSTED`)与毒丸升级只写 `PROC_STATE` 一条记录,终态与回填意图由同一条语句写入。完整边界见 [design.md](design.md) §3.3 表 1。 +- **写事务边界**:业务型终态与航班变更、待发事件、回填意图同处一个 PG 事务;其余终态与毒丸升级只写 `PROC_STATE`(终态与回填意图同一条 UPDATE)。完整边界见 [design.md](design.md) §3.3 表 1。 - 信箱处理标记写在该事务**之后**,由与终态同行的回填意图驱动重试与超期补写。 ## 3. 处理阶段 @@ -72,7 +72,7 @@ MSG_EVENT(自有 PG,投递侧,见 design.md §5) | 阶段 | 输入 | 输出 | 幂等依据 | 失败处理 | |---|---|---|---|---| -| 回填扫描(每 30 秒,**调度周期 ≠ 完成时限**) | 终态 ∧ 未确认标记 ∧ 未放弃 ∧(退避到期 ∨ 已达 `NOW − R`) | 追平标记,或判定放弃自动重试 | 消息 ID(仅空标生效) | 退避 30 秒起步、封顶 15 分钟;越过 `R` 后每轮都试;**确认行不存在**或达尝试上限则停止自动重试(可人工恢复) | +| 回填扫描(每 30 秒,**调度周期 ≠ 完成时限**) | 扫描谓词见 §5.2 | 追平标记,或判定放弃自动重试 | 消息 ID(仅空标生效) | 退避、超期强制补写与放弃口径见 §5.2 | 主泵**不做**回填:终态与回填意图由同一条 UPDATE 落库,回填一律由扫描驱动(跨库写不能占用 FIFO 关键路径)。 @@ -99,7 +99,7 @@ MSG_EVENT(自有 PG,投递侧,见 design.md §5) └─ 回填扫描(每 30 秒)按表 2 的谓词重试,越过 R 后强制补写 ``` -领域决策逻辑不执行 I/O。Processor 作为事务协调器:业务型终态在持有 `PIPELINE_LOCK` 的事务内写入自有 PG;非业务型终态与毒丸升级只写 `PROC_STATE` 单条记录。所有外部副作用(读信箱原文、写信箱标记)都发生在 PG 事务边界之外,跨库不共享事务。 +领域决策逻辑不执行 I/O;业务型/非业务型终态的写入边界与全部外部副作用(读信箱原文、写信箱标记)都在 PG 事务之外。跨库不共享事务。 ## 4. 故障恢复 @@ -111,19 +111,13 @@ MSG_EVENT(自有 PG,投递侧,见 design.md §5) | 水位卡在空洞 | `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「事务与外部副作用」一致),不再是缺口: - -1. **提交后回填前崩溃**:回填意图与处理终态是同一条记录的同一次写入、同一个事务;重启后扫描按该意图继续补写。 -2. **回填意图二次落账失败**:不存在第二处落账——意图就在终态行上。 - -仍属未闭环的有两处:回填本身在跨库单写期间的持续失败,由退避重试与 §5.2 的超期期限 `R` 兜底;以及 §5.2 的"信箱行已不存在"没有终态(登记于 §12)。两者均登记在 design.md §10 与 user-stories.md US-09。 +两个历史窗口(提交后回填前崩溃、回填意图二次落账失败)已随"终态与回填意图同体同行"消除——意图不存在第二处落账(口径与 design.md §10「事务与外部副作用」一致)。仍属未闭环的只有回填本身在跨库单写期间的持续失败:由退避重试与 §5.2 的超期期限 `R` 兜底,登记于 design.md §10 与 user-stories.md US-09;"信箱行不存在"已有终态(§5.2 放弃语义)。 ## 5. 消费水位、超期补写与历史积压 @@ -162,7 +156,7 @@ MSG_EVENT(自有 PG,投递侧,见 design.md §5) | 迟到检测(只读,阶段 0) | 复查**被放行的空洞 ID** 是否后来真的出现在信箱 | 进程内监视队列 + `existingIds` 批量存在性检查(`late-detect-period`) | 【现状】已实现:只计数与告警,**不补入队** | | 补偿扫描(阶段 1) | 发现"提交晚于水位推进"的迟到行并**安全补入队** | 按周期重扫 `W` 之前一个窗口(宽度由最大提交时延决定)内的 ID 区间 | 【缺口】未实现 | -【缺口】在补偿扫描交付之前,"较小 ID 迟提交"**没有补入队机制**:快路径只读 `ID > W`,一旦水位越过某个 ID,该 ID 之后到达的消息永远不会被发现。阶段 0 的**只读迟到检测**已实现——监视"被放行的空洞 ID"是否后来真的出现,命中即计数(`msgx.pipeline.late_arrival.detected.total`)并告警,把原先的**静默丢失**变成可观测事实;但它**不补入队**,因此当前设计仍**不能对外声明"迟到不越序"**([design.md](design.md) §8/§10)。检测是进程内尽力而为(重启丢失监视),且只覆盖"曾被放行过的空洞 ID"——这已足够:水位只会越过连续存在的行与被判定为空洞的 ID。 +【缺口】在补偿扫描交付之前,"较小 ID 迟提交"**没有补入队机制**:快路径只读 `ID > W`,一旦水位越过某个 ID,该 ID 之后到达的消息永远不会被发现。阶段 0 的**只读迟到检测**已实现——监视"被放行的空洞 ID"是否后来真的出现,命中即计数(`msgx.pipeline.late_arrival.detected.total`)并告警,把原先的**静默丢失**变成可观测事实;但它**不补入队**,因此当前设计仍**不能对外声明"迟到不越序"**([design.md](design.md) §8/§10)。检测是进程内尽力而为(重启丢失监视;监视队列有上界,被放行 ID 超量后只保留后段),且只覆盖"曾被放行过的空洞 ID"——这已足够:水位只会越过连续存在的行与被判定为空洞的 ID。 **水位与其余两条写路径的关系**: @@ -177,20 +171,22 @@ MSG_EVENT(自有 PG,投递侧,见 design.md §5) ### 5.2 超期标记补写 -所有处理终态都只依赖消息 ID 执行回填;报文残缺或缺少 META 不妨碍正常回填。本节规定回填的三种结果与"长期失败"时的超期兜底。 +所有处理终态都只依赖消息 ID 执行回填;报文残缺或缺少 META 不妨碍正常回填。本节规定回填的四种结果、"长期失败"时的超期兜底,以及 `R` 对重放窗口的实际作用。 **处理规则**:已达 `PROC_STATE` 终态、且接收时间超期(信箱 `DATE_RECEIVED`,落库为 `PROC_STATE.RECEIVED_AT`;判据 `RECEIVED_AT < NOW − 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 +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 ( RECEIVED_AT IS NOT NULL -- 接收时间可用 + AND RECEIVED_AT < NOW − R ) ) -- 已达超期期限,覆盖退避(一旦成立恒成立) +ORDER BY BACKFILL_ATTEMPTS ASC, MSG_ID ASC -- 公平轮转,永久失败行不占满批次 +LIMIT backfill-batch ``` **回填的可能结果**: @@ -210,17 +206,9 @@ ORDER BY MSG_ID ASC LIMIT backfill-batch - 重放窗口的保护只能来自下面两条路之一,且 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` 与审计日志,不在信箱新增枚举。 -- 补写只针对空标记;已有值不回撤、不覆盖,重复执行无副作用。 - `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 中"积压挡批"场景)、`R` 与 `R_keep` 的约束关系可核验(§5.2/§6/§9)。**它仍不保证"§6 的清除边界能在有限时间内达到"**:放弃行永远不满足"边界内全部行已打标",而"打标即清除"语义(§12 G10)会让该前提根本不成立。 ### 5.3 历史积压消息的处理方案 @@ -233,14 +221,12 @@ ORDER BY MSG_ID ASC LIMIT backfill-batch 1. 入队:`InboxPoller` 按升序有限批次将历史行全部建为 `PROC_STATE(PENDING)`,水位随之推到积压末端(连续推进以 §5.1 的 Q2 承诺为前提)。此阶段只写自有 PG,不触碰信箱标记。 2. 处理:主泵按最小未完成 ID 顺序消化。顺序、阈限与日常完全相同:不加速、不分流、不走旁路。积压期间到达的实时消息 ID 更大,自然排在积压之后。 -**入队与处理可以并发**:两者由不同线程驱动,FIFO 由"队头取最小未完成 `MSG_ID`"保证,并发不会破坏顺序。因此"先入队后处理"是**可选的运维规程**(例如先手动跑一段只收报的窗口,便于提前把水位推到完整上界),**不是正确性前提**,系统也不需要提供"只入队"模式。 - -这样安排的收益:入队廉价、处理昂贵;即使并发,入队阶段的中断恢复也只涉及重扫(§4 第一行),水位能尽早到达完整上界,快路径与补偿扫描窗口立即生效。 +**入队与处理可以并发**:两者由不同线程驱动,FIFO 由"队头取最小未完成 `MSG_ID`"保证。因此"先入队后处理"是**可选的运维规程**、**不是正确性前提**,系统也不提供"只入队"模式;入队阶段的中断恢复只涉及重扫(§4 第一行)。 **第三步:等待过程中的四个边界。** - 不插队:主泵按 FIFO 消化,积压期间到达的实时消息排在积压之后。本设计不允许并行队头,也不允许实时通道跳过积压。 -- 不失控:队头滞留上限(`msgx.pipeline.head-deadline`,当前 10 分钟)对积压同样生效;队头长期失败按既有规则转 `DEAD(EXHAUSTED)` 并告警,不会因积压而延长容忍。判据以首次处理时写入的 `PROCESSING_STARTED_AT` 为稳定起点、经可注入 `Clock` 判定(design.md §3.2、§10);长期积压与人工重放的期限口径仍由 Q6 定案。 +- 不失控:队头滞留上限(`msgx.pipeline.head-deadline`,当前 10 分钟)对积压同样生效;队头长期失败按既有规则转 `DEAD(EXHAUSTED)` 并告警,不会因积压而延长容忍(判据口径见 design.md §3.2);长期积压与人工重放的期限口径仍由 Q6 定案。 - 报文类型的处理方式不变:旧的全量日计划(`SCHD-DNLD`)按最新一份**合并**即可收敛(缺席不删、缺失字段保留;字段级冲突见 flight-state.md §3.1 与 Q13),但仍逐条执行;动态增量(`FLOP-*` / `FDEL` / `ADFT`)持有时序语义,必须逐条。 - 可放弃但必须留痕:摸底批中经库方确认"不再处理"的行,处置方式为——`PROC_STATE` 置 `SKIPPED` 并记录跳过原因,到达终态后走 §5.2 通道补写标记;本方案不支持任何"整段 DELETE"的快速通道。 @@ -255,7 +241,7 @@ ORDER BY MSG_ID ASC LIMIT backfill-batch **方案 A:按日分区(首选)**,适用于库方可以为表增加分区的场合。 1. `CMINMSGS` 按 `DATE_RECEIVED` 建立日粒度 RANGE 分区; -2. 某分区到达保留期时,确认该分区全部行已持有处理标记(§5.2 只在回填可持续成功时保证该条件;存在回填饥饿或"打标即清除"语义时见 §12 G9/G10); +2. 某分区到达保留期时,确认该分区全部行已持有处理标记(§5.2 只在回填可持续成功时保证该条件;放弃行与"打标即清除"语义都会破坏该前提,见 §12 G10); 3. 将该分区复制入历史表:`INSERT INTO CMINMSGS_HST SELECT`(`NOT EXISTS` 判重); 4. `TRUNCATE / DROP PARTITION` 执行清除:DDL 级操作,无行锁竞争,页外大字段(`CMINMSGS_CLOB_MSG` 列)整块释放;碎片与 binlog 影响待现场验证。 @@ -274,7 +260,7 @@ ORDER BY MSG_ID ASC LIMIT backfill-batch 方案共同前提: -- **方案要求库方按"标记 + 保留期"清除**(打标本身不触发清除):处理标记只是可清除的**必要条件**,触发条件是"到达保留期 `R_keep`" ∧ "边界内全部行已持有处理标记"。该语义属 Q7/Q9 确认范围;若确认结果是"打标即可清除",则**本节前提不成立,且增大 `R` 无法补救**(原因见 §5.2:回填不等 `R`,打标时刻与 `R` 解耦),必须另行约定保留期或引入独立原文保留通道(§12 G10)。 +- **方案要求库方按"标记 + 保留期"清除**(打标本身不触发清除):处理标记只是可清除的**必要条件**,触发条件是"到达保留期 `R_keep`" ∧ "边界内全部行已持有处理标记"。该语义属 Q7/Q9 确认范围;若确认结果是"打标即可清除",则**本节前提不成立,且增大 `R` 无法补救**(原因唯一见 §5.2),必须另行约定保留期或引入独立原文保留通道(§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 的可见性时延。 @@ -308,9 +294,9 @@ ORDER BY MSG_ID ASC LIMIT backfill-batch ## 9. 重放与原文可用性 - 可重放的错误类别为 `CODEC_ERROR / UNSUPPORTED / INFRA / EXHAUSTED` 四类(以 design.md §6.1 为准);重放处理当前状态,不恢复历史顺序。 -- 重放的前提是信箱原文仍可读取。**该前提只在"标记 + 保留期"清除语义下才能由配置满足**:`R_keep` 必须覆盖(人工重放期限 + 人工处置期限)(§6)。若库方是"打标即清除",则需独立的原文保留通道(§5.2、§12 G10)——**增大 `R` 无效**。 +- 重放的前提是信箱原文仍可读取。该前提仅在"标记 + 保留期"清除语义下可由配置满足(`R_keep` 覆盖重放与处置期限,§6);"打标即清除"语义下的处置唯一见 §5.2 与 §12 G10。 - 确认信箱行或原文已缺失时,按 `MALFORMED → DEAD` 处理;共享库超时、连接失败等暂时性读取异常按 `INFRA` 退避重试,不得伪装成原文缺失。若原文缺失源于库方违反保留契约提前清除,按契约违例走运维追责通道(契约条款见 Q7/Q9);该消息的死信处置本身不变。 -- **回填阶段的"信箱行缺失"是另一回事**(§5.2):它表示"本系统已经处理完、却找不到要标记的那一行"。此时不改变处理终态,也不重新处理业务;需要区分"该 ID 本就属于被判定的永久空洞"与"行被提前清除"。【缺口】当前无终态与告警,见 §12。 +- **回填阶段的"信箱行缺失"是另一回事**(§5.2):它表示"本系统已经处理完、却找不到要标记的那一行"。此时不改变处理终态,也不重新处理业务;需要区分"该 ID 本就属于被判定的永久空洞"与"行被提前清除"。处置按 §5.2:确认 `MISSING_ROW` 即放弃自动重试并告警,人工对账后才能重排回填。 ## 10. 开放问题索引 @@ -351,13 +337,12 @@ Q2、Q6–Q14 的完整登记见 [user-stories.md](user-stories.md) §6 的 Q | 编号 | 缺口 | 影响 | 相关章节 | |---|---|---|---| -| 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 | +| G1 | **窗口补偿扫描未实现**(阶段 0 只读检测已交付) | 迟提交的小 ID 无补入队机制,不能声明"迟到不越序";检测覆盖与口径唯一见 §5.1 | §5.1、§4 | | G4 | **`RECEIVED_AT` 为 NULL 时 R 兜底失效** | 该行只能靠退避重试;V2 只对存量行用 `UPDATED_AT` 兜底,新行没有 | §5.2、§4 | | G5 | **毒丸升级不在 `MessageLifecycleGate` 内** | 与人工重放并发时,存在"先标 DEAD 并写标记、再被重放拨回 `PENDING`"的窗口。**注意**:回填与重放**共用同一个 gate 单例**这一点已由装配断言守住(`PipelineSmokeTest`),但毒丸路径仍在 gate 之外 | §2、§5.3;design.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 | + +已关闭(现行行为并入正文,不再列为缺口):G2 兼容入口与水位领取(§5.1)、G3 回填缺失行放弃语义(§5.2)、G9 回填扫描公平轮转(§5.2 谓词)。 diff --git a/docs/user-stories.md b/docs/user-stories.md index 4d95b17..b9d71d7 100644 --- a/docs/user-stories.md +++ b/docs/user-stories.md @@ -135,7 +135,7 @@ 1. RESP/DNLD 共用流式解析、整包校验和规范化;校验失败不发布半包,旧快照保持可用。 2. RESP 仅匹配未过期、已发送的开放 RQFD;`DTTM < SENT_AT`、已过期、已被替代或无匹配时,不写业务状态,记录跳过原因。 3. 在自有 PG 单事务内,批处理写入已校验的 `FLIGHT_SCHD` 航班状态与资源明细;本次日计划中未出现的航班不因此被删除。 -4. 在同一 PG 事务中提交 `FLIGHT_SCHD` 变更、`MSG_EVENT` 待发通知与 `PROC_STATE(SUCCEEDED)`;匹配 RESP 同事务完成请求并置 `DONE`;事务提交后执行信箱回填。 +4. 在同一 PG 事务中提交 `FLIGHT_SCHD` 变更、`MSG_EVENT` 待发通知与 `PROC_STATE(SUCCEEDED)`;匹配 RESP 同事务完成请求并置 `DONE`;提交后信箱回填由扫描承接。 5. 相同报文重放不二次写入或重复发事件;单事务崩溃整体回滚,重放幂等。 **当前基础与落点**:DNLD 与 RESP 已共同路由到 `processing/ScheduleProcessor.applyScheduleRecords`,整包校验、归属日冲突整包拒绝与单事务写入已实现;`REQ_TRACK` 表与 `ReqTrackRepository` 已建。仍需补 RESP 应答守卫(开放 RQFD 匹配、时间比对)与请求完成关联逻辑。 @@ -190,7 +190,7 @@ 4. 重复补偿效果幂等,保留稳定的完成时间与审计;重放后的新处理结果不能被旧回填任务覆盖。非法报文缺 META 时也有明确回填方式。 5. 影子模式禁写,双跑仅一个系统持有标记写权;暴露 PG 终态、回填状态、积压、最老年龄与持续失败告警。 -**当前基础与落点**:回填意图与处理终态同体同行(`PROC_STATE.BACKFILL_*`),随业务事务提交,`BACKFILL_TODO` 已随 V2 迁移下线;提交后立即尝试一次,失败由 `BackfillService.sweep` 到期重试(每 30 秒、指数退避 30 秒起步封顶 15 分钟),接收时间超过超期期限 `R` 时强制补写(§5.2)。死信同样可补写——回填只需消息 ID,不依赖 META。剩余:Q7 的标记值集与写权限书面确认;影子环境禁写尚未实装。 +**当前基础与落点**:回填意图与处理终态同体同行(`PROC_STATE.BACKFILL_*`),随业务事务提交,`BACKFILL_TODO` 已随 V2 迁移下线;终态落库后回填一律由扫描驱动(`BackfillService.sweep` 每 30 秒、指数退避 30 秒起步封顶 15 分钟,处理关键路径不做跨库写),接收时间超过超期期限 `R` 时强制补写(§5.2)。死信同样可补写——回填只需消息 ID,不依赖 META。剩余:Q7 的标记值集与写权限书面确认;影子环境禁写尚未实装。 **前置**:US-03 终态接口;Q7、共享库更新权限。覆盖四类终态、事务回滚、重复补偿和重放竞争;生命周期、超期补写与清除口径以 [message-lifecycle.md](message-lifecycle.md) §4/§5.2/§6 为准。