chore(db): 压缩 Flyway 迁移链为单基线

原 V1–V10 的净结构全部合并进 V1__flight_state_baseline.sql,迁移链收敛为一条。
表集合、逐列类型/非空/默认值、表级约束、索引与种子行经重放比对与旧链条一致;
仅丢弃 5 条只对既有库有意义的存量数据语句(RECEIVED_AT/ENQUEUED_AT 回填、
schd 与开放请求去重、SENT_AT 回写)。

已按旧链迁移过的库版本链与校验和对不上,必须重建 schema 或删除数据卷后重跑,
禁止手工 repair 或改写 flyway_schema_history。

- 删除 V2–V10 共 9 个迁移文件
- 注释统一范式并改用稳定 ID(INV-x/C-x/G-x),修正 KAFKA_MSG/KAFKA_SCHD 与实际值不符
- 同步 docs/reference.md、docs/user-stories.md、oracle11g/README.md、README、application.yml、compose.yaml
- FlywayMigrationTest 改为断言单基线,并补 MSG_EVENT.SENT_AT 覆盖
This commit is contained in:
windyboy
2026-09-13 16:31:14 +08:00
parent 62f61ad9e4
commit 800d7617f2
17 changed files with 163 additions and 356 deletions
+3 -1
View File
@@ -64,7 +64,9 @@
## 数据库初始化(ACM2-12 口径) ## 数据库初始化(ACM2-12 口径)
**自有 PostgreSQL**(唯一自有库):Flyway 执行 `db/migration/V1__flight_state_baseline.sql`,建立 **自有 PostgreSQL**(唯一自有库):Flyway 执行 `db/migration/V1__flight_state_baseline.sql`,建立
航班当前态、明细表、处理终态与 outbox 等表(PG 方言)。全新库直接执行即可,无 legacy 前置。 航班当前态、明细表、处理终态与 outbox 等表(PG 方言)。这是**单基线**:原 V2–V10 的净结构已
合并其中,全新库直接执行即可,无 legacy 前置。已按旧链(V1–V10)迁移过的库版本链与校验和都
对不上,必须重建 schema 或删除数据卷后重跑,禁止手工 `repair` 或改写 `flyway_schema_history`
**共享 MySQLcdairport,他人系统库)**:本系统**不建表/schema**,仅信箱 DML——上游外部写 **共享 MySQLcdairport,他人系统库)**:本系统**不建表/schema**,仅信箱 DML——上游外部写
`CMINMSGS`;本系统 JDBC 轮询读 + 处理回填;出站写 `COUTMSGS`(他人读取发送);表结构与 `CMINMSGS`;本系统 JDBC 轮询读 + 处理回填;出站写 `COUTMSGS`(他人读取发送);表结构与
+1 -1
View File
@@ -30,7 +30,7 @@ services:
retries: 10 retries: 10
start_period: 15s start_period: 15s
# 由 Flyway 自动执行 V1.0.0 与 V1.1.0 迁移 # 由 Flyway 自动执行单基线迁移(db/migration/V1__flight_state_baseline.sql
postgres: postgres:
image: postgres:17-alpine image: postgres:17-alpine
container_name: msgx-dev-postgres container_name: msgx-dev-postgres
+1 -1
View File
@@ -103,7 +103,7 @@
| 持久化与恢复 | `infra/persistence/``infra/retry/``ProcFailure` / `ReplayService` / `FailureScheduler` | | 持久化与恢复 | `infra/persistence/``infra/retry/``ProcFailure` / `ReplayService` / `FailureScheduler` |
| 启停与配置 | `PipelineLifecycle.kt``config/PipelineProps.kt``config/HistoryProps.kt``config/OperationDayProps.kt`(含运营日时区启动自检) | | 启停与配置 | `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` | | 指标与健康 | `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. 错误分类与重放白名单 ## 4. 错误分类与重放白名单
+1 -1
View File
@@ -190,7 +190,7 @@
4. 重复补偿效果幂等,保留稳定的完成时间与审计;重放后的新处理结果不能被旧回填任务覆盖。非法报文缺 META 时也有明确回填方式。 4. 重复补偿效果幂等,保留稳定的完成时间与审计;重放后的新处理结果不能被旧回填任务覆盖。非法报文缺 META 时也有明确回填方式。
5. 影子模式禁写,双跑仅一个系统持有标记写权;暴露 PG 终态、回填状态、积压、最老年龄与持续失败告警。 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)「保留与清除」为准。 **前置**:US-03 终态接口;Q7、共享库更新权限。覆盖四类终态、事务回滚、重复补偿和重放竞争;生命周期与超期补写以 [design.md](design.md)「中断恢复」「回填」为准,清除口径以 [contracts.md](contracts.md)「保留与清除」为准。
+2 -2
View File
@@ -47,8 +47,8 @@ micronaut:
# 管理端点(U03/N32):/env、/beans 默认 sensitive;仅开发/影子环境放开——见 application-dev.yml # 管理端点(U03/N32):/env、/beans 默认 sensitive;仅开发/影子环境放开——见 application-dev.yml
# 自有 PostgreSQL:全部内部状态(PIPELINE_LOCK/PROC_STATE/MSG_EVENT/REQ_TRACK/INBOX_CURSOR/ # 自有 PostgreSQL:全部内部状态(PIPELINE_LOCK/PROC_STATE/MSG_EVENT/REQ_TRACK/INBOX_CURSOR/
# FLIGHT_SCHD 及明细表/SCHD_SNAP_LOG;迁移 db/migration/V1__flight_state_baseline.sql + # FLIGHT_SCHD 及明细表/SCHD_SNAP_LOG;迁移为单基线 db/migration/V1__flight_state_baseline.sql
# V2__inbox_lifecycle.sql——V2 起回填意图并入 PROC_STATEBACKFILL_TODO 已下线)。 # 原 V2–V10 的净结构已合并其中,回填事实与处理终态同表同行)。
# 数据层实装前 enabled=falsestub 模式不建连)。 # 数据层实装前 enabled=falsestub 模式不建连)。
datasources: datasources:
default: default:
@@ -1,13 +0,0 @@
-- =====================================================================
-- V10MSG_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';
@@ -1,23 +1,34 @@
-- ===================================================================== -- =====================================================================
-- 航班运行数据接入与当前状态管理 · 全量基线 -- 自有 PostgreSQL 全量基线(单迁移链条目)
-- --------------------------------------------------------------------- -- ---------------------------------------------------------------------
-- 权威设计:docs/flight-state.md(审计定稿)。本脚本整体取代历史迁移 -- 范围:本迁移只作用于自有 PostgreSQL。共享 MySQL 信箱(CMINMSGS/COUTMSGS)是他人
-- V1.0.0V1.4.0 已删除),按 §3.2 表职责建立全部分层: -- 系统库,只做契约内 DML,不建表、不改结构(C-14)。
-- · 决策层:FLIGHT_SCHD + 8 张资源明细表 + FLIGHT_ROUTE_POINT —— 承担正确性; --
-- · 管道层:PIPELINE_LOCK / PROC_STATE / MSG_EVENT(outbox) / REQ_TRACK / BACKFILL_TODO -- 本文件原为 V1 基线;原 V2–V10 的净结构已全部合并进来,版本链收敛为一条:
-- · 留痕层:SCHD_SNAP_LOG —— 只追加、不参与决策、可重建 -- · INBOX_CURSOR 及 SEEDED_AT(原 V2/V5
-- · 证据层:报文原文归档 —— 冷路径,尚未交付(ARCHIVE_KEY 仅留引用位)。 -- · PROC_STATE 的收信/入队时间与回填事实列,独立待办表 BACKFILL_TODO 不再存在
-- 核心定案(对照文档章节): -- (原 V2/V3/V4/V6/V7:回填意图并入 PROC_STATE,处理开始时间列加后即删);
-- · 身份:FLID 主键;OPERATION_DAY 一经确定不可变(§3.1/§5.3,应用层校验); -- · REQ_TRACK 开放态部分唯一索引 uq_req_open(原 V8);
-- STATE ∈ {ACTIVE, DELETED},无 ARCHIVED —— 物理清除只发生在历史归档成功之后(§8.2); -- · MSG_EVENT 的 schd 单行唯一索引 uq_schd_event、SENT_AT,以及已无消费者的
-- · 无名单差删:主链路无集合删除;删除入口只有 FDEL(§6.2)与历史归档(§8.2); -- idx_evt_flid 下线(原 V9/V10)。
-- · 版本:每次成功写入推进 STATE_VERSIONKafka 判旧 = (FLID, STATE_VERSION, UPDATED_AT)(§7.3); -- 原 V2–V10 中只对既有库有意义的数据搬迁语句(存量 RECEIVED_AT / ENQUEUED_AT 回填、
-- · 出站:COUTMSGS 是共享 MySQL 信箱(他人系统库,本系统不建表,迁移不覆盖)。 -- 重复开放请求与重复 schd 事件去重、存量 SENT_AT 回写)一律不再保留——基线只服务
-- 类型口径:TIMESTAMP(6) WITH TIME ZONE 统一 UTC 语义;无 JSONB/CLOB 依赖; -- 全新库,而 Flyway 的迁移链一旦发布即不可改写,故这些语句的语义只能在此说明。
-- BIGSERIAL 为 PostgreSQL 方言(Oracle 适配见 §10,未定前不预设)。 --
-- 升级路径:已应用过 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 ( CREATE TABLE PIPELINE_LOCK (
LOCK_ID INT NOT NULL PRIMARY KEY, -- 恒为 1 LOCK_ID INT NOT NULL PRIMARY KEY, -- 恒为 1
UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL 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()); INSERT INTO PIPELINE_LOCK (LOCK_ID, UPDATED_AT) VALUES (1, now());
-- ② 处理伴生状态:与共享信箱 CMINMSGS 一一对应(MSG_ID = CMINMSGS_ID)。 -- ② 处理伴生状态:与共享信箱 CMINMSGS 一一对应(MSG_ID = CMINMSGS_ID)。
-- 兼作快照重放判定(§5.1 步骤 2):MSG_ID 已有成功终态重放,直接记幂等成功。 -- 一信一行、一身份一记录(INV-9);MSG_ID 已有成功终态即判重放,直接记幂等成功。
-- 处理终态与回填意图在本表同表同行:不存在第二处落账,也不存在独立的回填待办表(INV-8)。
-- 状态不可逆:已提交的 SUCCEEDED 不因回填或投递失败回改(INV-6)。
CREATE TABLE PROC_STATE ( CREATE TABLE PROC_STATE (
MSG_ID BIGINT NOT NULL PRIMARY KEY, MSG_ID BIGINT NOT NULL PRIMARY KEY,
STATE VARCHAR(16) NOT NULL, -- PENDING/FAILED/SUCCEEDED/SKIPPED/DEAD STATE VARCHAR(16) NOT NULL, -- PENDING/FAILED/SUCCEEDED/SKIPPED/DEADARCHIVED 属未建归档能力 G-PROC-HST
IDENTITY_KEY VARCHAR(200), -- SNDR|TYPE|STYP|SEQNdecode 后首次绑定;SEQN 同会话顺序观测 §5.5) IDENTITY_KEY VARCHAR(200), -- SNDR|TYPE|STYP|SEQN,解码后首次绑定
ATTEMPTS INT NOT NULL DEFAULT 0, ATTEMPTS INT NOT NULL DEFAULT 0,
NEXT_ATTEMPT_AT TIMESTAMP(6) WITH TIME ZONE, NEXT_ATTEMPT_AT TIMESTAMP(6) WITH TIME ZONE, -- 退避到期时刻;未到期不得被后续消息越过(INV-3)
ERROR_CLASS VARCHAR(20), -- MALFORMED/PROTOCOL/INFRA/UNSUPPORTED/EXHAUSTED(§9 ERROR_CLASS VARCHAR(20), -- MALFORMED/PROTOCOL/INFRA/UNSUPPORTED/EXHAUSTED
LAST_ERROR VARCHAR(1000), LAST_ERROR VARCHAR(1000),
UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, 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) 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); -- 主泵队头查询(严格 FIFOINV-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 ( CREATE TABLE FLIGHT_SCHD (
FLID VARCHAR(32) NOT NULL PRIMARY KEY, -- AODB 实例 IDNumber(1-12)SIS §3.16.2 FLID VARCHAR(32) NOT NULL PRIMARY KEY, -- AODB 实例 IDNumber(1-12));不得由航班号或资源号推断
OPERATION_DAY DATE NULL, -- 运营保障日;一经确定不可变(§3.5);NULL=尚未被快照收录 OPERATION_DAY DATE NULL, -- 运营保障日:未由日计划收录时为 NULL;非空后不可变(INV-12
STATE VARCHAR(8) NOT NULL, -- ACTIVE/DELETED(§3.1,无第三态) STATE VARCHAR(8) NOT NULL, -- ACTIVE/DELETED,无第三态;删除只由 FDEL 或受控历史清理触发(INV-15
STATE_VERSION BIGINT NOT NULL DEFAULT 0, -- 每次成功写入 +1(§4 STATE_VERSION BIGINT NOT NULL DEFAULT 0, -- 每次成功写入 +1,重复消息不重复推进(INV-13
LAST_MSG_ID BIGINT NULL, -- 最近一次成功写入的消息 ID LAST_MSG_ID BIGINT NULL, -- 最近一次成功写入的消息 ID
-- ==== SCHD.FLTR 标量字段(保持 AODB 字符串原样;缺失语义由处理层决定,§5.4==== -- ==== SCHD.FLTR 标量字段(保持 AODB 字符串原样;未携带不隐式清空,INV-14====
ALCD VARCHAR(64) NULL, -- 航空公司代码 ALCD VARCHAR(64) NULL, -- 航空公司代码
ALSC VARCHAR(64) NULL, -- 航空公司简称 ALSC VARCHAR(64) NULL, -- 航空公司简称
FLNO VARCHAR(64) NULL, -- 航班号(展示用,不参与身份判定 §3.1 FLNO VARCHAR(64) NULL, -- 航班号(展示用,不参与身份判定)
MVIN VARCHAR(64) NULL, -- 进离港标识(A/D MVIN VARCHAR(64) NULL, -- 进离港标识(A/D
SODT VARCHAR(64) NULL, -- 计划运行时间(ddMMMyyHHmmOPERATION_DAY 计算源 §3.5 SODT VARCHAR(64) NULL, -- 计划运行时间(ddMMMyyHHmmOPERATION_DAY 的推导源
FLTY VARCHAR(64) NULL, -- 航班类型 FLTY VARCHAR(64) NULL, -- 航班类型
FLIN VARCHAR(64) NULL, -- 国内/国际/混合 FLIN VARCHAR(64) NULL, -- 国内/国际/混合
ACFT VARCHAR(64) NULL, -- 机型 ACFT VARCHAR(64) NULL, -- 机型
RENO VARCHAR(64) NULL, -- 机尾号 RENO VARCHAR(64) NULL, -- 机尾号
TAOP VARCHAR(64) NULL, -- 过站实际承运人代码(§3.4过站关联,非代码共享) TAOP VARCHAR(64) NULL, -- 过站实际承运人代码(过站关联,非代码共享)
TAFL VARCHAR(64) NULL, -- 过站实际承运航班号 TAFL VARCHAR(64) NULL, -- 过站实际承运航班号
TAID VARCHAR(64) NULL, -- 过站关联航班 FLID TAID VARCHAR(64) NULL, -- 过站关联航班 FLID
TRML VARCHAR(64) NULL, -- 航站楼 TRML VARCHAR(64) NULL, -- 航站楼
MAXP VARCHAR(64) NULL, -- 最大旅客数 MAXP VARCHAR(64) NULL, -- 最大旅客数
CSOP VARCHAR(64) NULL, -- 共享承运人代码(§3.4 CSOP VARCHAR(64) NULL, -- 共享承运人代码
CSFT 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, -- 预计时间 ESTT VARCHAR(64) NULL, -- 预计时间
ACTT VARCHAR(64) NULL, -- 实际时间 ACTT VARCHAR(64) NULL, -- 实际时间
STND VARCHAR(64) NULL, -- 备降站/机位 STND VARCHAR(64) NULL, -- 备降站/机位
PHAG 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, -- 备注 REMC VARCHAR(256) NULL, -- 备注
BOTM VARCHAR(64) NULL, -- 摆渡车时间 BOTM VARCHAR(64) NULL, -- 摆渡车时间
LACL VARCHAR(64) NULL, -- 行李确认时间 LACL VARCHAR(64) NULL, -- 行李确认时间
@@ -87,22 +128,23 @@ CREATE TABLE FLIGHT_SCHD (
EXSR VARCHAR(256) NULL, -- 例外原因 EXSR VARCHAR(256) NULL, -- 例外原因
FTSS VARCHAR(64) NULL, -- 航班状态 FTSS VARCHAR(64) NULL, -- 航班状态
PEDT VARCHAR(64) NULL, -- 前序估计时间 PEDT VARCHAR(64) NULL, -- 前序估计时间
NAAT VARCHAR(64) NULL, -- 到港终态时间(§8.1业务含义待术语表确认 §10 NAAT VARCHAR(64) NULL, -- 到港终态时间(业务含义待术语表确认)
NEAT VARCHAR(64) NULL, -- 离港终态时间(§8.1业务含义待术语表确认 §10 NEAT VARCHAR(64) NULL, -- 离港终态时间(业务含义待术语表确认)
PADT VARCHAR(64) NULL, -- 旅客登机时间 PADT VARCHAR(64) NULL, -- 旅客登机时间
CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL,
UPDATED_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_opday ON FLIGHT_SCHD (OPERATION_DAY);
CREATE INDEX idx_flight_schd_state ON FLIGHT_SCHD (STATE); 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,行级条件更新即满足; -- WHERE OPERATION_DAY IS NULL OR OPERATION_DAY = :day,行级条件更新即满足;
-- 不使用触发器(Oracle/PG 双方言成本)。 -- 不使用触发器(Oracle/PG 双方言成本)。
-- 资源明细表 ×8 + 路线点表每 FLID 多行;明细集合按组先删后插 §3.3 -- 资源明细表 ×8 + 路线点表每 FLID 多行;每次完整状态写入按该 FLID 先删后插,
-- ORDINAL 保留输入顺序,SOURCE_SEQ 只存上游序号 §3.2 -- 以完整合并结果为准(INV-14)。ORDINAL 是持久化顺序(从 1 起),SOURCE_SEQ 是上游
-- 删除一律标记在主行 STATE,明细物理清除只发生在历史归档 §8.2) -- 序号(允许为空或重复);相同资源号不代表同一条分配,禁止按资源号去重。删除一律
-- ④-1 登机门 GTDT -- 标记在主行 STATE,明细物理清除只发生在受控历史归档(INV-15)。
-- ⑤-1 登机门 GTDT
CREATE TABLE FLIGHT_GATE ( CREATE TABLE FLIGHT_GATE (
FLID VARCHAR(32) NOT NULL, FLID VARCHAR(32) NOT NULL,
ORDINAL INT NOT NULL, ORDINAL INT NOT NULL,
@@ -117,7 +159,7 @@ CREATE TABLE FLIGHT_GATE (
UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL,
PRIMARY KEY (FLID, ORDINAL) PRIMARY KEY (FLID, ORDINAL)
); );
-- -2 值机柜台 CKDT -- -2 值机柜台 CKDT
CREATE TABLE FLIGHT_CHECKIN ( CREATE TABLE FLIGHT_CHECKIN (
FLID VARCHAR(32) NOT NULL, FLID VARCHAR(32) NOT NULL,
ORDINAL INT NOT NULL, ORDINAL INT NOT NULL,
@@ -133,7 +175,7 @@ CREATE TABLE FLIGHT_CHECKIN (
UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL,
PRIMARY KEY (FLID, ORDINAL) PRIMARY KEY (FLID, ORDINAL)
); );
-- -3 行李转盘 CLDT -- -3 行李转盘 CLDT
CREATE TABLE FLIGHT_BELT ( CREATE TABLE FLIGHT_BELT (
FLID VARCHAR(32) NOT NULL, FLID VARCHAR(32) NOT NULL,
ORDINAL INT NOT NULL, ORDINAL INT NOT NULL,
@@ -149,7 +191,7 @@ CREATE TABLE FLIGHT_BELT (
UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL,
PRIMARY KEY (FLID, ORDINAL) PRIMARY KEY (FLID, ORDINAL)
); );
-- -4 计划机位 PSDT -- -4 计划机位 PSDT
CREATE TABLE FLIGHT_STAND_PLAN ( CREATE TABLE FLIGHT_STAND_PLAN (
FLID VARCHAR(32) NOT NULL, FLID VARCHAR(32) NOT NULL,
ORDINAL INT NOT NULL, ORDINAL INT NOT NULL,
@@ -161,23 +203,23 @@ CREATE TABLE FLIGHT_STAND_PLAN (
UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL,
PRIMARY KEY (FLID, ORDINAL) PRIMARY KEY (FLID, ORDINAL)
); );
-- -5 滑槽 CHDT -- -5 滑槽 CHDT
CREATE TABLE FLIGHT_CHUTE ( CREATE TABLE FLIGHT_CHUTE (
FLID VARCHAR(32) NOT NULL, FLID VARCHAR(32) NOT NULL,
ORDINAL INT NOT NULL, ORDINAL INT NOT NULL,
SOURCE_SEQ VARCHAR(64), SOURCE_SEQ VARCHAR(64),
CHUT VARCHAR(64), CHUT VARCHAR(64),
CHCLS VARCHAR(64), CHCLS VARCHAR(64),
PCBT VARCHAR(64), -- 计划开始SIS 线格式 PCBT PCBT VARCHAR(64), -- 计划开始
PCET VARCHAR(64), -- 计划结束PCET PCET VARCHAR(64), -- 计划结束
CBTM VARCHAR(64), -- 实际开始CBTM CBTM VARCHAR(64), -- 实际开始
CETM VARCHAR(64), -- 实际结束CETM CETM VARCHAR(64), -- 实际结束
CHTYP VARCHAR(64), CHTYP VARCHAR(64),
CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL,
UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL,
PRIMARY KEY (FLID, ORDINAL) PRIMARY KEY (FLID, ORDINAL)
); );
-- -6 延误 DELY(业务上任意时刻仅 1 个有效延误,保留多行能力以无损承接) -- -6 延误 DELY(业务上任意时刻仅 1 个有效延误,保留多行能力以无损承接)
CREATE TABLE FLIGHT_DELAY ( CREATE TABLE FLIGHT_DELAY (
FLID VARCHAR(32) NOT NULL, FLID VARCHAR(32) NOT NULL,
ORDINAL INT NOT NULL, ORDINAL INT NOT NULL,
@@ -190,7 +232,7 @@ CREATE TABLE FLIGHT_DELAY (
UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL,
PRIMARY KEY (FLID, ORDINAL) PRIMARY KEY (FLID, ORDINAL)
); );
-- -7 靠撤桥 ABTM(桥号 ABDGA/D 各一条) -- -7 靠撤桥 ABTM(桥号 ABDGA/D 各一条)
CREATE TABLE FLIGHT_BRIDGE_OP ( CREATE TABLE FLIGHT_BRIDGE_OP (
FLID VARCHAR(32) NOT NULL, FLID VARCHAR(32) NOT NULL,
ORDINAL INT NOT NULL, ORDINAL INT NOT NULL,
@@ -202,7 +244,7 @@ CREATE TABLE FLIGHT_BRIDGE_OP (
UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL,
PRIMARY KEY (FLID, ORDINAL) PRIMARY KEY (FLID, ORDINAL)
); );
-- -8 轮挡 CHOT(机位 CHIDON/OFF 各一条) -- -8 轮挡 CHOT(机位 CHIDON/OFF 各一条)
CREATE TABLE FLIGHT_CHOCK_OP ( CREATE TABLE FLIGHT_CHOCK_OP (
FLID VARCHAR(32) NOT NULL, FLID VARCHAR(32) NOT NULL,
ORDINAL INT NOT NULL, ORDINAL INT NOT NULL,
@@ -214,7 +256,7 @@ CREATE TABLE FLIGHT_CHOCK_OP (
UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL,
PRIMARY KEY (FLID, ORDINAL) PRIMARY KEY (FLID, ORDINAL)
); );
-- -9 路线点 ROUT/ERUT 共用(ROUTE_KIND 区分 §3.210 类集合 9 张表 -- -9 路线点 ROUT/ERUT 共用(ROUTE_KIND 区分两类;主键含该列,避免序号冲突
CREATE TABLE FLIGHT_ROUTE_POINT ( CREATE TABLE FLIGHT_ROUTE_POINT (
FLID VARCHAR(32) NOT NULL, FLID VARCHAR(32) NOT NULL,
ORDINAL INT 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_bridge_flid ON FLIGHT_BRIDGE_OP (FLID);
CREATE INDEX idx_flight_chock_flid ON FLIGHT_CHOCK_OP (FLID); CREATE INDEX idx_flight_chock_flid ON FLIGHT_CHOCK_OP (FLID);
-- 统一投递事件 outbox§7.3KAFKA_MSG 只通知变化;KAFKA_SCHD 发整态 -- 统一投递事件 outboxINV-10、INV-17):KAFKA:schd 发整态、KAFKA:msg 只通知变化
-- tombstone 仅在 ACTIVE→DELETED 时与删除同事务登记,投递失败持续重试 -- tombstone 仅在 ACTIVE→DELETED 时与删除同事务登记,投递失败按退避持续重试
-- 目标级全序投递是当前实现,按 FLID 保序见 CLM-7。
CREATE TABLE MSG_EVENT ( CREATE TABLE MSG_EVENT (
EVENT_ID BIGSERIAL PRIMARY KEY, EVENT_ID BIGSERIAL PRIMARY KEY, -- 对 KAFKA:msg 是稳定事件身份并决定投递顺序;对 KAFKA:schd 是每次接受 upsert 时替换的写代次
TARGET VARCHAR(30) NOT NULL, -- KAFKA_MSG / KAFKA_SCHD TARGET VARCHAR(30) NOT NULL, -- KAFKA:msg / KAFKA:schd
PARTITION_KEY VARCHAR(32) NOT NULL, -- 恒为 FLID PARTITION_KEY VARCHAR(32) NOT NULL, -- 当前恒为 FLIDQ4 定案前为假定,C-29
EVENT_TYPE VARCHAR(16) NOT NULL DEFAULT 'UPSERT', -- UPSERT / TOMBSTONE EVENT_TYPE VARCHAR(16) NOT NULL DEFAULT 'UPSERT', -- UPSERT / TOMBSTONE
STATE_VERSION BIGINT NOT NULL, -- 发布时航班版本;聚合按 FLID 取最新(§7.3 STATE_VERSION BIGINT NOT NULL, -- 发布时航班版本;schd 聚合按 FLID 只进不退
PAYLOAD_JSON TEXT NOT NULL, -- TOMBSTONE 时至少含 FLID/STATE_VERSION/DELETED PAYLOAD_JSON TEXT NOT NULL, -- TOMBSTONE 时至少含 FLID/STATE_VERSION/DELETED
STATE VARCHAR(16) NOT NULL, -- PENDING/SENT/DEAD STATE VARCHAR(16) NOT NULL, -- PENDING/SENT/DEAD
ATTEMPTS INT NOT NULL DEFAULT 0, ATTEMPTS INT NOT NULL DEFAULT 0,
NEXT_ATTEMPT_AT TIMESTAMP(6) WITH TIME ZONE, NEXT_ATTEMPT_AT TIMESTAMP(6) WITH TIME ZONE,
ERROR_CLASS VARCHAR(20), ERROR_CLASS VARCHAR(20),
LAST_ERROR VARCHAR(1000), 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_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 请求,按(运营日、发送方、请求类型) -- 请求状态机:只有 RESP 完成 RQFD 请求,按(运营日、发送方、请求类型)匹配最新一条
-- 匹配最新一条 PENDING同类请求只留一条有效,新请求置旧为 EXPIRED -- 开放请求。同类请求只留一条有效,新请求置旧请求为 EXPIRED。登记、超时与应答匹配
-- 尚未实现(G-REQ-TRACK、C-24)。
CREATE TABLE REQ_TRACK ( CREATE TABLE REQ_TRACK (
REQ_ID BIGSERIAL PRIMARY KEY, REQ_ID BIGSERIAL PRIMARY KEY,
REQ_TYPE VARCHAR(20) NOT NULL, REQ_TYPE VARCHAR(20) NOT NULL,
@@ -268,26 +316,14 @@ CREATE TABLE REQ_TRACK (
COMPLETED_AT TIMESTAMP(6) WITH TIME ZONE, COMPLETED_AT TIMESTAMP(6) WITH TIME ZONE,
CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL 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 -- ⑧ SCHD 快照留痕:事务外追加,只追加留痕、不参与决策;一行=一次尝试,重放也记。
-- 目标形态 = 业务事务内预登记,消除「提交后写待办」的崩溃窗口,§10 当前偏差) -- RESULT 与 FLAGS 分列(可「成功且告警」);写失败只记指标。
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) 清理
CREATE TABLE SCHD_SNAP_LOG ( CREATE TABLE SCHD_SNAP_LOG (
LOG_ID BIGSERIAL PRIMARY KEY, LOG_ID BIGSERIAL PRIMARY KEY,
MSG_ID BIGINT NOT NULL, MSG_ID BIGINT NOT NULL,
@@ -302,5 +338,5 @@ CREATE TABLE SCHD_SNAP_LOG (
ARCHIVE_KEY VARCHAR(200), -- 原文归档引用(证据层,尚未交付) ARCHIVE_KEY VARCHAR(200), -- 原文归档引用(证据层,尚未交付)
CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL 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); CREATE INDEX idx_snaplog_msg ON SCHD_SNAP_LOG (MSG_ID);
@@ -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;
@@ -1,3 +0,0 @@
-- HOL deadline 使用首次开始处理的稳定时刻,不再被每次重试更新的 UPDATED_AT 推后。
ALTER TABLE PROC_STATE ADD COLUMN PROCESSING_STARTED_AT TIMESTAMP(6) WITH TIME ZONE;
@@ -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;
@@ -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)(跳过当前可见存量)、<id> = 显式边界;
-- · 代码**不做默认选择**,也不会自动退化成 max;
-- · 升级实例(已有水位或已有处理记录)**拒绝重新播种**,重新切流必须是显式操作;
-- · SEEDED_AT 记录"已播种"这一事实;注意它**为 NULL 不等于"从未消费"**——
-- 已有库新增列后同样是 NULL,因此判据必须叠加"确实没有消费过"。
-- =====================================================================
ALTER TABLE INBOX_CURSOR ADD COLUMN SEEDED_AT TIMESTAMP(6) WITH TIME ZONE;
@@ -1,22 +0,0 @@
-- =====================================================================
-- V6PROC_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();
@@ -1,14 +0,0 @@
-- =====================================================================
-- V7PROC_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;
@@ -1,34 +0,0 @@
-- =====================================================================
-- V8REQ_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');
@@ -1,27 +0,0 @@
-- =====================================================================
-- V9KAFKA: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;
@@ -7,19 +7,21 @@ Flyway 配置**PG 路径使用 `classpath:db/migration`,两者互不混用
1. 现场 11.2 补丁级别、数据库字符集、DBA 权限清单拿到,且可提供可测试的目标库。 1. 现场 11.2 补丁级别、数据库字符集、DBA 权限清单拿到,且可提供可测试的目标库。
2. JDK 25 × ojdbc 驱动(具体版本)× Flyway Oracle 支持 × 连接池组合在目标库实测通过 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 的 upsertMERGE INTO)与绑定顺序适配完成并通过与 PG 同粒度的集成测试后, 3. 11g 的 upsertMERGE INTO)与绑定顺序适配完成并通过与 PG 同粒度的集成测试后,
才允许把仓储 SQL 切到 11g 方言。先前编译级交付的 `SqlDialect` 方言接缝已随 才允许把仓储 SQL 切到 11g 方言。先前编译级交付的 `SqlDialect` 方言接缝已随
V1 基线列名更替(fday/last_message_id → operation_day/last_msg_id)退役删除; 基线列名更替(fday/last_message_id → operation_day/last_msg_id)退役删除;
接入时必须按 V1 基线现有列名重建,禁止直接复用旧接缝。 接入时必须按基线现有列名重建,禁止直接复用旧接缝。
## 计划内容 ## 计划内容
- V1 起`V1__flight_state_baseline.sql` 的 11g 等价 DDL—— - 单基线`V1__flight_state_baseline.sql` 的 11g 等价 DDL——
PIPELINE_LOCK/PROC_STATE/MSG_EVENT/REQ_TRACK/BACKFILL_TODO/FLIGHT_SCHD + 9 张明细表/SCHD_SNAP_LOG 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 语义钉死。 `NUMBER`/序列替代 `BIGSERIAL``TIMESTAMP WITH TIME ZONE``VARCHAR2` BYTE/CHAR 语义钉死。
- `INSERT ... ON CONFLICT` 改 11g MERGEOPERATION_DAY 不可变条件,flight-state.md §2.1)。 - `INSERT ... ON CONFLICT` 改 11g MERGEOPERATION_DAY 不可变条件,`INV-12`)。
- 空串按 NULL 的语义回归:显式清空的 presence 信息不得被 11g 空串语义吞掉 - 空串按 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 各自独立推进,禁止复用版本号语义。 版本号与 PG location 各自独立推进,禁止复用版本号语义。
@@ -10,9 +10,10 @@ import org.junit.jupiter.api.Test
import java.sql.DriverManager import java.sql.DriverManager
/** /**
* 在真实 PostgreSQL 上跑一遍迁移确认结果符合预期 * 在真实 PostgreSQL 上跑一遍迁移确认结果符合预期单基线 `V1__flight_state_baseline.sql`
* V1 基线V2 信箱生命周期和 V3 稳定处理起点都能成功执行迁移记录显示成功该建的表和单行种子 * 执行成功且是唯一的迁移记录该建的表和单行种子PIPELINE_LOCKINBOX_CURSOR都在
* PIPELINE_LOCKINBOX_CURSOR都在回填相关字段进了 PROC_STATEBACKFILL_TODO 已下线 * 回填事实与收报水位都落在基线里`BACKFILL_TODO``idx_evt_flid``PROC_STATE` 的处理开始
* 时间列都不复存在这些正是原 V2V10 合并后的净结构
* *
* 没有可用的 PostgreSQL 时跳过不假装通过 * 没有可用的 PostgreSQL 时跳过不假装通过
*/ */
@@ -41,7 +42,7 @@ class FlywayMigrationTest {
DriverManager.getConnection(url, user, pass).use { conn -> DriverManager.getConnection(url, user, pass).use { conn ->
conn.createStatement().use { stmt -> conn.createStatement().use { stmt ->
// 迁移记录:V1 基线 + V2 信箱生命周期 + V3 稳定处理起点 + V4 回填闭环 + V5 切流播种 // 单基线:原 V1–V10 已合并为一条迁移,历史里只应留下 V1
stmt.executeQuery( stmt.executeQuery(
"SELECT version, script, success FROM flyway_schema_history ORDER BY installed_rank ASC", "SELECT version, script, success FROM flyway_schema_history ORDER BY installed_rank ASC",
).use { rs -> ).use { rs ->
@@ -49,29 +50,13 @@ class FlywayMigrationTest {
while (rs.next()) { while (rs.next()) {
records.add(Triple(rs.getString("version"), rs.getString("script"), rs.getBoolean("success"))) 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("1", records[0].first)
assertEquals("V1__flight_state_baseline.sql", records[0].second) 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 }) assertTrue(records.all { it.third })
} }
// 全表就绪(V2 后 backfill_todo 下线 // 全表就绪(回填事实并入 PROC_STATE,基线不再建 BACKFILL_TODO
stmt.executeQuery( stmt.executeQuery(
"SELECT table_name FROM information_schema.tables WHERE table_schema = 'public'", "SELECT table_name FROM information_schema.tables WHERE table_schema = 'public'",
).use { rs -> ).use { rs ->
@@ -110,20 +95,20 @@ class FlywayMigrationTest {
) )
} }
// V7PROC_STATE 不再保留处理开始时间(无判据消费 // 基线不含处理开始时间(原 V3 引入、V7 删除,合并后不得残留
stmt.executeQuery( stmt.executeQuery(
"SELECT count(*) FROM information_schema.columns WHERE table_name = 'proc_state' " + "SELECT count(*) FROM information_schema.columns WHERE table_name = 'proc_state' " +
"AND column_name = 'processing_started_at'", "AND column_name = 'processing_started_at'",
).use { rs -> ).use { rs ->
assertTrue(rs.next()) 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( stmt.executeQuery(
"SELECT indexdef FROM pg_indexes WHERE tablename = 'req_track' AND indexname = 'uq_req_open'", "SELECT indexdef FROM pg_indexes WHERE tablename = 'req_track' AND indexname = 'uq_req_open'",
).use { rs -> ).use { rs ->
assertTrue(rs.next(), "V8 必须建立开放态唯一索引 uq_req_open") assertTrue(rs.next(), "基线必须建立开放态唯一索引 uq_req_open")
val indexDef = rs.getString(1) val indexDef = rs.getString(1)
assertTrue(indexDef.contains("UNIQUE"), "uq_req_open 必须是唯一索引:$indexDef") assertTrue(indexDef.contains("UNIQUE"), "uq_req_open 必须是唯一索引:$indexDef")
assertTrue(indexDef.contains("WHERE"), "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())", "VALUES ('RQFD-NONE', DATE '2026-09-12', 'RMS', 'PENDING', now())",
) )
// V9KAFKA:schd 每个 FLID 单行(部分唯一索引),KAFKA:msg 仍可多行 // 基线KAFKA:schd 每个 FLID 单行(部分唯一索引),KAFKA:msg 仍可多行
stmt.executeQuery( stmt.executeQuery(
"SELECT indexdef FROM pg_indexes WHERE tablename = 'msg_event' AND indexname = 'uq_schd_event'", "SELECT indexdef FROM pg_indexes WHERE tablename = 'msg_event' AND indexname = 'uq_schd_event'",
).use { rs -> ).use { rs ->
assertTrue(rs.next(), "V9 必须建立 schd 单行唯一索引 uq_schd_event") assertTrue(rs.next(), "基线必须建立 schd 单行唯一索引 uq_schd_event")
val indexDef = rs.getString(1) val indexDef = rs.getString(1)
assertTrue(indexDef.contains("UNIQUE"), "uq_schd_event 必须是唯一索引:$indexDef") assertTrue(indexDef.contains("UNIQUE"), "uq_schd_event 必须是唯一索引:$indexDef")
assertTrue(indexDef.contains("WHERE"), "uq_schd_event 必须是仅约束 schd 的部分索引:$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'", "SELECT count(*) FROM pg_indexes WHERE tablename = 'msg_event' AND indexname = 'idx_evt_flid'",
).use { rs -> ).use { rs ->
assertTrue(rs.next()) 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( stmt.executeQuery(
"SELECT is_nullable FROM information_schema.columns " + "SELECT is_nullable FROM information_schema.columns " +
"WHERE table_name = 'proc_state' AND column_name = 'enqueued_at'", "WHERE table_name = 'proc_state' AND column_name = 'enqueued_at'",
).use { rs -> ).use { rs ->
assertTrue(rs.next(), "V6 必须已加上 ENQUEUED_AT 列") assertTrue(rs.next(), "基线必须已加上 ENQUEUED_AT 列")
assertEquals("NO", rs.getString("is_nullable"), "ENQUEUED_AT 必须非空") assertEquals("NO", rs.getString("is_nullable"), "ENQUEUED_AT 必须非空")
} }