diff --git a/README.md b/README.md index 8dc1c91..2ca75cb 100644 --- a/README.md +++ b/README.md @@ -64,7 +64,9 @@ ## 数据库初始化(ACM2-12 口径) **自有 PostgreSQL**(唯一自有库):Flyway 执行 `db/migration/V1__flight_state_baseline.sql`,建立 -航班当前态、明细表、处理终态与 outbox 等表(PG 方言)。全新库直接执行即可,无 legacy 前置。 +航班当前态、明细表、处理终态与 outbox 等表(PG 方言)。这是**单基线**:原 V2–V10 的净结构已 +合并其中,全新库直接执行即可,无 legacy 前置。已按旧链(V1–V10)迁移过的库版本链与校验和都 +对不上,必须重建 schema 或删除数据卷后重跑,禁止手工 `repair` 或改写 `flyway_schema_history`。 **共享 MySQL(cdairport,他人系统库)**:本系统**不建表/schema**,仅信箱 DML——上游外部写 `CMINMSGS`;本系统 JDBC 轮询读 + 处理回填;出站写 `COUTMSGS`(他人读取发送);表结构与 diff --git a/compose.yaml b/compose.yaml index 0a12137..96a84fa 100644 --- a/compose.yaml +++ b/compose.yaml @@ -30,7 +30,7 @@ services: retries: 10 start_period: 15s - # 由 Flyway 自动执行 V1.0.0 与 V1.1.0 迁移 + # 由 Flyway 自动执行单基线迁移(db/migration/V1__flight_state_baseline.sql) postgres: image: postgres:17-alpine container_name: msgx-dev-postgres diff --git a/docs/reference.md b/docs/reference.md index eda010b..ec9d831 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -103,7 +103,7 @@ | 持久化与恢复 | `infra/persistence/`、`infra/retry/`(`ProcFailure` / `ReplayService` / `FailureScheduler`) | | 启停与配置 | `PipelineLifecycle.kt`、`config/PipelineProps.kt`、`config/HistoryProps.kt`、`config/OperationDayProps.kt`(含运营日时区启动自检) | | 指标与健康 | `infra/metrics/PipelineMetrics.kt`、`infra/metrics/JobActivity.kt`、`infra/health/BacklogSnapshotProvider.kt`、`infra/health/JobRunnerHealthIndicator.kt` | -| 迁移 | `src/main/resources/db/migration/`(V1 基线、V2 生命周期、V3 稳定处理起点、V4 回填闭环、V5 切流播种、V6 本地入队时间、V7 删除处理起点、V8 `REQ_TRACK` 开放态唯一索引、V9 `schd` 单行化;`oracle11g/` 为占位) | +| 迁移 | `src/main/resources/db/migration/`(单基线 `V1__flight_state_baseline.sql`:原 V1–V10 的净结构已合并,迁移链收敛为一条;`oracle11g/` 为占位) | ## 4. 错误分类与重放白名单 diff --git a/docs/user-stories.md b/docs/user-stories.md index 137da14..f09ec39 100644 --- a/docs/user-stories.md +++ b/docs/user-stories.md @@ -190,7 +190,7 @@ 4. 重复补偿效果幂等,保留稳定的完成时间与审计;重放后的新处理结果不能被旧回填任务覆盖。非法报文缺 META 时也有明确回填方式。 5. 影子模式禁写,双跑仅一个系统持有标记写权;暴露 PG 终态、回填状态、积压、最老年龄与持续失败告警。 -**当前基础与落点**:回填意图与处理终态同体同行(`PROC_STATE.BACKFILL_*`),随业务事务提交,`BACKFILL_TODO` 已随 V2 迁移下线;终态落库后回填一律由定时扫描驱动(`BackfillService.sweep`,扫描周期与退避取值见 [reference.md](reference.md);独立退避键 `backfill-backoff-ms` / `backfill-backoff-cap-ms` 已落地),处理关键路径不做跨库写;本地入队时间(`ENQUEUED_AT`)超过超期期限 `R` 时强制补写(见 design「回填」)。死信同样可补写——回填只需消息 ID,不依赖 META。剩余:Q7 的标记值集与写权限书面确认;影子环境禁写尚未实装。 +**当前基础与落点**:回填意图与处理终态同体同行(`PROC_STATE.BACKFILL_*`),随业务事务提交,`BACKFILL_TODO` 已在单基线中下线(原 V2 迁移的净结果已并入 `V1__flight_state_baseline.sql`);终态落库后回填一律由定时扫描驱动(`BackfillService.sweep`,扫描周期与退避取值见 [reference.md](reference.md);独立退避键 `backfill-backoff-ms` / `backfill-backoff-cap-ms` 已落地),处理关键路径不做跨库写;本地入队时间(`ENQUEUED_AT`)超过超期期限 `R` 时强制补写(见 design「回填」)。死信同样可补写——回填只需消息 ID,不依赖 META。剩余:Q7 的标记值集与写权限书面确认;影子环境禁写尚未实装。 **前置**:US-03 终态接口;Q7、共享库更新权限。覆盖四类终态、事务回滚、重复补偿和重放竞争;生命周期与超期补写以 [design.md](design.md)「中断恢复」「回填」为准,清除口径以 [contracts.md](contracts.md)「保留与清除」为准。 diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index 41938c4..19c4e46 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -47,8 +47,8 @@ micronaut: # 管理端点(U03/N32):/env、/beans 默认 sensitive;仅开发/影子环境放开——见 application-dev.yml # 自有 PostgreSQL:全部内部状态(PIPELINE_LOCK/PROC_STATE/MSG_EVENT/REQ_TRACK/INBOX_CURSOR/ -# FLIGHT_SCHD 及明细表/SCHD_SNAP_LOG;迁移 db/migration/V1__flight_state_baseline.sql + -# V2__inbox_lifecycle.sql——V2 起回填意图并入 PROC_STATE,BACKFILL_TODO 已下线)。 +# FLIGHT_SCHD 及明细表/SCHD_SNAP_LOG;迁移为单基线 db/migration/V1__flight_state_baseline.sql, +# 原 V2–V10 的净结构已合并其中,回填事实与处理终态同表同行)。 # 数据层实装前 enabled=false(stub 模式不建连)。 datasources: default: diff --git a/src/main/resources/db/migration/V10__event_sent_at.sql b/src/main/resources/db/migration/V10__event_sent_at.sql deleted file mode 100644 index cf7cfbc..0000000 --- a/src/main/resources/db/migration/V10__event_sent_at.sql +++ /dev/null @@ -1,13 +0,0 @@ --- ===================================================================== --- V10:MSG_EVENT 增加 SENT_AT,为有界清理提供时间基准(G-EVENT-RETENTION) --- --- 所有投递确认操作在设置 STATE='SENT' 的同一条 UPDATE 中写入 SENT_AT; --- 清理作业只删除 STATE='SENT' AND SENT_AT IS NOT NULL AND SENT_AT < now - retention。 --- --- 存量 SENT 行的 SENT_AT 设为迁移时刻 CURRENT_TIMESTAMP, --- 避免迁移完成后它们立即过期被清理掉。 --- ===================================================================== - -ALTER TABLE msg_event ADD COLUMN sent_at TIMESTAMP(6) WITH TIME ZONE; - -UPDATE msg_event SET sent_at = CURRENT_TIMESTAMP WHERE state = 'SENT'; diff --git a/src/main/resources/db/migration/V1__flight_state_baseline.sql b/src/main/resources/db/migration/V1__flight_state_baseline.sql index 761bda9..a06c32a 100644 --- a/src/main/resources/db/migration/V1__flight_state_baseline.sql +++ b/src/main/resources/db/migration/V1__flight_state_baseline.sql @@ -1,23 +1,34 @@ -- ===================================================================== --- 航班运行数据接入与当前状态管理 · 全量基线 +-- 自有 PostgreSQL 全量基线(单迁移链条目) -- --------------------------------------------------------------------- --- 权威设计:docs/flight-state.md(审计定稿)。本脚本整体取代历史迁移 --- (V1.0.0–V1.4.0 已删除),按 §3.2 表职责建立全部分层: --- · 决策层:FLIGHT_SCHD + 8 张资源明细表 + FLIGHT_ROUTE_POINT —— 承担正确性; --- · 管道层:PIPELINE_LOCK / PROC_STATE / MSG_EVENT(outbox) / REQ_TRACK / BACKFILL_TODO; --- · 留痕层:SCHD_SNAP_LOG —— 只追加、不参与决策、可重建; --- · 证据层:报文原文归档 —— 冷路径,尚未交付(ARCHIVE_KEY 仅留引用位)。 --- 核心定案(对照文档章节): --- · 身份:FLID 主键;OPERATION_DAY 一经确定不可变(§3.1/§5.3,应用层校验); --- STATE ∈ {ACTIVE, DELETED},无 ARCHIVED —— 物理清除只发生在历史归档成功之后(§8.2); --- · 无名单差删:主链路无集合删除;删除入口只有 FDEL(§6.2)与历史归档(§8.2); --- · 版本:每次成功写入推进 STATE_VERSION;Kafka 判旧 = (FLID, STATE_VERSION, UPDATED_AT)(§7.3); --- · 出站:COUTMSGS 是共享 MySQL 信箱(他人系统库,本系统不建表,迁移不覆盖)。 --- 类型口径:TIMESTAMP(6) WITH TIME ZONE 统一 UTC 语义;无 JSONB/CLOB 依赖; --- BIGSERIAL 为 PostgreSQL 方言(Oracle 适配见 §10,未定前不预设)。 +-- 范围:本迁移只作用于自有 PostgreSQL。共享 MySQL 信箱(CMINMSGS/COUTMSGS)是他人 +-- 系统库,只做契约内 DML,不建表、不改结构(C-14)。 +-- +-- 本文件原为 V1 基线;原 V2–V10 的净结构已全部合并进来,版本链收敛为一条: +-- · INBOX_CURSOR 及 SEEDED_AT(原 V2/V5); +-- · PROC_STATE 的收信/入队时间与回填事实列,独立待办表 BACKFILL_TODO 不再存在 +-- (原 V2/V3/V4/V6/V7:回填意图并入 PROC_STATE,处理开始时间列加后即删); +-- · REQ_TRACK 开放态部分唯一索引 uq_req_open(原 V8); +-- · MSG_EVENT 的 schd 单行唯一索引 uq_schd_event、SENT_AT,以及已无消费者的 +-- idx_evt_flid 下线(原 V9/V10)。 +-- 原 V2–V10 中只对既有库有意义的数据搬迁语句(存量 RECEIVED_AT / ENQUEUED_AT 回填、 +-- 重复开放请求与重复 schd 事件去重、存量 SENT_AT 回写)一律不再保留——基线只服务 +-- 全新库,而 Flyway 的迁移链一旦发布即不可改写,故这些语句的语义只能在此说明。 +-- +-- 升级路径:已应用过 V1–V10 的库无法直接升级(版本链与校验和都对不上)。重建 schema +-- 或删除数据卷后重跑迁移;禁止对生产库手工 repair 或改写 flyway_schema_history。 +-- +-- 分层(职责与写者见 docs/flight-state.md、docs/design.md): +-- · 决策层:FLIGHT_SCHD + 8 张资源明细表 + FLIGHT_ROUTE_POINT —— 权威当前态(INV-11); +-- · 管道层:PIPELINE_LOCK / INBOX_CURSOR / PROC_STATE / MSG_EVENT / REQ_TRACK; +-- · 留痕层:SCHD_SNAP_LOG —— 只追加、可重建、不参与决策; +-- · 证据层:报文原文归档 —— 冷路径,尚未交付(ARCHIVE_KEY 仅留引用位,G-REPLAY-CHANNEL)。 +-- 类型口径:TIMESTAMP(6) WITH TIME ZONE 统一 UTC 语义;BIGSERIAL 为 PostgreSQL 方言 +-- (Oracle 11g 等价 DDL 见同目录 oracle11g/README.md)。 -- ===================================================================== --- ① 单行锁:串行化状态写事务(§1 边界:只串行化 DB 事务,不代替消息认领/选主/故障切换) +-- ① 单行锁:事务内串行化状态写事务(只串行化本地 DB 事务,不代替消息认领、选主或 +-- 故障切换,PRE-5)。历史清理与主泵处理器共用它互斥(INV-18)。 CREATE TABLE PIPELINE_LOCK ( LOCK_ID INT NOT NULL PRIMARY KEY, -- 恒为 1 UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL @@ -25,50 +36,80 @@ CREATE TABLE PIPELINE_LOCK ( INSERT INTO PIPELINE_LOCK (LOCK_ID, UPDATED_AT) VALUES (1, now()); -- ② 处理伴生状态:与共享信箱 CMINMSGS 一一对应(MSG_ID = CMINMSGS_ID)。 --- 兼作快照重放判定(§5.1 步骤 2):MSG_ID 已有成功终态 → 重放,直接记幂等成功。 +-- 一信一行、一身份一记录(INV-9);MSG_ID 已有成功终态即判重放,直接记幂等成功。 +-- 处理终态与回填意图在本表同表同行:不存在第二处落账,也不存在独立的回填待办表(INV-8)。 +-- 状态不可逆:已提交的 SUCCEEDED 不因回填或投递失败回改(INV-6)。 CREATE TABLE PROC_STATE ( MSG_ID BIGINT NOT NULL PRIMARY KEY, - STATE VARCHAR(16) NOT NULL, -- PENDING/FAILED/SUCCEEDED/SKIPPED/DEAD - IDENTITY_KEY VARCHAR(200), -- SNDR|TYPE|STYP|SEQN(decode 后首次绑定;SEQN 同会话顺序观测 §5.5) + STATE VARCHAR(16) NOT NULL, -- PENDING/FAILED/SUCCEEDED/SKIPPED/DEAD(ARCHIVED 属未建归档能力 G-PROC-HST) + IDENTITY_KEY VARCHAR(200), -- SNDR|TYPE|STYP|SEQN,解码后首次绑定 ATTEMPTS INT NOT NULL DEFAULT 0, - NEXT_ATTEMPT_AT TIMESTAMP(6) WITH TIME ZONE, - ERROR_CLASS VARCHAR(20), -- MALFORMED/PROTOCOL/INFRA/UNSUPPORTED/EXHAUSTED(§9) + NEXT_ATTEMPT_AT TIMESTAMP(6) WITH TIME ZONE, -- 退避到期时刻;未到期不得被后续消息越过(INV-3) + ERROR_CLASS VARCHAR(20), -- MALFORMED/PROTOCOL/INFRA/UNSUPPORTED/EXHAUSTED LAST_ERROR VARCHAR(1000), UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, + -- ==== 时间基准 ==== + RECEIVED_AT TIMESTAMP(6) WITH TIME ZONE, -- 复制自信箱 DATE_RECEIVED:库方时钟、可为 NULL,仅用于对账与展示(PRE-4) + ENQUEUED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL DEFAULT now(), -- 本地入队时间:非空、与判据 NOW 同源,是超期期限 R 的唯一比较对象(PRE-4、CLM-4) + -- ==== 回填事实(INV-7/INV-8)==== + BACKFILL_AT TIMESTAMP(6) WITH TIME ZONE, -- 非空 = 处理标记已确认写入 + BACKFILL_NEXT_AT TIMESTAMP(6) WITH TIME ZONE, -- 非空 = 还欠一次回填;与终态同语句写下 + BACKFILL_ATTEMPTS INT NOT NULL DEFAULT 0, + BACKFILL_ERROR VARCHAR(512), + BACKFILL_ABANDONED_AT TIMESTAMP(6) WITH TIME ZONE, -- 非空 = 已停止自动重试;**不等于**标记已确认(C-8、C-16) + BACKFILL_ABANDONED_REASON VARCHAR(64), -- MISSING_ROW(确定性)/ TRANSIENT_DEADLINE(暂时性耗尽 R) CONSTRAINT uk_proc_identity UNIQUE (IDENTITY_KEY) ); -CREATE INDEX idx_proc_head ON PROC_STATE (STATE, MSG_ID); -- 主泵队头查询(严格 FIFO) +CREATE INDEX idx_proc_head ON PROC_STATE (STATE, MSG_ID); -- 主泵队头查询(严格 FIFO,INV-3) +-- 回填扫描:只覆盖仍需自动回填的行;按 (ATTEMPTS, MSG_ID) 轮转,避免最旧失败行长期占满批次。 +CREATE INDEX idx_proc_backfill_due ON PROC_STATE (BACKFILL_ATTEMPTS, MSG_ID) + WHERE BACKFILL_AT IS NULL AND BACKFILL_ABANDONED_AT IS NULL; --- ③ 航班当前态主行:一行一 FLID(决策层权威;展示视图是投影,不是权威 §3.3) +-- ③ 消费水位:全表只有一行。W 只随新 ID 成功入队推进、只增不减,遇空洞即停; +-- 入队与水位推进同事务(INV-2),因此不存在「水位已推进、消息未入队」的持久化状态。 +-- HOLE_SINCE 持久化空洞观测时刻,进程重启不丢计时;老化阈值取 +-- PARAM:msgx.pipeline.max-commit-delay,超期只放行空洞本身、不越过任何已存在的行。 +CREATE TABLE INBOX_CURSOR ( + CURSOR_ID INT NOT NULL PRIMARY KEY, -- 固定为 1 + COMMITTED_UP_TO BIGINT NOT NULL, -- 水位 W:到哪个 ID 为止已全部读进自有库 + HOLE_SINCE TIMESTAMP(6) WITH TIME ZONE, -- 后续缺号最早被发现的时间;不缺号时为 NULL + UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, + SEEDED_AT TIMESTAMP(6) WITH TIME ZONE -- 记录「已按 PARAM:msgx.pipeline.cutover-watermark 播种」这一事实; + -- 为 NULL **不等于**从未消费——已有库新增列后同样为 NULL +); +-- 初值 W=0:首轮把信箱现存行全部重新读一遍,重复登记不会建出第二行。 +INSERT INTO INBOX_CURSOR (CURSOR_ID, COMMITTED_UP_TO, HOLE_SINCE, UPDATED_AT) VALUES (1, 0, NULL, now()); + +-- ④ 航班当前态主行:一行一 FLID(决策层权威;展示视图是投影,不是权威,INV-11) CREATE TABLE FLIGHT_SCHD ( - FLID VARCHAR(32) NOT NULL PRIMARY KEY, -- AODB 实例 ID,Number(1-12)(SIS §3.16.2) - OPERATION_DAY DATE NULL, -- 运营保障日;一经确定不可变(§3.5);NULL=尚未被快照收录 - STATE VARCHAR(8) NOT NULL, -- ACTIVE/DELETED(§3.1,无第三态) - STATE_VERSION BIGINT NOT NULL DEFAULT 0, -- 每次成功写入 +1(§4) + FLID VARCHAR(32) NOT NULL PRIMARY KEY, -- AODB 实例 ID(Number(1-12));不得由航班号或资源号推断 + OPERATION_DAY DATE NULL, -- 运营保障日:未由日计划收录时为 NULL;非空后不可改变(INV-12) + STATE VARCHAR(8) NOT NULL, -- ACTIVE/DELETED,无第三态;删除只由 FDEL 或受控历史清理触发(INV-15) + STATE_VERSION BIGINT NOT NULL DEFAULT 0, -- 每次成功写入 +1,重复消息不重复推进(INV-13) LAST_MSG_ID BIGINT NULL, -- 最近一次成功写入的消息 ID - -- ==== SCHD.FLTR 标量字段(值保持 AODB 字符串原样;缺失语义由处理层决定,§5.4)==== + -- ==== SCHD.FLTR 标量字段(保持 AODB 字符串原样;未携带不隐式清空,INV-14)==== ALCD VARCHAR(64) NULL, -- 航空公司代码 ALSC VARCHAR(64) NULL, -- 航空公司简称 - FLNO VARCHAR(64) NULL, -- 航班号(展示用,不参与身份判定 §3.1) + FLNO VARCHAR(64) NULL, -- 航班号(展示用,不参与身份判定) MVIN VARCHAR(64) NULL, -- 进离港标识(A/D) - SODT VARCHAR(64) NULL, -- 计划运行时间(ddMMMyyHHmm;OPERATION_DAY 计算源 §3.5) + SODT VARCHAR(64) NULL, -- 计划运行时间(ddMMMyyHHmm;OPERATION_DAY 的推导源) FLTY VARCHAR(64) NULL, -- 航班类型 FLIN VARCHAR(64) NULL, -- 国内/国际/混合 ACFT VARCHAR(64) NULL, -- 机型 RENO VARCHAR(64) NULL, -- 机尾号 - TAOP VARCHAR(64) NULL, -- 过站实际承运人代码(§3.4:过站关联,非代码共享) + TAOP VARCHAR(64) NULL, -- 过站实际承运人代码(过站关联,非代码共享) TAFL VARCHAR(64) NULL, -- 过站实际承运航班号 TAID VARCHAR(64) NULL, -- 过站关联航班 FLID TRML VARCHAR(64) NULL, -- 航站楼 MAXP VARCHAR(64) NULL, -- 最大旅客数 - CSOP VARCHAR(64) NULL, -- 共享承运人代码(§3.4) + CSOP VARCHAR(64) NULL, -- 共享承运人代码 CSFT VARCHAR(64) NULL, -- 共享航班号 - MAID VARCHAR(32) NULL, -- 共享主航班 FLID + MAID VARCHAR(32) NULL, -- 共享主航班 FLID:主/共享关系的事实来源;MAFL 由它派生(INV-21、INV-22) ESTT VARCHAR(64) NULL, -- 预计时间 ACTT VARCHAR(64) NULL, -- 实际时间 STND VARCHAR(64) NULL, -- 备降站/机位 PHAG VARCHAR(64) NULL, -- 地服代理(值机) - CNCL VARCHAR(64) NULL, -- 取消时间(§8.1 历史判定依据) + CNCL VARCHAR(64) NULL, -- 取消时间;受控历史清理的判定依据之一(INV-15、INV-18) REMC VARCHAR(256) NULL, -- 备注 BOTM VARCHAR(64) NULL, -- 摆渡车时间 LACL VARCHAR(64) NULL, -- 行李确认时间 @@ -87,22 +128,23 @@ CREATE TABLE FLIGHT_SCHD ( EXSR VARCHAR(256) NULL, -- 例外原因 FTSS VARCHAR(64) NULL, -- 航班状态 PEDT VARCHAR(64) NULL, -- 前序估计时间 - NAAT VARCHAR(64) NULL, -- 到港终态时间(§8.1;业务含义待术语表确认 §10) - NEAT VARCHAR(64) NULL, -- 离港终态时间(§8.1;业务含义待术语表确认 §10) + NAAT VARCHAR(64) NULL, -- 到港终态时间(业务含义待术语表确认) + NEAT VARCHAR(64) NULL, -- 离港终态时间(业务含义待术语表确认) PADT VARCHAR(64) NULL, -- 旅客登机时间 CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL ); CREATE INDEX idx_flight_schd_opday ON FLIGHT_SCHD (OPERATION_DAY); CREATE INDEX idx_flight_schd_state ON FLIGHT_SCHD (STATE); --- OPERATION_DAY 不可变的数据库层强化(§7.4):快照 upsert 一律带 +-- OPERATION_DAY 不可变(INV-12)的库层强化方式:快照 upsert 一律带 -- WHERE OPERATION_DAY IS NULL OR OPERATION_DAY = :day,行级条件更新即满足; -- 不使用触发器(Oracle/PG 双方言成本)。 --- ④ 资源明细表 ×8 + 路线点表(每 FLID 多行;明细集合按组先删后插 §3.3; --- ORDINAL 保留输入顺序,SOURCE_SEQ 只存上游序号 §3.2; --- 删除一律标记在主行 STATE,明细物理清除只发生在历史归档 §8.2) --- ④-1 登机门 GTDT +-- ⑤ 资源明细表 ×8 + 路线点表:每 FLID 多行;每次完整状态写入按该 FLID 先删后插, +-- 以完整合并结果为准(INV-14)。ORDINAL 是持久化顺序(从 1 起),SOURCE_SEQ 是上游 +-- 序号(允许为空或重复);相同资源号不代表同一条分配,禁止按资源号去重。删除一律 +-- 标记在主行 STATE,明细物理清除只发生在受控历史归档(INV-15)。 +-- ⑤-1 登机门 GTDT CREATE TABLE FLIGHT_GATE ( FLID VARCHAR(32) NOT NULL, ORDINAL INT NOT NULL, @@ -117,7 +159,7 @@ CREATE TABLE FLIGHT_GATE ( UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, PRIMARY KEY (FLID, ORDINAL) ); --- ④-2 值机柜台 CKDT +-- ⑤-2 值机柜台 CKDT CREATE TABLE FLIGHT_CHECKIN ( FLID VARCHAR(32) NOT NULL, ORDINAL INT NOT NULL, @@ -133,7 +175,7 @@ CREATE TABLE FLIGHT_CHECKIN ( UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, PRIMARY KEY (FLID, ORDINAL) ); --- ④-3 行李转盘 CLDT +-- ⑤-3 行李转盘 CLDT CREATE TABLE FLIGHT_BELT ( FLID VARCHAR(32) NOT NULL, ORDINAL INT NOT NULL, @@ -149,7 +191,7 @@ CREATE TABLE FLIGHT_BELT ( UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, PRIMARY KEY (FLID, ORDINAL) ); --- ④-4 计划机位 PSDT +-- ⑤-4 计划机位 PSDT CREATE TABLE FLIGHT_STAND_PLAN ( FLID VARCHAR(32) NOT NULL, ORDINAL INT NOT NULL, @@ -161,23 +203,23 @@ CREATE TABLE FLIGHT_STAND_PLAN ( UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, PRIMARY KEY (FLID, ORDINAL) ); --- ④-5 滑槽 CHDT +-- ⑤-5 滑槽 CHDT CREATE TABLE FLIGHT_CHUTE ( FLID VARCHAR(32) NOT NULL, ORDINAL INT NOT NULL, SOURCE_SEQ VARCHAR(64), CHUT VARCHAR(64), CHCLS VARCHAR(64), - PCBT VARCHAR(64), -- 计划开始(SIS 线格式 PCBT) - PCET VARCHAR(64), -- 计划结束(PCET) - CBTM VARCHAR(64), -- 实际开始(CBTM) - CETM VARCHAR(64), -- 实际结束(CETM) + PCBT VARCHAR(64), -- 计划开始 + PCET VARCHAR(64), -- 计划结束 + CBTM VARCHAR(64), -- 实际开始 + CETM VARCHAR(64), -- 实际结束 CHTYP VARCHAR(64), CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, PRIMARY KEY (FLID, ORDINAL) ); --- ④-6 延误 DELY(业务上任意时刻仅 1 个有效延误,保留多行能力以无损承接) +-- ⑤-6 延误 DELY(业务上任意时刻仅 1 个有效延误,保留多行能力以无损承接) CREATE TABLE FLIGHT_DELAY ( FLID VARCHAR(32) NOT NULL, ORDINAL INT NOT NULL, @@ -190,7 +232,7 @@ CREATE TABLE FLIGHT_DELAY ( UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, PRIMARY KEY (FLID, ORDINAL) ); --- ④-7 靠撤桥 ABTM(桥号 ABDG;A/D 各一条) +-- ⑤-7 靠撤桥 ABTM(桥号 ABDG;A/D 各一条) CREATE TABLE FLIGHT_BRIDGE_OP ( FLID VARCHAR(32) NOT NULL, ORDINAL INT NOT NULL, @@ -202,7 +244,7 @@ CREATE TABLE FLIGHT_BRIDGE_OP ( UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, PRIMARY KEY (FLID, ORDINAL) ); --- ④-8 轮挡 CHOT(机位 CHID;ON/OFF 各一条) +-- ⑤-8 轮挡 CHOT(机位 CHID;ON/OFF 各一条) CREATE TABLE FLIGHT_CHOCK_OP ( FLID VARCHAR(32) NOT NULL, ORDINAL INT NOT NULL, @@ -214,7 +256,7 @@ CREATE TABLE FLIGHT_CHOCK_OP ( UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, PRIMARY KEY (FLID, ORDINAL) ); --- ④-9 路线点 ROUT/ERUT 共用(ROUTE_KIND 区分 §3.2:10 类集合 9 张表) +-- ⑤-9 路线点 ROUT/ERUT 共用(ROUTE_KIND 区分两类;主键含该列,避免序号冲突) CREATE TABLE FLIGHT_ROUTE_POINT ( FLID VARCHAR(32) NOT NULL, ORDINAL INT NOT NULL, @@ -236,27 +278,33 @@ CREATE INDEX idx_flight_delay_flid ON FLIGHT_DELAY (FLID); CREATE INDEX idx_flight_bridge_flid ON FLIGHT_BRIDGE_OP (FLID); CREATE INDEX idx_flight_chock_flid ON FLIGHT_CHOCK_OP (FLID); --- ⑤ 统一投递事件 outbox(§7.3:KAFKA_MSG 只通知变化;KAFKA_SCHD 发整态; --- tombstone 仅在 ACTIVE→DELETED 时与删除同事务登记,投递失败持续重试) +-- ⑥ 统一投递事件 outbox(INV-10、INV-17):KAFKA:schd 发整态、KAFKA:msg 只通知变化; +-- tombstone 仅在 ACTIVE→DELETED 时与删除同事务登记,投递失败按退避持续重试。 +-- 目标级全序投递是当前实现,按 FLID 保序见 CLM-7。 CREATE TABLE MSG_EVENT ( - EVENT_ID BIGSERIAL PRIMARY KEY, - TARGET VARCHAR(30) NOT NULL, -- KAFKA_MSG / KAFKA_SCHD - PARTITION_KEY VARCHAR(32) NOT NULL, -- 恒为 FLID + EVENT_ID BIGSERIAL PRIMARY KEY, -- 对 KAFKA:msg 是稳定事件身份并决定投递顺序;对 KAFKA:schd 是每次接受 upsert 时替换的写代次 + TARGET VARCHAR(30) NOT NULL, -- KAFKA:msg / KAFKA:schd + PARTITION_KEY VARCHAR(32) NOT NULL, -- 当前恒为 FLID(Q4 定案前为假定,C-29) EVENT_TYPE VARCHAR(16) NOT NULL DEFAULT 'UPSERT', -- UPSERT / TOMBSTONE - STATE_VERSION BIGINT NOT NULL, -- 发布时航班版本;聚合按 FLID 取最新(§7.3) - PAYLOAD_JSON TEXT NOT NULL, -- TOMBSTONE 时至少含 FLID/STATE_VERSION/DELETED - STATE VARCHAR(16) NOT NULL, -- PENDING/SENT/DEAD + STATE_VERSION BIGINT NOT NULL, -- 发布时航班版本;schd 聚合按 FLID 只进不退 + PAYLOAD_JSON TEXT NOT NULL, -- TOMBSTONE 时至少含 FLID/STATE_VERSION/DELETED + STATE VARCHAR(16) NOT NULL, -- PENDING/SENT/DEAD ATTEMPTS INT NOT NULL DEFAULT 0, NEXT_ATTEMPT_AT TIMESTAMP(6) WITH TIME ZONE, ERROR_CLASS VARCHAR(20), LAST_ERROR VARCHAR(1000), - CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL + CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, + SENT_AT TIMESTAMP(6) WITH TIME ZONE -- 投递确认的同一条 UPDATE 内写入;是保留期判定的唯一基准(G-EVENT-RETENTION) ); CREATE INDEX idx_evt_head ON MSG_EVENT (TARGET, STATE, EVENT_ID); -- 每 target 队头 -CREATE INDEX idx_evt_flid ON MSG_EVENT (PARTITION_KEY, STATE_VERSION); -- 同 FLID 未发事件合并 +-- KAFKA:schd 按 FLID 单行(只保留最新 STATE_VERSION);KAFKA:msg 仍是多行 append-log,不受此约束。 +CREATE UNIQUE INDEX uq_schd_event + ON MSG_EVENT (TARGET, PARTITION_KEY) + WHERE TARGET = 'KAFKA:schd'; --- ⑥ 请求状态机(§7.1:只有 RESP 完成 RQFD 请求,按(运营日、发送方、请求类型) --- 匹配最新一条 PENDING;同类请求只留一条有效,新请求置旧为 EXPIRED) +-- ⑦ 请求状态机:只有 RESP 完成 RQFD 请求,按(运营日、发送方、请求类型)匹配最新一条 +-- 开放请求。同类请求只留一条有效,新请求置旧请求为 EXPIRED。登记、超时与应答匹配 +-- 尚未实现(G-REQ-TRACK、C-24)。 CREATE TABLE REQ_TRACK ( REQ_ID BIGSERIAL PRIMARY KEY, REQ_TYPE VARCHAR(20) NOT NULL, @@ -268,26 +316,14 @@ CREATE TABLE REQ_TRACK ( COMPLETED_AT TIMESTAMP(6) WITH TIME ZONE, CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL ); -CREATE INDEX idx_req_open ON REQ_TRACK (REQ_TYPE, OPERATION_DAY, SENDER, STATE); +CREATE INDEX idx_req_open ON REQ_TRACK (REQ_TYPE, OPERATION_DAY, SENDER, STATE); -- 仍服务含关闭态的查询 +-- 开放态(PENDING/SENT)内 (REQ_TYPE, OPERATION_DAY, SENDER) 唯一;关闭态行可并存。 +CREATE UNIQUE INDEX uq_req_open + ON REQ_TRACK (REQ_TYPE, OPERATION_DAY, SENDER) + WHERE STATE IN ('PENDING', 'SENT'); --- ⑦ 共享信箱回填补偿待办(§7.2:提交后回填;回填失败不得把 SUCCEEDED 改回 FAILED; --- 目标形态 = 业务事务内预登记,消除「提交后写待办」的崩溃窗口,§10 当前偏差) -CREATE TABLE BACKFILL_TODO ( - MSG_ID BIGINT NOT NULL PRIMARY KEY, - SNDR VARCHAR(64) NOT NULL, - TYPE VARCHAR(32) NOT NULL, - STYP VARCHAR(32), - SEQN BIGINT, - ATTEMPTS INT NOT NULL DEFAULT 0, - NEXT_ATTEMPT_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, - LAST_ERROR VARCHAR(512), - CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, - UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL -); -CREATE INDEX idx_backfill_todo_due ON BACKFILL_TODO (NEXT_ATTEMPT_AT); - --- ⑧ SCHD 快照留痕(§5.5):事务外追加,只追加留痕、不参与决策;一行=一次尝试,重放也记; --- RESULT 与 FLAGS 分列(可「成功且告警」);写失败只记指标;保留 90 天,按 (SCOPE_END, RECV_AT) 清理 +-- ⑧ SCHD 快照留痕:事务外追加,只追加留痕、不参与决策;一行=一次尝试,重放也记。 +-- RESULT 与 FLAGS 分列(可「成功且告警」);写失败只记指标。 CREATE TABLE SCHD_SNAP_LOG ( LOG_ID BIGSERIAL PRIMARY KEY, MSG_ID BIGINT NOT NULL, @@ -302,5 +338,5 @@ CREATE TABLE SCHD_SNAP_LOG ( ARCHIVE_KEY VARCHAR(200), -- 原文归档引用(证据层,尚未交付) CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL ); -CREATE INDEX idx_snaplog_cleanup ON SCHD_SNAP_LOG (SCOPE_END, RECV_AT); +CREATE INDEX idx_snaplog_cleanup ON SCHD_SNAP_LOG (SCOPE_END, RECV_AT); -- 保留期清理扫描 CREATE INDEX idx_snaplog_msg ON SCHD_SNAP_LOG (MSG_ID); diff --git a/src/main/resources/db/migration/V2__inbox_lifecycle.sql b/src/main/resources/db/migration/V2__inbox_lifecycle.sql deleted file mode 100644 index e7f86c8..0000000 --- a/src/main/resources/db/migration/V2__inbox_lifecycle.sql +++ /dev/null @@ -1,60 +0,0 @@ --- ===================================================================== --- V2:收报进度与回填记录的存放方式调整 --- --------------------------------------------------------------------- --- 这次迁移解决三个问题,全部只动自有 PostgreSQL。共享 MySQL 那边不建表、 --- 不改结构,边界没有变化。 --- --- 1) 收报读到哪儿了要能记住(新增 INBOX_CURSOR) --- 原先每一轮都从 0 开始扫信箱,靠"有没有处理标记"判断该不该取。这个做法有问题: --- 处理完但还没把标记写回信箱的行(以及永远不会回填的死信)会一直占着每批的 --- 名额,攒够一批之后新消息就再也读不到了。 --- 改成记住"读到哪个 ID 了"(水位),每轮只往后读,并且登记和水位推进放在同一个 --- 事务里——中途崩溃时水位没动,重启后重扫一遍就补齐了。 --- 水位还有个附带规则:如果后面的 ID 缺号,先停下来等(可能是上游还没提交完), --- 等太久就认定它不会来了、跳过去继续,否则水位会卡在第一个空位上再也不动。 --- --- 2) 回填记录并进 PROC_STATE,不再单独建表 --- 处理完有两件事要做:记下终态、把处理标记写回信箱。原先第二步靠一张独立的 --- 待办表,两张表两次写入,中间崩溃就会出现"业务处理完了却没人记得回填"。 --- 现在这两件事是同一条 UPDATE:在一个事务里提交,要么都成、要么都不做。 --- BACKFILL_TODO 因此下线——它保存的 sndr/type/styp/seqn 在回填时从来没用上, --- 回填只需要一个消息 ID。 --- --- 3) 记下收信时间(新增 RECEIVED_AT) --- 有两个用途:判断一条消息等了多久还是没能回填(超期就强制补写), --- 以及统计"最老一条未处理消息收了多久",供运维观察积压。 --- ===================================================================== - --- ① 收报水位:全表只有一行 -CREATE TABLE INBOX_CURSOR ( - CURSOR_ID INT NOT NULL PRIMARY KEY, -- 固定为 1 - COMMITTED_UP_TO BIGINT NOT NULL, -- 水位:到哪个 ID 为止已经全部读进自有库 - HOLE_SINCE TIMESTAMP(6) WITH TIME ZONE, -- 后面那个缺号最早是什么时候发现的;不缺号时为 NULL - UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL -); - --- 初值 0:上线后第一轮会把信箱里所有行重新读一遍,重复登记不会建出第二行, --- 所以顺手把历史上"已入队但没回填"造成的漏读一并补上。 -INSERT INTO INBOX_CURSOR (CURSOR_ID, COMMITTED_UP_TO, HOLE_SINCE, UPDATED_AT) VALUES (1, 0, NULL, now()); - --- ② PROC_STATE:加上收信时间与回填进度(取代 BACKFILL_TODO) -ALTER TABLE PROC_STATE ADD COLUMN RECEIVED_AT TIMESTAMP(6) WITH TIME ZONE; -- 信箱里的接收时间 -ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_AT TIMESTAMP(6) WITH TIME ZONE; -- 有值 = 已确认信箱行带上了标记 -ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_NEXT_AT TIMESTAMP(6) WITH TIME ZONE; -- 有值 = 还欠一次回填,写终态时一起写下 -ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_ATTEMPTS INT NOT NULL DEFAULT 0; -ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_ERROR VARCHAR(512); -- 回填失败的原因,供排查 - --- 存量数据补收信时间:跨库读不到信箱里的 DATE_RECEIVED,只能拿入队时间当兜底。 --- 这个值只会偏晚,宁可晚一点补写标记,也不会提前把还会重放的消息标掉。 -UPDATE PROC_STATE SET RECEIVED_AT = UPDATED_AT WHERE RECEIVED_AT IS NULL; - --- 存量里已经处理完的行补一条回填待办,交给回填扫描处理。 --- 其中已经打过标记的行会被"只写空标记"的条件挡住,重复执行没有副作用。 -UPDATE PROC_STATE SET BACKFILL_NEXT_AT = now() -WHERE STATE IN ('SUCCEEDED', 'SKIPPED', 'DEAD') AND BACKFILL_AT IS NULL; - --- 回填扫描用的索引,只覆盖还没确认回填的行 -CREATE INDEX idx_proc_backfill_due ON PROC_STATE (BACKFILL_NEXT_AT) WHERE BACKFILL_AT IS NULL; - --- ③ BACKFILL_TODO 下线:回填进度已经并进 PROC_STATE,机制只剩一套,索引随表一起释放 -DROP TABLE BACKFILL_TODO; diff --git a/src/main/resources/db/migration/V3__stable_processing_start.sql b/src/main/resources/db/migration/V3__stable_processing_start.sql deleted file mode 100644 index ac200df..0000000 --- a/src/main/resources/db/migration/V3__stable_processing_start.sql +++ /dev/null @@ -1,3 +0,0 @@ --- HOL deadline 使用首次开始处理的稳定时刻,不再被每次重试更新的 UPDATED_AT 推后。 -ALTER TABLE PROC_STATE ADD COLUMN PROCESSING_STARTED_AT TIMESTAMP(6) WITH TIME ZONE; - diff --git a/src/main/resources/db/migration/V4__backfill_closure.sql b/src/main/resources/db/migration/V4__backfill_closure.sql deleted file mode 100644 index f9b161b..0000000 --- a/src/main/resources/db/migration/V4__backfill_closure.sql +++ /dev/null @@ -1,33 +0,0 @@ --- ===================================================================== --- V4:回填闭环 --- --------------------------------------------------------------------- --- 只动自有 PostgreSQL;共享 MySQL 不建表、不改结构。 --- --- 这次迁移解决两个问题: --- --- 1) "放弃回填"需要显式语义(新增 BACKFILL_ABANDONED_AT / BACKFILL_ABANDONED_REASON) --- 原先"信箱行不存在"(MISSING)被当成可重试失败:记录会每 30 秒重试且永不收敛, --- unmarkedTerminal 指标只涨不落。现在把两类情况分开: --- · 运行时查询确认行不存在(确定性结论,重试不改变结果)→ 立即放弃自动重试; --- · 超时/连接失败(暂时性)→ 继续按退避重试,达到上限后放弃自动重试。 --- **放弃 ≠ 标记已确认**:BACKFILL_AT 仍为空,因此不满足"边界内全部行已打标"的 --- 清除前提,库方不应据此清除。放弃行留有原因字段,并支持人工恢复(清标记后重排一次)。 --- --- 2) 扫描公平性(重建索引) --- 原扫描按 MSG_ID 升序取批:最旧的一批永久失败行会持续占满批次,后面的记录永远 --- 轮不到(全局回填饥饿)。现在按 (BACKFILL_ATTEMPTS, MSG_ID) 轮转,并把已放弃的 --- 行排除在扫描之外,保证新行一定能拿到名额。 --- --- 明确不做:**不按年龄做任何存量推断**。"终态 + 从未尝试 + 接收时间早于保留窗"不能 --- 证明信箱行已删除,V2 用 UPDATED_AT 兜底也不构成删除证据;按年龄批量置 abandoned --- 会把仍存在且未标记的旧行永久排除出扫描,反而阻断"边界内全部行已打标"。 --- 存量行一律保留原语义,交由运行时 MISSING 或人工对账处置。 --- ===================================================================== - -ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_ABANDONED_AT TIMESTAMP(6) WITH TIME ZONE; -ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_ABANDONED_REASON VARCHAR(64); - --- 回填扫描:只覆盖仍需自动回填的行,并支持 (attempts, msg_id) 的公平轮转顺序。 -DROP INDEX IF EXISTS idx_proc_backfill_due; -CREATE INDEX idx_proc_backfill_due ON PROC_STATE (BACKFILL_ATTEMPTS, MSG_ID) - WHERE BACKFILL_AT IS NULL AND BACKFILL_ABANDONED_AT IS NULL; diff --git a/src/main/resources/db/migration/V5__cutover_seed.sql b/src/main/resources/db/migration/V5__cutover_seed.sql deleted file mode 100644 index 20f87af..0000000 --- a/src/main/resources/db/migration/V5__cutover_seed.sql +++ /dev/null @@ -1,21 +0,0 @@ --- ===================================================================== --- V5:切流水位播种(一次性、显式) --- --------------------------------------------------------------------- --- 只动自有 PostgreSQL;共享 MySQL 不建表、不改结构。 --- --- 背景:`max-commit-delay` 的空洞老化只在"水位后面出现缺号"时起作用。若第一次对着一个 --- 已有数据的信箱启动(尤其最老分区已被 DROP、MIN(ID) 远大于 1),水位从 0 起会把 ID=1 --- 判成空洞,白等一个老化窗口才前进;而 W=0 又意味着会把保留期内全部存量重新入队。 --- 两种后果都不是代码能替业务决定的,因此这里只提供**显式的、一次性的**播种机制。 --- --- 语义: --- · 默认(未配置 msgx.pipeline.cutover-watermark)不播种,保持既有行为; --- · 四种模式严格区分:min = W:MIN(ID)-1(读当前全部现存行)、zero = W:0(按空洞规则从 0 扫)、 --- max = W:MAX(ID)(跳过当前可见存量)、 = 显式边界; --- · 代码**不做默认选择**,也不会自动退化成 max; --- · 升级实例(已有水位或已有处理记录)**拒绝重新播种**,重新切流必须是显式操作; --- · SEEDED_AT 记录"已播种"这一事实;注意它**为 NULL 不等于"从未消费"**—— --- 已有库新增列后同样是 NULL,因此判据必须叠加"确实没有消费过"。 --- ===================================================================== - -ALTER TABLE INBOX_CURSOR ADD COLUMN SEEDED_AT TIMESTAMP(6) WITH TIME ZONE; diff --git a/src/main/resources/db/migration/V6__enqueued_at.sql b/src/main/resources/db/migration/V6__enqueued_at.sql deleted file mode 100644 index ca95338..0000000 --- a/src/main/resources/db/migration/V6__enqueued_at.sql +++ /dev/null @@ -1,22 +0,0 @@ --- ===================================================================== --- V6:PROC_STATE 增加本地入队时间 ENQUEUED_AT --- --------------------------------------------------------------------- --- 只动自有 PostgreSQL;共享 MySQL 不建表、不改结构。 --- --- 背景(门禁裁决 D6):超期补写期限 `R` 原先比较 `RECEIVED_AT`,那是**库方写入的异地时钟**。 --- 库方时钟前偏时 `RECEIVED_AT < NOW − R` 会提前成立,在「打标即可清除」语义下正好制造 --- 原文被提前清除的丢失窗口(`PRE-4`)。改为比较本系统自己写入的 `ENQUEUED_AT`: --- 本地、非空、与判据用的 `NOW` 同源。`RECEIVED_AT` 退回「仅用于对账与展示」,并允许为 NULL。 --- --- 存量行没有历史入队时间,只能一次性尽力回填:优先信箱接收时间,其次最后一次更新时间。 --- 这是一次推断,不回写业务、不影响任何已有终态,也不改已发布的迁移历史。 --- ===================================================================== - -ALTER TABLE PROC_STATE ADD COLUMN ENQUEUED_AT TIMESTAMP(6) WITH TIME ZONE; - -UPDATE PROC_STATE - SET ENQUEUED_AT = COALESCE(RECEIVED_AT, UPDATED_AT, now()) - WHERE ENQUEUED_AT IS NULL; - -ALTER TABLE PROC_STATE ALTER COLUMN ENQUEUED_AT SET NOT NULL; -ALTER TABLE PROC_STATE ALTER COLUMN ENQUEUED_AT SET DEFAULT now(); diff --git a/src/main/resources/db/migration/V7__drop_processing_started_at.sql b/src/main/resources/db/migration/V7__drop_processing_started_at.sql deleted file mode 100644 index 991aec4..0000000 --- a/src/main/resources/db/migration/V7__drop_processing_started_at.sql +++ /dev/null @@ -1,14 +0,0 @@ --- ===================================================================== --- V7:PROC_STATE 删除 PROCESSING_STARTED_AT --- --------------------------------------------------------------------- --- 只动自有 PostgreSQL;共享 MySQL 不建表、不改结构。 --- --- 该列由 V3 引入(原 HOL deadline 的稳定处理起点),但代码里没有任何判据消费它: --- 处理终态只看尝试上限(`PARAM:msgx.pipeline.max-attempts`)。按精简原则删除列与 --- 全部读写路径,不留"以后也许有用"的死字段;将来若为 `CLM-9` 需要处理开始时间, --- 先写文档再实现。 --- --- 已发布的迁移历史(V1–V6)保持不变。 --- ===================================================================== - -ALTER TABLE PROC_STATE DROP COLUMN PROCESSING_STARTED_AT; diff --git a/src/main/resources/db/migration/V8__req_track_open_unique.sql b/src/main/resources/db/migration/V8__req_track_open_unique.sql deleted file mode 100644 index 46a44aa..0000000 --- a/src/main/resources/db/migration/V8__req_track_open_unique.sql +++ /dev/null @@ -1,34 +0,0 @@ --- ===================================================================== --- V8:REQ_TRACK 开放态部分唯一约束 --- --------------------------------------------------------------------- --- 只动自有 PostgreSQL;共享 MySQL 不建表、不改结构。 --- --- design「记录模型」:同一 (REQ_TYPE, OPERATION_DAY, SENDER) 只允许一个开放请求, --- 且仅对开放状态(PENDING / SENT)生效。这里只建立数据库防线;「新请求原子过期旧请求 --- 并登记自身」的运行时流程仍归 US-08 协调器([G-REQ-TRACK])。 --- --- 存量若已有重复开放组,不静默删数据:带业务键的诊断直接让迁移失败,由人工处置。 --- 保留既有普通索引 idx_req_open(仍服务含关闭态的查询)。 --- ===================================================================== - -DO $$ -DECLARE offenders text; -BEGIN - SELECT string_agg(format('%s/%s/%s', req_type, operation_day, sender), ', ') - INTO offenders - FROM ( - SELECT req_type, operation_day, sender - FROM req_track - WHERE state IN ('PENDING', 'SENT') - GROUP BY req_type, operation_day, sender - HAVING count(*) > 1 - LIMIT 20 - ) dup; - IF offenders IS NOT NULL THEN - RAISE EXCEPTION 'V8 aborted: duplicate open REQ_TRACK rows for (req_type/operation_day/sender): %', offenders; - END IF; -END $$; - -CREATE UNIQUE INDEX uq_req_open - ON req_track (req_type, operation_day, sender) - WHERE state IN ('PENDING', 'SENT'); diff --git a/src/main/resources/db/migration/V9__schd_single_row.sql b/src/main/resources/db/migration/V9__schd_single_row.sql deleted file mode 100644 index 4f9a23a..0000000 --- a/src/main/resources/db/migration/V9__schd_single_row.sql +++ /dev/null @@ -1,27 +0,0 @@ --- ===================================================================== --- V9:KAFKA:schd outbox 单行化(按 FLID) --- --------------------------------------------------------------------- --- 只动自有 PostgreSQL;共享 MySQL 不建表、不改结构。 --- --- design「事件投递」:`KAFKA:schd` 只提供最新状态,outbox 按 `FLID` 单行 upsert。 --- 存量是按版本追加的多行,必须先收敛再建部分唯一约束,否则升级直接失败: --- 每个 (target='KAFKA:schd', partition_key) 留下 `STATE_VERSION DESC, EVENT_ID DESC` --- 的第一行,与运行时的「只进不退」规则一致;`KAFKA:msg` 仍是多行 append-log,不受约束。 --- ===================================================================== - -DELETE FROM msg_event a - USING msg_event b - WHERE a.target = 'KAFKA:schd' - AND b.target = 'KAFKA:schd' - AND a.partition_key = b.partition_key - AND ( - a.state_version < b.state_version - OR (a.state_version = b.state_version AND a.event_id < b.event_id) - ); - -CREATE UNIQUE INDEX uq_schd_event - ON msg_event (target, partition_key) - WHERE target = 'KAFKA:schd'; - --- 读端不再按 (PARTITION_KEY, STATE_VERSION) 合并(单行化 + 条件确认),idx_evt_flid 已无消费者。 -DROP INDEX IF EXISTS idx_evt_flid; diff --git a/src/main/resources/db/migration/oracle11g/README.md b/src/main/resources/db/migration/oracle11g/README.md index a9b386c..ecdbddb 100644 --- a/src/main/resources/db/migration/oracle11g/README.md +++ b/src/main/resources/db/migration/oracle11g/README.md @@ -7,19 +7,21 @@ Flyway 配置**;PG 路径使用 `classpath:db/migration`,两者互不混用 1. 现场 11.2 补丁级别、数据库字符集、DBA 权限清单拿到,且可提供可测试的目标库。 2. JDK 25 × ojdbc 驱动(具体版本)× Flyway Oracle 支持 × 连接池组合在目标库实测通过 - ——不能以"PG 通过"代替 Oracle 验收(Oracle 11g 适配为开放项,docs/flight-state.md §6)。 + ——不能以"PG 通过"代替 Oracle 验收(Oracle 11g 适配为开放项,见 `docs/flight-state.md`)。 3. 11g 的 upsert(MERGE INTO)与绑定顺序适配完成并通过与 PG 同粒度的集成测试后, 才允许把仓储 SQL 切到 11g 方言。先前编译级交付的 `SqlDialect` 方言接缝已随 - V1 基线列名更替(fday/last_message_id → operation_day/last_msg_id)退役删除; - 接入时必须按 V1 基线现有列名重建,禁止直接复用旧接缝。 + 基线列名更替(fday/last_message_id → operation_day/last_msg_id)退役删除; + 接入时必须按基线现有列名重建,禁止直接复用旧接缝。 ## 计划内容 -- V1 起:`V1__flight_state_baseline.sql` 的 11g 等价 DDL—— - PIPELINE_LOCK/PROC_STATE/MSG_EVENT/REQ_TRACK/BACKFILL_TODO/FLIGHT_SCHD + 9 张明细表/SCHD_SNAP_LOG; +- 单基线:`V1__flight_state_baseline.sql` 的 11g 等价 DDL—— + PIPELINE_LOCK/INBOX_CURSOR/PROC_STATE/MSG_EVENT/REQ_TRACK/FLIGHT_SCHD + 9 张明细表/SCHD_SNAP_LOG; `NUMBER`/序列替代 `BIGSERIAL`、`TIMESTAMP WITH TIME ZONE`、`VARCHAR2` BYTE/CHAR 语义钉死。 -- `INSERT ... ON CONFLICT` 改 11g MERGE(OPERATION_DAY 不可变条件,flight-state.md §2.1)。 +- `INSERT ... ON CONFLICT` 改 11g MERGE(OPERATION_DAY 不可变条件,`INV-12`)。 - 空串按 NULL 的语义回归:显式清空的 presence 信息不得被 11g 空串语义吞掉 - (flight-state.md §3.1 显式清空口径)。 + (字段清空语义见 `Q13`;定案前按「未携带不清空」实现,`INV-14`)。 +- 11g 无部分索引:`uq_req_open` / `uq_schd_event` / `idx_proc_backfill_due` 三个带 `WHERE` + 的索引必须换成等价的函数索引或冗余列方案,不能照搬 PG 定义。 版本号与 PG location 各自独立推进,禁止复用版本号语义。 diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt index abfff44..433d001 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt @@ -10,9 +10,10 @@ import org.junit.jupiter.api.Test import java.sql.DriverManager /** - * 在真实 PostgreSQL 上跑一遍迁移,确认结果符合预期: - * V1 基线、V2 信箱生命周期和 V3 稳定处理起点都能成功执行,迁移记录显示成功,该建的表和单行种子 - * (PIPELINE_LOCK、INBOX_CURSOR)都在,回填相关字段进了 PROC_STATE、BACKFILL_TODO 已下线。 + * 在真实 PostgreSQL 上跑一遍迁移,确认结果符合预期:单基线 `V1__flight_state_baseline.sql` + * 执行成功且是唯一的迁移记录,该建的表和单行种子(PIPELINE_LOCK、INBOX_CURSOR)都在, + * 回填事实与收报水位都落在基线里,`BACKFILL_TODO`、`idx_evt_flid`、`PROC_STATE` 的处理开始 + * 时间列都不复存在——这些正是原 V2–V10 合并后的净结构。 * * 没有可用的 PostgreSQL 时跳过(不假装通过)。 */ @@ -41,7 +42,7 @@ class FlywayMigrationTest { DriverManager.getConnection(url, user, pass).use { conn -> conn.createStatement().use { stmt -> - // 迁移记录:V1 基线 + V2 信箱生命周期 + V3 稳定处理起点 + V4 回填闭环 + V5 切流播种 + // 单基线:原 V1–V10 已合并为一条迁移,历史里只应留下 V1 stmt.executeQuery( "SELECT version, script, success FROM flyway_schema_history ORDER BY installed_rank ASC", ).use { rs -> @@ -49,29 +50,13 @@ class FlywayMigrationTest { while (rs.next()) { records.add(Triple(rs.getString("version"), rs.getString("script"), rs.getBoolean("success"))) } - assertTrue(records.size >= 9, "flyway_schema_history must record all migrations") + assertEquals(1, records.size, "迁移链必须收敛为单基线,不得再有 V2+ 条目") assertEquals("1", records[0].first) assertEquals("V1__flight_state_baseline.sql", records[0].second) - assertEquals("2", records[1].first) - assertEquals("V2__inbox_lifecycle.sql", records[1].second) - assertEquals("3", records[2].first) - assertEquals("V3__stable_processing_start.sql", records[2].second) - assertEquals("4", records[3].first) - assertEquals("V4__backfill_closure.sql", records[3].second) - assertEquals("5", records[4].first) - assertEquals("V5__cutover_seed.sql", records[4].second) - assertEquals("6", records[5].first) - assertEquals("V6__enqueued_at.sql", records[5].second) - assertEquals("7", records[6].first) - assertEquals("V7__drop_processing_started_at.sql", records[6].second) - assertEquals("8", records[7].first) - assertEquals("V8__req_track_open_unique.sql", records[7].second) - assertEquals("9", records[8].first) - assertEquals("V9__schd_single_row.sql", records[8].second) assertTrue(records.all { it.third }) } - // 全表就绪(V2 后 backfill_todo 下线) + // 全表就绪(回填事实并入 PROC_STATE,基线不再建 BACKFILL_TODO) stmt.executeQuery( "SELECT table_name FROM information_schema.tables WHERE table_schema = 'public'", ).use { rs -> @@ -110,20 +95,20 @@ class FlywayMigrationTest { ) } - // V7:PROC_STATE 不再保留处理开始时间(无判据消费) + // 基线不含处理开始时间(原 V3 引入、V7 删除,合并后不得残留) stmt.executeQuery( "SELECT count(*) FROM information_schema.columns WHERE table_name = 'proc_state' " + "AND column_name = 'processing_started_at'", ).use { rs -> assertTrue(rs.next()) - assertEquals(0, rs.getInt(1), "V7 必须已删除 PROCESSING_STARTED_AT") + assertEquals(0, rs.getInt(1), "基线必须不含 PROCESSING_STARTED_AT") } - // V8:REQ_TRACK 开放态部分唯一约束(仅 PENDING/SENT;关闭态可并存) + // 基线:REQ_TRACK 开放态部分唯一约束(仅 PENDING/SENT;关闭态可并存) stmt.executeQuery( "SELECT indexdef FROM pg_indexes WHERE tablename = 'req_track' AND indexname = 'uq_req_open'", ).use { rs -> - assertTrue(rs.next(), "V8 必须建立开放态唯一索引 uq_req_open") + assertTrue(rs.next(), "基线必须建立开放态唯一索引 uq_req_open") val indexDef = rs.getString(1) assertTrue(indexDef.contains("UNIQUE"), "uq_req_open 必须是唯一索引:$indexDef") assertTrue(indexDef.contains("WHERE"), "uq_req_open 必须是仅约束开放态的部分索引:$indexDef") @@ -150,11 +135,11 @@ class FlywayMigrationTest { "VALUES ('RQFD-NONE', DATE '2026-09-12', 'RMS', 'PENDING', now())", ) - // V9:KAFKA:schd 每个 FLID 单行(部分唯一索引),KAFKA:msg 仍可多行 + // 基线:KAFKA:schd 每个 FLID 单行(部分唯一索引),KAFKA:msg 仍可多行 stmt.executeQuery( "SELECT indexdef FROM pg_indexes WHERE tablename = 'msg_event' AND indexname = 'uq_schd_event'", ).use { rs -> - assertTrue(rs.next(), "V9 必须建立 schd 单行唯一索引 uq_schd_event") + assertTrue(rs.next(), "基线必须建立 schd 单行唯一索引 uq_schd_event") val indexDef = rs.getString(1) assertTrue(indexDef.contains("UNIQUE"), "uq_schd_event 必须是唯一索引:$indexDef") assertTrue(indexDef.contains("WHERE"), "uq_schd_event 必须是仅约束 schd 的部分索引:$indexDef") @@ -181,15 +166,24 @@ class FlywayMigrationTest { "SELECT count(*) FROM pg_indexes WHERE tablename = 'msg_event' AND indexname = 'idx_evt_flid'", ).use { rs -> assertTrue(rs.next()) - assertEquals(0, rs.getInt(1), "V9 必须删除已无消费者的 idx_evt_flid") + assertEquals(0, rs.getInt(1), "基线必须不含已无消费者的 idx_evt_flid") } - // V6:入队时间是超期判据 R 的比较对象,必须非空(received_at 则允许为 NULL) + // 基线:投递确认时间 SENT_AT 是保留期判定的唯一基准 + stmt.executeQuery( + "SELECT count(*) FROM information_schema.columns WHERE table_name = 'msg_event' " + + "AND column_name = 'sent_at'", + ).use { rs -> + assertTrue(rs.next()) + assertEquals(1, rs.getInt(1), "基线必须含 MSG_EVENT.SENT_AT") + } + + // 基线:入队时间是超期判据 R 的比较对象,必须非空(received_at 则允许为 NULL) stmt.executeQuery( "SELECT is_nullable FROM information_schema.columns " + "WHERE table_name = 'proc_state' AND column_name = 'enqueued_at'", ).use { rs -> - assertTrue(rs.next(), "V6 必须已加上 ENQUEUED_AT 列") + assertTrue(rs.next(), "基线必须已加上 ENQUEUED_AT 列") assertEquals("NO", rs.getString("is_nullable"), "ENQUEUED_AT 必须非空") }