diff --git a/docs/decision-flight-state.md b/docs/decision-flight-state.md index 785bc48..36daa3d 100644 --- a/docs/decision-flight-state.md +++ b/docs/decision-flight-state.md @@ -95,3 +95,25 @@ A 是最小代价回退位;B 不建议。 - **FS6** U29/U09 不变量测试门禁:崩溃幂等无自增、CAS 防并发、ADFT 存活保障、UTC 方言测试、ES 历史清场 5 场景全绿。 - **FS7** 影子对拍 Diff 工具:`FlightStoreDiffTool` 递归 AST 字段归一化比较与已知合法偏离识别。 - **FS8** 文档回改:架构、设计、用户故事与项目规范同步收敛。 + +## 8. 修订:11g 现场约束下废除 JSON 存储(本节追加,不改写上文定案历史) + +> **触发**:现场环境只提供 Oracle 11g(无任何 JSON 能力——无 `JSON_VALUE`、无 `IS JSON` 约束、无 JSON 类型, +> JSON 支持自 12.1.0.2 才引入)。Q2 定案的「FLTR JSONB 全量」存储形态在该约束下不成立: +> 整文档存储丧失库端校验、字段级索引与直接 SQL 可查性,违背「运营航班可被 DBA/运营直接使用」的初衷。 + +**修订结论**: + +| # | 问题 | 修订结论 | +|---|---|---| +| R1 | 表形态 | 废除 `FLTR_JSON` JSONB 整文档,改为**一行一航班的宽表**:FLID 主键 + FDAY 代列 + SCHD.FLTR 标量字段列(锚定 `unisysaodbsis.xsd` 契约与 legacy 派生字段 ABDG/LPSDT/ABN)+ 1:N 明细集合序列化文本列。字段即列,天然可索引、可直查。 | +| R2 | 集合字段 | 登机口/柜台/转盘/桥/延误等 1:N 明细(单航班可达 99 条)暂存序列化文本列(`*_TXT`,11g 移植为 CLOB);阶段 2/3 Handler 钉死语义后按需升独立子表。 | +| R3 | 值语义 | 全部字段保持 legacy 字符串原样(不做库端类型转换),保证影子对拍逐字段保真;范围查询需要的类型化列(如 SODT→DATE)按查询需求逐列后置提升。 | +| R4 | SCHD_GEN | `FLIDS_JSON` 同步废除,展开为 `SCHD_GEN_FLID(FDAY, FLID)` 行;版本 CAS 语义不变。 | +| R5 | 值机视图 | Handler 在线视图契约由 `flid → FLTR_JSON` 改为 `flid → 字段集`(与 legacy `hgetAllFlightInfo` hash 同构);增量写语义=字段级合并(与 legacy hmset 一致),快照=整体替换。 | +| R6 | 对拍口径 | FS7 Diff 工具改为 PG 宽表列值 vs legacy FLTR JSON 的逐字段归一化比对(数值精度/空值等价规则保留),比原 JSONB AST 比对更精确到字段。 | + +**未尽事项(11g 方言移植不在本修订范围,另行决策)**:存储模型已 11g 兼容,但仓储层仍存 PG 方言 +(`ON CONFLICT` upsert 需改 `MERGE`)、I5 advisory lock 的 11g 替代(`DBMS_LOCK`)、Flyway/驱动 +对 11.2 的支持矩阵。若「自有 PG → 现场 11g」成为确定部署形态,需按 ACM2-28 体例开独立决策记录 +重开「自有库」平台定案,评估点在方言与并发原语,不在本修订已解决的存储形态。 diff --git a/docs/design.md b/docs/design.md index 9b81100..6081604 100644 --- a/docs/design.md +++ b/docs/design.md @@ -19,8 +19,9 @@ | `PUMP_JOB` | 持久化维护作业 | 状态为 `QUEUED / RUNNING / DONE / FAILED`;不与业务消息共用排序序号。 | | `REQ_TRACK` | 上游请求及应答关联 | 保存请求类型、参数、出站信箱 ID、发送和完成时间;同类只允许一个开放请求。 | | `REF_MASTER` | 静态参考数据 | `(RTYPE, RKEY)` 唯一,`SOURCE` 记录数据来源。 | -| `FLIGHT_SCHD` | 运营航班当前权威状态(SCHD 快照 + FLOP/ADFT 增量合并) | `FLID` 主键,`FDAY` 所属日代(可空),`FLTR_JSON` 全量,`TIMESTAMPTZ` 时间戳。 | -| `SCHD_GEN` | 日计划代版本与当前代有效航班集合(差删依据) | `FDAY` 主键,`VERSION` 版本号,`FLIDS_JSON` 航班集合,`TIMESTAMPTZ`。 | +| `FLIGHT_SCHD` | 运营航班当前权威状态(SCHD 快照 + FLOP/ADFT 增量合并),一行一航班宽表 | `FLID` 主键,`FDAY` 所属日代(可空),SCHD.FLTR 标量字段列 + 1:N 明细集合序列化文本列(11g 修订,废除 `FLTR_JSON` JSONB 整文档),`TIMESTAMPTZ` 时间戳。 | +| `SCHD_GEN` | 日计划代版本(差删依据与版本 CAS) | `FDAY` 主键,`VERSION` 版本号,`TIMESTAMPTZ`。 | +| `SCHD_GEN_FLID` | 当前代有效航班 FLID 集合(原 `FLIDS_JSON` 展开为行) | `(FDAY, FLID)` 主键,`FDAY` 外键级联删除。 | | `PROC_STATE_HST` | 终态处理记录的归档目标 | 属于目标设计,当前迁移尚未建表;不得改写为共享库历史表。 | 字段与索引定义以 `src/main/resources/db/migration/` 为准(含 `V1.1.0__flight_schd.sql`)。报文原文仍从共享信箱读取,因此必须协调原文保留期,不能在消息尚需处理或重放时提前清理。 diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/SchdAggregation.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/SchdAggregation.kt index 28fb62c..32f17cd 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/SchdAggregation.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/SchdAggregation.kt @@ -18,5 +18,5 @@ object SchdAggregation { fun latestPerFlightOfPush(pushes: List): List = pushes .groupBy { it.flid } - .map { (_, ps) -> ps.maxBy { it.eventSeq }.fltrJson } + .map { (_, ps) -> ps.maxBy { it.eventSeq }.payloadJson } } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/Decision.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/Decision.kt index ff68c86..2681ab8 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/Decision.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/Decision.kt @@ -12,10 +12,14 @@ data class Decision( val refUpserts: List = emptyList(), // → 静态主数据(独立 PG reference 库,ACM2-11) ) -/** 航班状态变更(阶段 A 落自有 PG FLIGHT_SCHD,与事件同事务原子提交;ACM2-28 定案)。 */ +/** + * 航班状态变更(阶段 A 落自有库 FLIGHT_SCHD 宽表,与事件同事务原子提交;ACM2-28 定案)。 + * fields = 本报文变更的字段集(field → value,与 legacy flightInfo hash 同构), + * 仓储按「字段级合并」落库(仅新增/覆盖,不删除缺失字段,与 legacy hmset 同语义)。 + */ data class FlightChange( val flid: String, - val payloadJson: String, + val fields: Map, val maid: String? = null, ) @@ -23,7 +27,7 @@ data class NotifyPayload(val payloadJson: String) data class SchdPush( val flid: String, - val fltrJson: String, + val payloadJson: String, // KAFKA_SCHD 出站载荷(报文线格式,与库内展开存储无关) val eventSeq: Long = 0, // 由事务插入时赋 EVENT_ID 语义序,聚合取 max ) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/FlightFieldsJson.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/FlightFieldsJson.kt new file mode 100644 index 0000000..875f96a --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/FlightFieldsJson.kt @@ -0,0 +1,19 @@ +package com.gzzn.omms.msgexchange.infra.persistence + +import com.fasterxml.jackson.databind.ObjectMapper +import com.fasterxml.jackson.databind.node.ObjectNode + +/** + * 航班字段集(FLIGHT_SCHD 宽表字段)→ KAFKA_SCHD 事件载荷序列化。 + * 仅用于 MSG_EVENT 出站载荷(报文线格式,legacy 下游按 JSON 消费),与库内存储无关; + * 键按字典序输出保证同状态产出字节级稳定载荷(重放/幂等判据不因遍历序漂移)。 + */ +object FlightFieldsJson { + private val mapper = ObjectMapper() + + fun toJson(fields: Map): String { + val node = mapper.createObjectNode() + fields.toSortedMap().forEach { (k, v) -> node.put(k, v) } + return node.toString() + } +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt index 9e1eecb..2063ac7 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt @@ -9,13 +9,15 @@ import com.gzzn.omms.msgexchange.domain.RefUpsert import java.time.Instant /** - * ACM2-12 仓储接口(接口驱动,主泵/调度循环可单测;Micronaut Data JDBC 实装属 U05 批次)。 + * ACM2-28 仓储接口(接口驱动,主泵/调度循环可单测;Micronaut Data JDBC 实装属 U05 批次)。 * 存储边界:自有 PostgreSQL(datasources.default)= 本文件除 CminmsgInboxRepository 外 * 的全部接口(消息管道 PROC_STATE/MSG_EVENT、PUMP_JOB、REQ_TRACK、21 类 REF_MASTER); * 共享 MySQL 信箱(CMINMSGS / COUTMSGS)经信箱封装访问,仅 DML、不建表; * 主路径=上游外部写 CMINMSGS → 本系统 JDBC 轮询读;compat=insertRaw HTTP 写; - * 自有 PG = 消息管道 + 运营航班 FLIGHT_SCHD/SCHD_GEN + 静态数据;Redis 已退出阶段 A 权威与写路径。 + * 自有库 = 消息管道 + 运营航班 FLIGHT_SCHD/SCHD_GEN + 静态数据;Redis 已退出阶段 A 权威与写路径。 */ +/** 单航班字段集(field → value,与 legacy flightInfo hash 同构)。 */ +typealias FlightFields = Map interface ProcStateRepository { fun insert(cminmsgsId: Long, state: ProcStatus = ProcStatus.PENDING) @@ -68,20 +70,22 @@ interface MsgEventRepository { } /** - * 阶段 A 运营航班权威与日计划代(ACM2-28 采纳选项 C 定案): - * - 表 FLIGHT_SCHD:当前运营航班全量权威态(SCHD 快照 + FLOP/ADFT 增量合并),落自有 PostgreSQL; - * - 表 SCHD_GEN:各日代版本与当前代有效航班全量集合(差删依据),由 Redis 回归自有 PG; - * - Redis 退出动态权威与全部写路径; - * - 事务 2 与快照发布全在自有 PG 内以单事务原子提交; - * - 增量更新(FLOP/ADFT):新插 FDAY=NULL,已有行通过 ON CONFLICT 保留原 FDAY; - * - 按代差删域化:DELETE FROM FLIGHT_SCHD WHERE FDAY = :day AND FLID = ANY(:diffSet) + * 阶段 A 运营航班权威与日计划代(ACM2-28 定案 + 11g 修订:废除 FLTR_JSON JSONB 整文档存储): + * - 表 FLIGHT_SCHD:一行一航班的运营航班宽表(FLID 主键 + FDAY 可空所属代 + SCHD.FLTR 标量字段列 + * + 1:N 明细集合序列化文本列),字段即列、天然可索引可直查,PG/Oracle 11g 方言一致; + * - 表 SCHD_GEN:各日代版本(SQL CAS);表 SCHD_GEN_FLID:当前代有效航班 FLID 集合(差删依据); + * - Redis 退出动态权威与全部写路径;事务 2 与快照发布全在自有库内单事务原子提交; + * - 增量更新(FLOP/ADFT):新插 FDAY=NULL,已有行保留原 FDAY,按字段列更新 + * (与 legacy flightInfo hash「仅新增/覆盖、不删除缺失字段」同语义); + * - 快照(DNLD):声明/更新 FDAY 归属,航班字段整体替换; + * - 按代差删域化:DELETE FROM FLIGHT_SCHD WHERE FDAY = :day AND FLID IN (:diffSet) * 仅删除仍属旧代的行,ADFT(FDAY=NULL)与已迁移至新代的同 FLID 行天然存活。 */ interface FlightSchdRepository { data class FlightRecord( val flid: String, val fday: String?, - val fltrJson: String, + val fields: FlightFields, val createdAt: Instant, val updatedAt: Instant, ) @@ -93,26 +97,26 @@ interface FlightSchdRepository { val updatedAt: Instant = Instant.now(), ) - /** 快照全量写入(DNLD):强行声明/更新 FDAY 归属,批处理写入。 */ - fun upsertSnapshotBatch(day: String, flights: List>, now: Instant = Instant.now()) + /** 快照全量写入(DNLD):强行声明/更新 FDAY 归属,批处理写入,航班字段集整体替换。 */ + fun upsertSnapshotBatch(day: String, flights: List>, now: Instant = Instant.now()) - /** 增量更新(FLOP/ADFT):新插 FDAY=NULL,已有行保留原 FDAY。 */ + /** 增量更新(FLOP/ADFT):新插 FDAY=NULL,已有行保留原 FDAY;字段级合并。 */ fun upsertIncremental(changes: List, now: Instant = Instant.now()) /** 按代差删域化:仅删除 FDAY = day 且在 delFlids 中的记录(ADFT 与跨代已迁移行受保护)。 */ fun deleteDiffByDay(day: String, delFlids: Collection): Int - /** 点查单航班 FLTR_JSON。 */ - fun findByFlid(flid: String): String? + /** 点查单航班字段集;无字段行(含航班不存在)返回 null。 */ + fun findByFlid(flid: String): FlightFields? - /** 点查多航班 FLTR_JSON。 */ - fun findByFlids(flids: Collection): Map + /** 点查多航班字段集:FLID → 字段集,仅含实际存在的航班。 */ + fun findByFlids(flids: Collection): Map /** 按计划日查询当前有效航班。 */ - fun findByDay(day: String): List> + fun findByDay(day: String): List> /** 全量查询(供影子对拍 / 一致性对账)。 */ - fun findAll(): Map + fun findAll(): Map /** * 历史清场删除:仅删除已确认归档至 ES 的 FLID 集合;空集合不执行;分批参数化删除。 @@ -209,9 +213,9 @@ interface PumpJobRepository { @Deprecated("Replaced by FlightSchdRepository in ACM2-28", ReplaceWith("FlightSchdRepository")) interface FlightStateRepository { /** 阶段 B 权威;replaceDay = 单事务删差集+写新代+版本提升。 */ - fun replaceDay(day: String, flights: List>) + fun replaceDay(day: String, flights: List>) - fun findByDay(day: String): List> + fun findByDay(day: String): List> } /** diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt index 424ed4b..cf14ed1 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt @@ -6,6 +6,7 @@ import com.gzzn.omms.msgexchange.domain.MsgEvent import com.gzzn.omms.msgexchange.domain.ProcState import com.gzzn.omms.msgexchange.domain.ProcStatus import com.gzzn.omms.msgexchange.domain.RefUpsert +import com.gzzn.omms.msgexchange.infra.persistence.FlightFields import com.gzzn.omms.msgexchange.infra.persistence.FlightSchdRepository import com.gzzn.omms.msgexchange.infra.persistence.FlightStateRepository import com.gzzn.omms.msgexchange.infra.persistence.MsgEventRepository @@ -304,45 +305,73 @@ class JdbcPipelineTransactionManager( class JdbcFlightSchdRepository( private val ds: DataSource, ) : FlightSchdRepository { - private val flidsTypeRef = object : com.fasterxml.jackson.core.type.TypeReference>() {} - private val mapper = com.fasterxml.jackson.databind.ObjectMapper() - private fun serializeFlids(flids: Collection): String = mapper.writeValueAsString(flids) + companion object { + /** SCHD.FLTR 标量列 + legacy 派生列(与 V1.1.0 迁移一致,全部可空 VARCHAR)。 */ + private val SCALAR_COLUMNS: List = listOf( + "ALCD", "ALSC", "FLNO", "MVIN", "SODT", "FLTY", "FLIN", "ACFT", "RENO", + "TAOP", "TAFL", "TAID", "TRML", "MAXP", "CSOP", "CSFT", "MAID", + "ESTT", "ACTT", "STND", "PHAG", "CNCL", "REMC", "BOTM", "LACL", "FINT", + "APPT", "EGSR", "EGST", "FHAG", "MHAG", "VIPP", "VIPR", "LBNO", "LBWT", + "PAXC", "EXSC", "EXSR", "FTSS", "PEDT", "NEAT", "PADT", "NAAT", + "ABDG", "LPSDT", "ABN", + ) - private fun parseFlids(json: String): Set { - if (json.isBlank()) return emptySet() - return try { - mapper.readValue(json, flidsTypeRef) - } catch (_: Exception) { - emptySet() - } + /** 1:N 明细集合键(库列名 = 键 + _TXT 后缀,序列化文本)。 */ + private val COLLECTION_KEYS: List = listOf( + "ROUT", "ERUT", "CHDT", "GTDT", "PSDT", "CKDT", "CLDT", "DELY", + "CHOT", "ABTM", "SRVT", "VIPF", "MAFL", "FDIV", "FRET", "FLAB", + ) + + /** 库列名全集(标量原样;集合列带 _TXT 后缀)。 */ + private val ALL_COLUMNS: List = SCALAR_COLUMNS + COLLECTION_KEYS.map { "${it}_TXT" } + + /** 字段键 → 库列名白名单;未知键拒绝写入(fail fast,经处理边界落 FAILED(INFRA))。 + * FLID 即主键列:写侧忽略该键(主键已承载),读侧合成返回,视图与 legacy hash 同构。 */ + private val FIELD_TO_COLUMN: Map = + SCALAR_COLUMNS.associateWith { it } + COLLECTION_KEYS.associateWith { "${it}_TXT" } + + private fun columnToFieldKey(column: String): String = + if (column.endsWith("_TXT")) column.removeSuffix("_TXT") else column } private fun toSqlDate(day: String): java.sql.Date = java.sql.Date.valueOf(day.trim().take(10)) - override fun upsertSnapshotBatch(day: String, flights: List>, now: Instant) { + private fun requireColumns(fields: FlightFields): List> = + fields.mapNotNull { (key, value) -> + if (key == "FLID") return@mapNotNull null // 主键列:写侧忽略 + val column = FIELD_TO_COLUMN[key] + ?: throw IllegalArgumentException("unknown flight field: $key") + column to value + } + + override fun upsertSnapshotBatch(day: String, flights: List>, now: Instant) { if (flights.isEmpty()) return - val sql = """ - INSERT INTO flight_schd (flid, fday, fltr_json, created_at, updated_at) - VALUES (?, ?, ?::jsonb, ?, ?) - ON CONFLICT (flid) DO UPDATE SET - fday = EXCLUDED.fday, - fltr_json = EXCLUDED.fltr_json, - updated_at = EXCLUDED.updated_at - """.trimIndent() + // 快照 = 航班字段整体替换:缺省字段置 NULL,批量 SQL 形状对全部行一致 + val sql = buildString { + append("INSERT INTO flight_schd (flid, fday, ") + append(ALL_COLUMNS.joinToString(", ")) + append(", created_at, updated_at) VALUES (?, ?, ") + append(ALL_COLUMNS.joinToString(", ") { "?" }) + append(", ?, ?) ON CONFLICT (flid) DO UPDATE SET ") + append((listOf("fday") + ALL_COLUMNS + listOf("updated_at")).joinToString(", ") { "$it = EXCLUDED.$it" }) + } val sqlDate = toSqlDate(day) val sqlTimestamp = now.toSqlTimestamp() val conn = ds.obtainConnection() try { conn.prepareStatement(sql).use { ps -> var count = 0 - for ((flid, json) in flights) { - ps.setString(1, flid) - ps.setDate(2, sqlDate) - ps.setString(3, json) - ps.setTimestamp(4, sqlTimestamp) - ps.setTimestamp(5, sqlTimestamp) + for ((flid, fields) in flights) { + var i = 1 + ps.setString(i++, flid) + ps.setDate(i++, sqlDate) + for (column in ALL_COLUMNS) { + ps.setString(i++, fields[columnToFieldKey(column)]) + } + ps.setTimestamp(i++, sqlTimestamp) + ps.setTimestamp(i++, sqlTimestamp) ps.addBatch() count++ if (count % 200 == 0) { @@ -360,23 +389,19 @@ class JdbcFlightSchdRepository( override fun upsertIncremental(changes: List, now: Instant) { if (changes.isEmpty()) return - val sql = """ - INSERT INTO flight_schd (flid, fday, fltr_json, created_at, updated_at) - VALUES (?, NULL, ?::jsonb, ?, ?) - ON CONFLICT (flid) DO UPDATE SET - fltr_json = EXCLUDED.fltr_json, - updated_at = EXCLUDED.updated_at - """.trimIndent() val sqlTimestamp = now.toSqlTimestamp() val conn = ds.obtainConnection() try { - conn.prepareStatement(sql).use { ps -> + // 父行:新插 FDAY=NULL;已有行保留原 FDAY,仅推进 updated_at + conn.prepareStatement( + "INSERT INTO flight_schd (flid, created_at, updated_at) VALUES (?, ?, ?) " + + "ON CONFLICT (flid) DO UPDATE SET updated_at = EXCLUDED.updated_at", + ).use { ps -> var count = 0 for (c in changes) { ps.setString(1, c.flid) - ps.setString(2, c.payloadJson) + ps.setTimestamp(2, sqlTimestamp) ps.setTimestamp(3, sqlTimestamp) - ps.setTimestamp(4, sqlTimestamp) ps.addBatch() count++ if (count % 200 == 0) { @@ -387,6 +412,20 @@ class JdbcFlightSchdRepository( ps.executeBatch() } } + // 字段列更新(legacy flightInfo hash 同语义:仅新增/覆盖,不删除缺失字段) + for (c in changes) { + val columns = requireColumns(c.fields) + if (columns.isEmpty()) continue + val sql = "UPDATE flight_schd SET " + + columns.joinToString(", ") { "${it.first} = ?" } + + ", updated_at = ? WHERE flid = ?" + ds.update(sql) { ps -> + var i = 1 + columns.forEach { (_, value) -> ps.setString(i++, value) } + ps.setTimestamp(i++, sqlTimestamp) + ps.setString(i, c.flid) + } + } } finally { conn.releaseIfNotInTransaction() } @@ -407,38 +446,52 @@ class JdbcFlightSchdRepository( return totalDeleted } - override fun findByFlid(flid: String): String? = - ds.queryOne( - "SELECT fltr_json FROM flight_schd WHERE flid = ?", - { ps -> ps.setString(1, flid) }, - ) { rs -> rs.getString("fltr_json") } + private fun mapFlightRow(rs: java.sql.ResultSet): Pair { + val flid = rs.getString("flid") + val fields = linkedMapOf() + fields["FLID"] = flid // 视图与 legacy flightInfo hash 同构:FLID 字段常在 + for (column in ALL_COLUMNS) { + val value = rs.getString(column) ?: continue + fields[columnToFieldKey(column)] = value + } + return flid to fields + } - override fun findByFlids(flids: Collection): Map { + override fun findByFlid(flid: String): FlightFields? { + val row = ds.queryOne( + "SELECT flid, ${ALL_COLUMNS.joinToString(", ")} FROM flight_schd WHERE flid = ?", + { ps -> ps.setString(1, flid) }, + ::mapFlightRow, + ) + return row?.second + } + + override fun findByFlids(flids: Collection): Map { if (flids.isEmpty()) return emptyMap() - val result = mutableMapOf() + val result = linkedMapOf() + val sql = "SELECT flid, ${ALL_COLUMNS.joinToString(", ")} FROM flight_schd WHERE flid IN (" for (chunk in flids.chunked(200)) { val placeholders = chunk.joinToString(",") { "?" } - val sql = "SELECT flid, fltr_json FROM flight_schd WHERE flid IN ($placeholders)" - val pairs = ds.query( - sql, + ds.query( + sql + placeholders + ")", { ps -> chunk.forEachIndexed { i, flid -> ps.setString(i + 1, flid) } }, - ) { rs -> rs.getString("flid") to rs.getString("fltr_json") } - result.putAll(pairs) + ) { rs -> mapFlightRow(rs) }.forEach { (flid, fields) -> result[flid] = fields } } return result } - override fun findByDay(day: String): List> = + override fun findByDay(day: String): List> = ds.query( - "SELECT flid, fltr_json FROM flight_schd WHERE fday = ? ORDER BY flid ASC", + "SELECT flid, ${ALL_COLUMNS.joinToString(", ")} FROM flight_schd WHERE fday = ? ORDER BY flid ASC", { ps -> ps.setDate(1, toSqlDate(day)) }, - ) { rs -> rs.getString("flid") to rs.getString("fltr_json") } + ::mapFlightRow, + ) - override fun findAll(): Map = + override fun findAll(): Map = ds.query( - "SELECT flid, fltr_json FROM flight_schd ORDER BY flid ASC", + "SELECT flid, ${ALL_COLUMNS.joinToString(", ")} FROM flight_schd ORDER BY flid ASC", {}, - ) { rs -> rs.getString("flid") to rs.getString("fltr_json") }.toMap() + ) { rs -> mapFlightRow(rs) }.toMap() override fun deleteByFlids(flids: Set): Int { if (flids.isEmpty()) return 0 @@ -453,53 +506,90 @@ class JdbcFlightSchdRepository( return totalDeleted } - override fun getGen(day: String): FlightSchdRepository.GenMeta? = - ds.queryOne( - "SELECT fday, version, flids_json, updated_at FROM schd_gen WHERE fday = ?", + override fun getGen(day: String): FlightSchdRepository.GenMeta? { + val meta = ds.queryOne( + "SELECT version, updated_at FROM schd_gen WHERE fday = ?", { ps -> ps.setDate(1, toSqlDate(day)) }, - ) { rs -> - FlightSchdRepository.GenMeta( - fday = rs.getDate("fday").toString(), - version = rs.getLong("version"), - flids = parseFlids(rs.getString("flids_json")), - updatedAt = rs.getInstant("updated_at") ?: Instant.now(), - ) + ) { rs -> rs.getLong("version") to (rs.getInstant("updated_at") ?: Instant.now()) } + ?: return null + val flids = ds.query( + "SELECT flid FROM schd_gen_flid WHERE fday = ? ORDER BY flid ASC", + { ps -> ps.setDate(1, toSqlDate(day)) }, + ) { rs -> rs.getString("flid") }.toSet() + return FlightSchdRepository.GenMeta( + fday = day, + version = meta.first, + flids = flids, + updatedAt = meta.second, + ) + } + + /** 当前代 FLID 集合整体替换(与 SCHD_GEN 版本推进同事务)。 */ + private fun replaceGenFlids(day: String, flids: Set) { + ds.update("DELETE FROM schd_gen_flid WHERE fday = ?") { ps -> + ps.setDate(1, toSqlDate(day)) } + if (flids.isEmpty()) return + val conn = ds.obtainConnection() + try { + conn.prepareStatement("INSERT INTO schd_gen_flid (fday, flid) VALUES (?, ?)").use { ps -> + var count = 0 + for (flid in flids.sorted()) { + ps.setDate(1, toSqlDate(day)) + ps.setString(2, flid) + ps.addBatch() + count++ + if (count % 500 == 0) { + ps.executeBatch() + } + } + if (count % 500 != 0) { + ps.executeBatch() + } + } + } finally { + conn.releaseIfNotInTransaction() + } + } override fun putGenIfVersion(day: String, expected: Long, newGen: FlightSchdRepository.GenMeta, now: Instant): Boolean { - val flidsJson = serializeFlids(newGen.flids) val sqlDate = toSqlDate(day) val sqlTimestamp = now.toSqlTimestamp() if (expected == 0L) { val inserted = ds.update( """ - INSERT INTO schd_gen (fday, version, flids_json, updated_at) - VALUES (?, ?, ?::jsonb, ?) + INSERT INTO schd_gen (fday, version, updated_at) + VALUES (?, ?, ?) ON CONFLICT (fday) DO NOTHING """.trimIndent(), ) { ps -> ps.setDate(1, sqlDate) ps.setLong(2, newGen.version) - ps.setString(3, flidsJson) - ps.setTimestamp(4, sqlTimestamp) + ps.setTimestamp(3, sqlTimestamp) + } + if (inserted == 1) { + replaceGenFlids(day, newGen.flids) + return true } - if (inserted == 1) return true } val updated = ds.update( """ UPDATE schd_gen - SET version = ?, flids_json = ?::jsonb, updated_at = ? + SET version = ?, updated_at = ? WHERE fday = ? AND version = ? """.trimIndent(), ) { ps -> ps.setLong(1, newGen.version) - ps.setString(2, flidsJson) - ps.setTimestamp(3, sqlTimestamp) - ps.setDate(4, sqlDate) - ps.setLong(5, expected) + ps.setTimestamp(2, sqlTimestamp) + ps.setDate(3, sqlDate) + ps.setLong(4, expected) } - return updated == 1 + if (updated == 1) { + replaceGenFlids(day, newGen.flids) + return true + } + return false } override fun deleteGenBefore(cutoffDay: String): Int = @@ -646,9 +736,9 @@ class JdbcReqTrackRepository( class JdbcFlightStateRepository( private val flightSchd: FlightSchdRepository, ) : FlightStateRepository { - override fun replaceDay(day: String, flights: List>) = + override fun replaceDay(day: String, flights: List>) = flightSchd.upsertSnapshotBatch(day, flights) - override fun findByDay(day: String): List> = + override fun findByDay(day: String): List> = flightSchd.findByDay(day) } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt index 023ee03..2dfde4c 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt @@ -7,6 +7,7 @@ import com.gzzn.omms.msgexchange.domain.ProcState import com.gzzn.omms.msgexchange.domain.ProcStatus import com.gzzn.omms.msgexchange.domain.RefUpsert import com.gzzn.omms.msgexchange.infra.persistence.CminmsgInboxRepository +import com.gzzn.omms.msgexchange.infra.persistence.FlightFields import com.gzzn.omms.msgexchange.infra.persistence.FlightSchdRepository import com.gzzn.omms.msgexchange.infra.persistence.FlightStateRepository import com.gzzn.omms.msgexchange.infra.persistence.MsgEventRepository @@ -201,12 +202,12 @@ class StubFlightSchd : FlightSchdRepository { data class Record( val flid: String, val fday: String?, - val fltrJson: String, + val fields: FlightFields, val createdAt: Instant, val updatedAt: Instant, ) - private val records = mutableMapOf() + private val records = linkedMapOf() private val gens = mutableMapOf() fun clear() { @@ -214,13 +215,13 @@ class StubFlightSchd : FlightSchdRepository { gens.clear() } - override fun upsertSnapshotBatch(day: String, flights: List>, now: Instant) { - for ((flid, json) in flights) { + override fun upsertSnapshotBatch(day: String, flights: List>, now: Instant) { + for ((flid, fields) in flights) { val existing = records[flid] records[flid] = Record( flid = flid, fday = day, - fltrJson = json, + fields = fields, createdAt = existing?.createdAt ?: now, updatedAt = now, ) @@ -233,7 +234,7 @@ class StubFlightSchd : FlightSchdRepository { records[c.flid] = Record( flid = c.flid, fday = existing?.fday, // preserve existing fday; null if new - fltrJson = c.payloadJson, + fields = (existing?.fields ?: emptyMap()) + c.fields, // 字段级合并 createdAt = existing?.createdAt ?: now, updatedAt = now, ) @@ -248,18 +249,18 @@ class StubFlightSchd : FlightSchdRepository { return toRemove.size } - override fun findByFlid(flid: String): String? = records[flid]?.fltrJson + override fun findByFlid(flid: String): FlightFields? = records[flid]?.fields - override fun findByFlids(flids: Collection): Map = - flids.mapNotNull { flid -> records[flid]?.let { flid to it.fltrJson } }.toMap() + override fun findByFlids(flids: Collection): Map = + flids.mapNotNull { flid -> records[flid]?.let { flid to it.fields } }.toMap() - override fun findByDay(day: String): List> = + override fun findByDay(day: String): List> = records.values.filter { it.fday == day } .sortedBy { it.flid } - .map { it.flid to it.fltrJson } + .map { it.flid to it.fields } - override fun findAll(): Map = - records.mapValues { it.value.fltrJson } + override fun findAll(): Map = + records.mapValues { it.value.fields } override fun deleteByFlids(flids: Set): Int { if (flids.isEmpty()) return 0 @@ -355,9 +356,9 @@ class StubReqTrack : ReqTrackRepository { class StubFlightState( private val flightSchd: FlightSchdRepository, ) : FlightStateRepository { - override fun replaceDay(day: String, flights: List>) = + override fun replaceDay(day: String, flights: List>) = flightSchd.upsertSnapshotBatch(day, flights) - override fun findByDay(day: String): List> = + override fun findByDay(day: String): List> = flightSchd.findByDay(day) } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/JobExecutor.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/JobExecutor.kt index 84e1fa6..c0ec732 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/JobExecutor.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/JobExecutor.kt @@ -1,5 +1,7 @@ package com.gzzn.omms.msgexchange.jobs +import com.gzzn.omms.msgexchange.infra.persistence.FlightFields +import com.gzzn.omms.msgexchange.infra.persistence.FlightFieldsJson import com.gzzn.omms.msgexchange.infra.persistence.FlightSchdRepository import com.gzzn.omms.msgexchange.infra.persistence.MsgEventRepository import com.gzzn.omms.msgexchange.infra.persistence.PumpJobRepository @@ -34,8 +36,8 @@ class HistorySweepJob( private val flightSchd: FlightSchdRepository, ) { companion object { - var historyPicker: ((Map) -> Map)? = null - var esArchiver: ((Map) -> Set)? = null + var historyPicker: ((Map) -> Map)? = null + var esArchiver: ((Map) -> Set)? = null var cutoffProvider: (() -> String?)? = null } @@ -54,7 +56,7 @@ class HistorySweepJob( /** U10/T07(修订):接入 ES success 集之前的门禁——占位默认返回空集; * 现役五条判史规则(SODT 3 天 / CNCL 1 小时 / 备降 / 离港 / 到港)golden 通过后才允许接线。 */ - private fun pickHistory(all: Map): Map = + private fun pickHistory(all: Map): Map = historyPicker?.invoke(all) ?: emptyMap() } @@ -77,9 +79,15 @@ class ProjectionRebuildJob( fun run() { for (day in activeDays()) { val flights = flightSchd.findByDay(day) - flights.forEach { (flid, payload) -> + flights.forEach { (flid, fields) -> msgEvents.insertSync( - listOf(MsgEvent(target = Targets.KAFKA_SCHD, partitionKey = flid, payloadJson = payload)), + listOf( + MsgEvent( + target = Targets.KAFKA_SCHD, + partitionKey = flid, + payloadJson = FlightFieldsJson.toJson(fields), + ), + ), ) } } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Handler.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Handler.kt index 81e7f55..7abe904 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Handler.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Handler.kt @@ -11,8 +11,11 @@ import com.gzzn.omms.msgexchange.domain.MsgKind interface Handler { val kind: MsgKind - /** 在线状态以只读视图传入(阶段 A 点查自有 PG FLIGHT_SCHD,由调用方装配,ACM2-28 定案)。 */ - fun decide(flightView: Map, msg: DecodedMessage): Decision + /** + * 在线状态以只读视图传入:FLID → 字段集(field → value,与 legacy flightInfo hash / + * hgetAllFlightInfo 同构;阶段 A 由调用方点查自有库 FLIGHT_SCHD 宽表装配,ACM2-28 定案)。 + */ + fun decide(flightView: Map>, msg: DecodedMessage): Decision } /** diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt index 5183b28..a3c9c96 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt @@ -219,7 +219,7 @@ class MessageProcessor( // upsert FLIGHT_SCHD(decision.flightChanges) + insert MSG_EVENT + PROC_STATE → SUCCEEDED val events = buildList { decision.msgNotifies.forEach { add(MsgEvent(target = Targets.KAFKA_MSG, payloadJson = it.payloadJson)) } - decision.schdPush.forEach { add(MsgEvent(target = Targets.KAFKA_SCHD, partitionKey = it.flid, payloadJson = it.fltrJson)) } + decision.schdPush.forEach { add(MsgEvent(target = Targets.KAFKA_SCHD, partitionKey = it.flid, payloadJson = it.payloadJson)) } } txManager.inTransaction { diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/SnapshotFlow.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/SnapshotFlow.kt index 2f8ce6a..cbe2e6d 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/SnapshotFlow.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/SnapshotFlow.kt @@ -8,6 +8,8 @@ import com.gzzn.omms.msgexchange.domain.MsgEvent import com.gzzn.omms.msgexchange.domain.MsgKind import com.gzzn.omms.msgexchange.domain.Targets import com.gzzn.omms.msgexchange.infra.persistence.CminmsgInboxRepository +import com.gzzn.omms.msgexchange.infra.persistence.FlightFields +import com.gzzn.omms.msgexchange.infra.persistence.FlightFieldsJson import com.gzzn.omms.msgexchange.infra.persistence.FlightSchdRepository import com.gzzn.omms.msgexchange.infra.persistence.MsgEventRepository import com.gzzn.omms.msgexchange.infra.persistence.PipelineTransactionManager @@ -42,7 +44,7 @@ class SnapshotFlow( return } val ok = staged as StageResult.Ok - val normalized = ok.flights // (flid, payloadJson) + val normalized = ok.flights // (flid, fields) val day = ok.day // 内存与超大包熔断防御(ACM2-28 评论 4 项 3) @@ -93,9 +95,9 @@ class SnapshotFlow( } } - // 构造并批量写入 schd 投递事件 - val events = normalized.map { (flid, json) -> - MsgEvent(target = Targets.KAFKA_SCHD, partitionKey = flid, payloadJson = json) + // 构造并批量写入 schd 投递事件(载荷 = 字段集序列化,KAFKA_SCHD 线格式) + val events = normalized.map { (flid, fields) -> + MsgEvent(target = Targets.KAFKA_SCHD, partitionKey = flid, payloadJson = FlightFieldsJson.toJson(fields)) } if (events.isNotEmpty()) { msgEvents.insertAll(events) @@ -116,7 +118,7 @@ class SnapshotFlow( /** staging 结果(骨架)。 */ sealed interface StageResult { - data class Ok(val day: String, val flights: List>) : StageResult + data class Ok(val day: String, val flights: List>) : StageResult data class Invalid(val reason: String) : StageResult companion object { diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/tools/FlightStoreDiffTool.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/tools/FlightStoreDiffTool.kt index 13bc8ef..118a53c 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/tools/FlightStoreDiffTool.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/tools/FlightStoreDiffTool.kt @@ -5,22 +5,25 @@ import com.fasterxml.jackson.databind.ObjectMapper import com.fasterxml.jackson.databind.node.ArrayNode import com.fasterxml.jackson.databind.node.ObjectNode import com.fasterxml.jackson.databind.node.ValueNode +import com.gzzn.omms.msgexchange.infra.persistence.FlightFields import java.math.BigDecimal /** - * ACM2-28 FS7:影子对拍跨存储 Diff 比较内核。 + * ACM2-28 FS7:影子对拍跨存储 Diff 工具。 * - * 核心定位: - * 本工具提供跨存储航班 AST 结构对比、字段规范化与已知合法偏离识别内核; - * 外部装配侧(从 PostgreSQL/Redis 读取数据并按 CMINMSGS_ID 水位对齐取数)属于影子期装配调用方职责。 - * 字段归一化规则(ACM2-28 评论 4 项 4): - * 1. 递归 AST/Map 结构对比,绝不做裸字符串比对(规避 PG jsonb 自动去重重排 key 的分叉); - * 2. 数值精度归一化(1.0 vs 1 视为等价); - * 3. 空值规范化(JSON null 与字段缺失视为等价,规避无语义偏离); + * 核心目标: + * 对齐 CMINMSGS_ID 水位后,比对 nextgen 自有库 FLIGHT_SCHD(11g 修订:宽表列值) + * vs legacy 现役 Redis (flightInfo)。 + * + * 字段归一化规则(ACM2-28 评论 4 项 4,随宽表存储形态修订): + * 1. PG 侧为宽表列值映射(字段名 → 字符串值),legacy 侧解析 FLTR JSON 后按字段逐一比对, + * 天然规避 JSONB/文本的 key 排序与序列化形态差异,绝不裸字符串比对; + * 2. 数值精度归一化(PG "100.0" vs legacy 100 视为等价;PG 值不可解析为数值则判不一致); + * 3. 空值规范化(legacy JSON null/字段缺失 与 PG 列 NULL 视为等价,规避无语义偏离); + * 4. 嵌套/集合值(legacy 为对象或数组):PG 侧序列化文本按 JSON 解析后递归比对。 * * 已知合法偏离声明(ACM2-28 方案六): - * ① FDAY 域化差删对跨代迁移行的保护(legacy 误删、nextgen 不删,属于修正性偏离); - * ② JSONB 与 Redis string 的 key 排序差异。 + * ① FDAY 域化差删对跨代迁移行的保护(legacy 误删、nextgen 不删,属于修正性偏离)。 */ class FlightStoreDiffTool( private val mapper: ObjectMapper = ObjectMapper(), @@ -85,12 +88,12 @@ class FlightStoreDiffTool( /** * 执行全量比对。 - * @param pgFlights nextgen FLIGHT_SCHD 行(flid -> fltr_json) + * @param pgFlights nextgen FLIGHT_SCHD 行(flid -> 宽表列值映射) * @param legacyFlights legacy Redis flightInfo 行(flid -> value string) * @param crossDayMigratedFlids 已知在多日之间迁移的航班 FLID 集合(用于识别合法偏差 ①) */ fun diff( - pgFlights: Map, + pgFlights: Map, legacyFlights: Map, crossDayMigratedFlids: Set = emptySet(), ): DiffReport { @@ -100,10 +103,10 @@ class FlightStoreDiffTool( val unexpected = mutableListOf() for (flid in allFlids) { - val pgRaw = pgFlights[flid] + val pgFields = pgFlights[flid] val legacyRaw = legacyFlights[flid] - if (pgRaw == null) { + if (pgFields == null) { unexpected += Deviation( flid = flid, kind = DeviationKind.MISSING_IN_PG, @@ -118,22 +121,22 @@ class FlightStoreDiffTool( known += Deviation( flid = flid, kind = DeviationKind.CROSS_DAY_PROTECTED, - pgValue = pgRaw, + pgValue = pgFields, detail = "Known Deviation 1: Cross-day migrated flight protected by FDAY domain-scoped delete in nextgen PG, erroneously deleted in legacy Redis", ) } else { unexpected += Deviation( flid = flid, kind = DeviationKind.UNEXPECTED_EXTRA_IN_PG, - pgValue = pgRaw, + pgValue = pgFields, detail = "Flight present in PG but missing in Redis without cross-day migration justification", ) } continue } - // 两侧均存在:进行 AST 递归规范化比较 - val mismatches = compareJsonTree(flid, pgRaw, legacyRaw) + // 两侧均存在:逐字段归一化比对(PG 列值 vs legacy JSON 解析树) + val mismatches = compareFlight(flid, pgFields, legacyRaw) if (mismatches.isEmpty()) { matched++ } else { @@ -151,26 +154,79 @@ class FlightStoreDiffTool( } /** - * 递归 AST 比较两段 JSON。 + * 单航班比对:PG 宽表列值映射 vs legacy FLTR JSON 文本。 + * 遍历两侧字段名并集,按归一化规则逐字段判定。 */ - internal fun compareJsonTree(flid: String, pgJson: String, legacyJson: String): List { - val pgNode = try { - mapper.readTree(pgJson) - } catch (e: Exception) { - return listOf(Deviation(flid, DeviationKind.FIELD_MISMATCH, detail = "PG JSON parse failed: ${e.message}")) - } - + internal fun compareFlight(flid: String, pgFields: FlightFields, legacyJson: String): List { val legacyNode = try { mapper.readTree(legacyJson) } catch (e: Exception) { return listOf(Deviation(flid, DeviationKind.FIELD_MISMATCH, detail = "Legacy JSON parse failed: ${e.message}")) } + if (!legacyNode.isObject) { + return listOf(Deviation(flid, DeviationKind.FIELD_MISMATCH, detail = "Legacy JSON is not an object")) + } + val legacyObj = legacyNode as ObjectNode val mismatches = mutableListOf() - compareNodes(flid, "", pgNode, legacyNode, mismatches) + val allKeys = (pgFields.keys + legacyObj.fieldNames().asSequence().toSet()).toSortedSet() + for (key in allKeys) { + val pgValue: String? = pgFields[key] + val legacyChild = legacyObj.get(key) + compareFieldNode(flid, key, pgValue, legacyChild, mismatches) + } return mismatches } + /** 单字段判定:PG 列值(字符串或缺失)vs legacy JSON 节点。 */ + private fun compareFieldNode(flid: String, path: String, pgValue: String?, legacyNode: JsonNode?, acc: MutableList) { + fun mismatch(pg: Any?, legacy: Any?, detail: String) { + acc += Deviation(flid, DeviationKind.FIELD_MISMATCH, path = path, pgValue = pg, legacyValue = legacy, detail = detail) + } + + val legacyAbsent = legacyNode == null || legacyNode.isNull + val pgAbsent = pgValue == null + + // 规范化:legacy JSON null/缺失 与 PG 列 NULL 等价 + if (legacyAbsent && pgAbsent) return + if (legacyAbsent != pgAbsent) { + mismatch(pgValue, legacyNode?.asText(), "One side is null/absent while other is present") + return + } + + val legacy = legacyNode!! + + // 数值归一化(PG 列值可解析为数值时按数值比较) + if (legacy.isNumber) { + val pgNumber = pgValue!!.trim().toBigDecimalOrNull() + val legacyNumber = BigDecimal(legacy.asText()).stripTrailingZeros() + if (pgNumber == null || pgNumber.stripTrailingZeros().compareTo(legacyNumber) != 0) { + mismatch(pgValue, legacy.asText(), "Numeric value mismatch: $pgValue != ${legacy.asText()}") + } + return + } + + // 嵌套/集合值:PG 序列化文本按 JSON 解析后递归比对 + if (legacy.isObject || legacy.isArray) { + val pgNode = try { + mapper.readTree(pgValue) + } catch (_: Exception) { + null + } + if (pgNode == null) { + mismatch(pgValue, "<${legacy.nodeType}>", "PG value is not parseable as JSON while legacy is ${legacy.nodeType}") + return + } + compareNodes(flid, path, pgNode, legacy, acc) + return + } + + // 基本类型:文本值比较 + if (pgValue != legacy.asText()) { + mismatch(pgValue, legacy.asText(), "Value mismatch: $pgValue != ${legacy.asText()}") + } + } + private fun compareNodes(flid: String, path: String, n1: JsonNode?, n2: JsonNode?, acc: MutableList) { if (n1 == null && n2 == null) return diff --git a/src/main/resources/db/migration/V1.1.0__flight_schd.sql b/src/main/resources/db/migration/V1.1.0__flight_schd.sql index 04523e6..f11d947 100644 --- a/src/main/resources/db/migration/V1.1.0__flight_schd.sql +++ b/src/main/resources/db/migration/V1.1.0__flight_schd.sql @@ -1,32 +1,115 @@ -- ===================================================================== --- 自有 PostgreSQL 库 · 阶段 A 运营航班表与计划代迁移 · ACM2-28 定案(选项 C) +-- 自有库 · 阶段 A 运营航班表与计划代迁移 · ACM2-28 定案(选项 C)· 11g 修订 -- --------------------------------------------------------------------- -- 存储边界(ACM2-28 定案,取代 ACM2-12「FLIGHT_STATE 缓做不落表、Redis 永续权威」口径): --- · FLIGHT_SCHD:SCHD 日计划快照 + FLOP/ADFT 增量合并后的当前运营航班权威态,落自有 PG。 --- · SCHD_GEN:各日代版本与当前代航班全量集合(差删依据),由 Redis 回归自有 PG。 --- · Redis:彻底退出动态权威与全部写路径;阶段 A 读写事务全在自有 PG 原子完成。 +-- · FLIGHT_SCHD:SCHD 日计划快照 + FLOP/ADFT 增量合并后的当前运营航班权威态,一行一航班。 +-- · SCHD_GEN / SCHD_GEN_FLID:各日代版本(SQL CAS)与当前代航班全量集合(差删依据)。 +-- · Redis:彻底退出动态权威与全部写路径;阶段 A 读写事务全在自有库原子完成。 -- · 增量 Upsert 语义: --- - 快照全量写入(DNLD):声明/更新 FDAY 归属,更新 FLTR_JSON; --- - 增量写入(FLOP/ADFT):新插置 NULL,已有行通过 ON CONFLICT 保留原 FDAY; --- - 按代差删域化:DELETE FROM FLIGHT_SCHD WHERE FDAY = :day AND FLID = ANY(:diffSet) +-- - 快照全量写入(DNLD):声明/更新 FDAY 归属,航班字段整体替换; +-- - 增量写入(FLOP/ADFT):新插置 NULL FDAY,已有行保留原 FDAY,按字段列更新 +-- (与 legacy flightInfo hash「仅新增/覆盖 field」同语义); +-- - 按代差删域化:DELETE FROM FLIGHT_SCHD WHERE FDAY = :day AND FLID IN (:diffSet) -- 仅删除仍属旧代的行,ADFT(FDAY=NULL)与已迁移至新代的同 FLID 行天然存活。 --- · 时间列规范:统一使用 TIMESTAMP(6) WITH TIME ZONE (TIMESTAMPTZ)。 +-- · 时间列规范:统一使用 TIMESTAMP(6) WITH TIME ZONE。 +-- --------------------------------------------------------------------- +-- 11g 修订(现场数据库版本 11g,废除 FLTR_JSON JSONB 整文档存储): +-- · 表形态 = 运营航班宽表:一行一航班,字段即列,天然可索引、可直查、可加约束; +-- · 列清单锚定 legacy SIS 契约(unisysaodbsis.xsd SCHD.FLTR 记录,ACMA-4 基线, +-- 32 个 Handler 的 KEEP/FIX 翻译即以此为闭集): +-- - 标量字段 → 真实列(约 45 列 + legacy 派生列 ABDG/LPSDT/ABN); +-- - 1:N 明细集合(登机口/柜台/转盘/桥/延误/航线等,单航班可达 99 条)各占一列 +-- 存序列化文本(TEXT;11g 移植为 CLOB),阶段 2/3 语义钉死后按需升独立子表; +-- · 值语义保持 legacy 字符串原样(不做库端类型转换),保证影子对拍逐字段保真; +-- 范围查询需要的类型化列(如 SODT→DATE)按查询需求逐列后置提升。 -- ===================================================================== --- ① 运营航班表(当前运营航班动态权威,FLID 与 legacy flightInfo 同键) +-- ① 运营航班表(当前运营航班动态权威,FLID 与 legacy flightInfo 同键;一行一航班) CREATE TABLE FLIGHT_SCHD ( - FLID VARCHAR(32) NOT NULL PRIMARY KEY, -- AODB 唯一航班标识 - FDAY DATE NULL, -- 所属日计划代;NULL = 非代所有(ADFT/增量增建) - FLTR_JSON JSONB NOT NULL, -- 航班全量记录(与 legacy flightInfo value 同构) - CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, - UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL + FLID VARCHAR(32) NOT NULL PRIMARY KEY, -- AODB 唯一航班标识 + FDAY DATE NULL, -- 所属日计划代;NULL = 非代所有(ADFT/增量增建) + -- ==== SCHD.FLTR 标量字段(unisysaodbsis.xsd,值保持 legacy 字符串原样)==== + ALCD VARCHAR(64) NULL, -- 航空公司代码 + ALSC VARCHAR(64) NULL, -- 航空公司简称 + FLNO VARCHAR(64) NULL, -- 航班号 + MVIN VARCHAR(64) NULL, -- 进离港标识(A=到港 D=离港) + SODT VARCHAR(64) NULL, -- 计划时间(ddMMMyyHHmm) + FLTY VARCHAR(64) NULL, -- 航班类型 + FLIN VARCHAR(64) NULL, -- 国内/国际/混合属性 + ACFT VARCHAR(64) NULL, -- 机型 + RENO VARCHAR(64) NULL, -- 机尾号 + 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, -- 共享承运人代码 + CSFT VARCHAR(64) NULL, -- 共享航班号 + MAID VARCHAR(32) NULL, -- 共享主航班 FLID + ESTT VARCHAR(64) NULL, -- 预计时间 + ACTT VARCHAR(64) NULL, -- 实际时间 + STND VARCHAR(64) NULL, -- 备降站 + PHAG VARCHAR(64) NULL, -- 地服代理(值机) + CNCL VARCHAR(64) NULL, -- 取消时间 + REMC VARCHAR(80) NULL, -- 备注 + BOTM VARCHAR(64) NULL, -- 摆渡车时间 + LACL VARCHAR(64) NULL, -- 行李确认时间 + FINT VARCHAR(64) NULL, -- 完成时间 + APPT VARCHAR(64) NULL, -- 旅客到达时间 + EGSR VARCHAR(64) NULL, -- 保障开始时间 + EGST VARCHAR(64) NULL, -- 保障结束时间 + FHAG VARCHAR(64) NULL, -- 地服代理(货运) + MHAG VARCHAR(64) NULL, -- 地服代理(机务) + VIPP VARCHAR(64) NULL, -- VIP 旅客数 + VIPR VARCHAR(64) NULL, -- VIP 等级 + LBNO VARCHAR(64) NULL, -- 行李件数 + LBWT VARCHAR(64) NULL, -- 行李重量 + PAXC VARCHAR(64) NULL, -- 旅客计数 + EXSC VARCHAR(64) NULL, -- 例外代码 + EXSR VARCHAR(256) NULL, -- 例外原因 + FTSS VARCHAR(64) NULL, -- 航班状态 + PEDT VARCHAR(64) NULL, -- 前序估计时间 + NEAT VARCHAR(64) NULL, -- 下机完成时间 + PADT VARCHAR(64) NULL, -- 旅客登机时间 + NAAT VARCHAR(64) NULL, -- 关门时间 + -- ==== legacy flightInfo 派生标量字段(user-stories §SCHD.FLTR)==== + ABDG VARCHAR(64) NULL, -- 当前登机桥(派生) + LPSDT VARCHAR(64) NULL, -- 历史计划机位(派生) + ABN VARCHAR(64) NULL, -- 异常标识(派生) + -- ==== 1:N 明细集合(序列化文本,一列一集合;阶段 2/3 按需升独立子表)==== + ROUT_TXT TEXT NULL, -- 航线(ROUTEDAILY ×6) + ERUT_TXT TEXT NULL, -- 扩展航线(ROUTEDAILY ×2..7) + CHDT_TXT TEXT NULL, -- 廊桥/ chute 数据(×99) + GTDT_TXT TEXT NULL, -- 登机口(×99) + PSDT_TXT TEXT NULL, -- 机位(×9) + CKDT_TXT TEXT NULL, -- 值机柜台(×99) + CLDT_TXT TEXT NULL, -- 行李转盘(×99) + DELY_TXT TEXT NULL, -- 延误(不限) + CHOT_TXT TEXT NULL, -- 客桥桥载(×99) + ABTM_TXT TEXT NULL, -- 登机桥(×99) + SRVT_TXT TEXT NULL, -- 服务(不限) + VIPF_TXT TEXT NULL, -- VIP(不限) + MAFL_TXT TEXT NULL, -- 共享航班列表(legacy 派生集合) + FDIV_TXT TEXT NULL, -- 备降明细 + FRET_TXT TEXT NULL, -- 返航明细 + FLAB_TXT TEXT NULL, -- 中止明细 + CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, + UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL ); CREATE INDEX idx_flight_schd_fday ON FLIGHT_SCHD (FDAY); -- 日代扫描与影子对拍分片 +-- 运营查询索引按需逐列追加(11g/PG 均为普通 B-tree,零方言成本) -- ② 日计划代与差删元数据表(代版本推进与版本 CAS) CREATE TABLE SCHD_GEN ( FDAY DATE NOT NULL PRIMARY KEY, -- 计划日(yyyy-MM-dd) VERSION BIGINT NOT NULL, -- 代版本号(单调递增) - FLIDS_JSON JSONB NOT NULL, -- 当前代有效航班 FLID 集合 UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL ); + +-- ③ 当前代有效航班 FLID 集合(差删依据) +CREATE TABLE SCHD_GEN_FLID ( + FDAY DATE NOT NULL, -- 所属计划日(FK 级联删除) + FLID VARCHAR(32) NOT NULL, -- 当前代有效航班 + CONSTRAINT pk_schd_gen_flid PRIMARY KEY (FDAY, FLID), + CONSTRAINT fk_schd_gen_flid FOREIGN KEY (FDAY) REFERENCES SCHD_GEN (FDAY) ON DELETE CASCADE +); diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlightSchdJdbcPgTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlightSchdJdbcPgTest.kt index bf327b7..3283526 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlightSchdJdbcPgTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlightSchdJdbcPgTest.kt @@ -77,7 +77,7 @@ class FlightSchdJdbcPgTest { val day = "2026-09-07" // 构造 250 条记录跨越 batchSize 200 门限 val flights = (1..250).map { i -> - "TEST_FL_$i" to """{"FLID":"TEST_FL_$i","air":"CA$i"}""" + "TEST_FL_$i" to mapOf("FLID" to "TEST_FL_$i", "FLNO" to "CA$i") } repo.upsertSnapshotBatch(day, flights) @@ -89,7 +89,7 @@ class FlightSchdJdbcPgTest { // 点查测试 val f1 = repo.findByFlid("TEST_FL_1") assertNotNull(f1) - assertTrue(f1!!.contains("CA1")) + assertEquals("CA1", f1!!["FLNO"]) // 多键点查测试(分批查) val subset = (1..50).map { "TEST_FL_$it" } @@ -107,26 +107,26 @@ class FlightSchdJdbcPgTest { fun `PG dialect - incremental upsert preserves existing FDAY and sets NULL for new`() { val day = "2026-09-07" // 1. 快照写入 REG_01,FDAY 为 2026-09-07 - repo.upsertSnapshotBatch(day, listOf("TEST_REG_01" to """{"FLID":"TEST_REG_01","v":1}""")) + repo.upsertSnapshotBatch(day, listOf("TEST_REG_01" to mapOf("FLID" to "TEST_REG_01", "REMC" to "1"))) // 2. 增量更新 REG_01(已有行)与 ADFT_01(新行) val changes = listOf( - FlightChange("TEST_REG_01", """{"FLID":"TEST_REG_01","v":2}"""), - FlightChange("TEST_ADFT_01", """{"FLID":"TEST_ADFT_01","v":1}"""), + FlightChange("TEST_REG_01", mapOf("FLID" to "TEST_REG_01", "REMC" to "2")), + FlightChange("TEST_ADFT_01", mapOf("FLID" to "TEST_ADFT_01", "REMC" to "1")), ) repo.upsertIncremental(changes) // 3. 验证 REG_01 内容更新但 FDAY 依然保留为 2026-09-07 - val reg01 = ds.queryOne("SELECT fday, fltr_json FROM flight_schd WHERE flid = 'TEST_REG_01'", {}) { rs -> - rs.getDate("fday")?.toString() to rs.getString("fltr_json") + val reg01 = ds.queryOne("SELECT fday, remc FROM flight_schd WHERE flid = 'TEST_REG_01'", {}) { rs -> + rs.getDate("fday")?.toString() to rs.getString("remc") } assertNotNull(reg01) assertEquals("2026-09-07", reg01!!.first) - assertTrue(reg01.second.contains("\"v\": 2") || reg01.second.contains("\"v\":2")) + assertEquals("2", reg01.second) // 4. 验证 ADFT_01 新插入行 FDAY 恒为 NULL - val adft01 = ds.queryOne("SELECT fday, fltr_json FROM flight_schd WHERE flid = 'TEST_ADFT_01'", {}) { rs -> - rs.getDate("fday")?.toString() to rs.getString("fltr_json") + val adft01 = ds.queryOne("SELECT fday, remc FROM flight_schd WHERE flid = 'TEST_ADFT_01'", {}) { rs -> + rs.getDate("fday")?.toString() to rs.getString("remc") } assertNotNull(adft01) assertNull(adft01!!.first) @@ -171,7 +171,7 @@ class FlightSchdJdbcPgTest { val day = "2026-09-07" try { ds.withTransaction { - repo.upsertSnapshotBatch(day, listOf("TEST_TX_01" to """{"tx":1}""")) + repo.upsertSnapshotBatch(day, listOf("TEST_TX_01" to mapOf("REMC" to "tx"))) repo.putGenIfVersion(day, 0L, FlightSchdRepository.GenMeta(day, 1L, setOf("TEST_TX_01"))) // 模拟事务内部抛出异常 throw IllegalStateException("forced-abort-for-rollback-test") @@ -189,7 +189,7 @@ class FlightSchdJdbcPgTest { fun `PG dialect - ES history sweep deleteByFlids idempotent batch execution`() { val day = "2026-09-07" val flids = (1..50).map { "TEST_SWEEP_$it" } - repo.upsertSnapshotBatch(day, flids.map { it to """{"flid":"$it"}""" }) + repo.upsertSnapshotBatch(day, flids.map { it to mapOf("FLID" to it) }) assertEquals(50, repo.findByFlids(flids).size) // 第一次分批删除 @@ -211,7 +211,7 @@ class FlightSchdJdbcPgTest { val day = "2026-09-07" val fixedInstant = Instant.parse("2026-09-07T08:30:15.123456Z") - repo.upsertSnapshotBatch(day, listOf("TEST_TZ_01" to """{"tz":1}"""), fixedInstant) + repo.upsertSnapshotBatch(day, listOf("TEST_TZ_01" to mapOf("FLID" to "TEST_TZ_01")), fixedInstant) repo.putGenIfVersion(day, 0L, FlightSchdRepository.GenMeta(day, 1L, setOf("TEST_TZ_01")), fixedInstant) // 东京时区下回读 @@ -249,8 +249,8 @@ class FlightSchdJdbcPgTest { conn.createStatement().use { it.execute("SET TIME ZONE 'Asia/Shanghai'") } conn.prepareStatement( """ - INSERT INTO flight_schd (flid, fday, fltr_json, created_at, updated_at) - VALUES ('TEST_SESS_01', '2026-09-07', '{"session":"shanghai"}'::jsonb, ?, ?) + INSERT INTO flight_schd (flid, fday, remc, created_at, updated_at) + VALUES ('TEST_SESS_01', '2026-09-07', 'shanghai', ?, ?) """.trimIndent(), ).use { ps -> ps.setTimestamp(1, writeInstant.toSqlTimestamp()) diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJobTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJobTest.kt index 774cbb4..6ec6dae 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJobTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJobTest.kt @@ -40,9 +40,9 @@ class HistorySweepJobTest { @Test fun `Scenario 1 - ES All Success removes all confirmed flights and cleans old gens`() { val flights = listOf( - "F1" to """{"FLID":"F1"}""", - "F2" to """{"FLID":"F2"}""", - "F3" to """{"FLID":"F3"}""", + "F1" to mapOf("FLID" to "F1"), + "F2" to mapOf("FLID" to "F2"), + "F3" to mapOf("FLID" to "F3"), ) flightSchd.upsertSnapshotBatch("2026-09-01", flights) flightSchd.putGenIfVersion("2026-09-01", 0L, FlightSchdRepository.GenMeta("2026-09-01", 1L, setOf("F1", "F2", "F3"))) @@ -68,9 +68,9 @@ class HistorySweepJobTest { @Test fun `Scenario 2 - ES Partial Success deletes only successful flights and keeps failed ones`() { val flights = listOf( - "F1" to """{"FLID":"F1"}""", - "F2" to """{"FLID":"F2"}""", - "F3" to """{"FLID":"F3"}""", + "F1" to mapOf("FLID" to "F1"), + "F2" to mapOf("FLID" to "F2"), + "F3" to mapOf("FLID" to "F3"), ) flightSchd.upsertSnapshotBatch("2026-09-01", flights) @@ -92,8 +92,8 @@ class HistorySweepJobTest { @Test fun `Scenario 3 - ES All Failure deletes nothing and keeps all candidates in PG`() { val flights = listOf( - "F1" to """{"FLID":"F1"}""", - "F2" to """{"FLID":"F2"}""", + "F1" to mapOf("FLID" to "F1"), + "F2" to mapOf("FLID" to "F2"), ) flightSchd.upsertSnapshotBatch("2026-09-01", flights) @@ -112,7 +112,7 @@ class HistorySweepJobTest { @Test fun `Scenario 4 - Crash after ES success before PG delete recovers idempotently on next run`() { val flights = listOf( - "F1" to """{"FLID":"F1"}""", + "F1" to mapOf("FLID" to "F1"), ) flightSchd.upsertSnapshotBatch("2026-09-01", flights) @@ -150,7 +150,7 @@ class HistorySweepJobTest { // 5. 场景五:PG 删除重复执行(无副作用与幂等性) @Test fun `Scenario 5 - Replaying deleteByFlids is completely idempotent with no side effects`() { - val flights = listOf("F1" to """{"FLID":"F1"}""") + val flights = listOf("F1" to mapOf("FLID" to "F1")) flightSchd.upsertSnapshotBatch("2026-09-01", flights) // 首次删除 diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FlightSchdInvariantTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FlightSchdInvariantTest.kt index e5d0447..07eac2f 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FlightSchdInvariantTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FlightSchdInvariantTest.kt @@ -115,11 +115,11 @@ class FlightSchdInvariantTest { val handler = object : Handler { override val kind: MsgKind = MsgKind.Flop("DELY") - override fun decide(flightView: Map, msg: DecodedMessage): Decision { + override fun decide(flightView: Map>, msg: DecodedMessage): Decision { return Decision( - flightChanges = listOf(FlightChange("F101", """{"FLID":"F101","status":"DELAYED"}""")), + flightChanges = listOf(FlightChange("F101", mapOf("FLID" to "F101", "STAT" to "DELAYED"))), msgNotifies = listOf(NotifyPayload("""{"flid":"F101","event":"DELAY"}""")), - schdPush = listOf(SchdPush("F101", """{"FLID":"F101","status":"DELAYED"}""")), + schdPush = listOf(SchdPush("F101", """{"FLID":"F101","STAT":"DELAYED"}""")), ) } } @@ -133,10 +133,10 @@ class FlightSchdInvariantTest { val head = procState.headUnfinished()!! processor.processOne(head) - // 断言:FLIGHT_SCHD 有更新 - val fltr = flightSchd.findByFlid("F101") - assertNotNull(fltr) - assertTrue(fltr!!.contains("DELAYED")) + // 断言:FLIGHT_SCHD 有更新(字段列 DELAYED 已落库) + val fields = flightSchd.findByFlid("F101") + assertNotNull(fields) + assertEquals("DELAYED", fields!!["STAT"]) // 断言:MSG_EVENT 写入了 KAFKA_MSG 和 KAFKA_SCHD 两条事件 val evtMsg = msgEvents.headUnsent(Targets.KAFKA_MSG) @@ -167,8 +167,8 @@ class FlightSchdInvariantTest { // 配置 staging 解析模拟产出 2 条航班 val flights = listOf( - "FL_01" to """{"FLID":"FL_01","air":"CA1234"}""", - "FL_02" to """{"FLID":"FL_02","air":"MU5678"}""", + "FL_01" to mapOf("FLID" to "FL_01", "FLNO" to "CA1234"), + "FL_02" to mapOf("FLID" to "FL_02", "FLNO" to "MU5678"), ) SnapshotFlow.StageResult.parser = { SnapshotFlow.StageResult.Ok(day, flights) } @@ -228,7 +228,7 @@ class FlightSchdInvariantTest { val meta = MetaFields("AODB", "SCHD", "DNLD", 2L, 20260907033000L) val decoded = DecodedMessage(meta, MsgKind.Schd(MsgKind.SchdSubtype.DNLD), "") - val flights = listOf("FL_01" to """{"FLID":"FL_01"}""") + val flights = listOf("FL_01" to mapOf("FLID" to "FL_01")) SnapshotFlow.StageResult.parser = { SnapshotFlow.StageResult.Ok(day, flights) } // 先预置日代版本为 5L(模拟另一并发实例已经推进了版本) @@ -268,11 +268,11 @@ class FlightSchdInvariantTest { val day = "2026-09-07" // 1. 增量更新写入一条临时加飞航班 ADFT(FDAY 为 NULL) - val adftChange = FlightChange(flid = "ADFT_888", payloadJson = """{"FLID":"ADFT_888","type":"ADFT"}""") + val adftChange = FlightChange(flid = "ADFT_888", fields = mapOf("FLID" to "ADFT_888", "FLTY" to "ADFT")) flightSchd.upsertIncremental(listOf(adftChange)) // 2. 写入旧代的一条定期计划航班 REG_OLD(FDAY 为 2026-09-07) - flightSchd.upsertSnapshotBatch(day, listOf("REG_OLD" to """{"FLID":"REG_OLD","type":"REG"}""")) + flightSchd.upsertSnapshotBatch(day, listOf("REG_OLD" to mapOf("FLID" to "REG_OLD", "FLTY" to "REG"))) flightSchd.putGenIfVersion(day, 0L, FlightSchdRepository.GenMeta(day, 1L, setOf("REG_OLD"))) // 3. 执行下一轮快照 DNLD,新代仅包含 REG_NEW(REG_OLD 不在新代中,属于待删差集;ADFT 也不在新代中) @@ -281,7 +281,7 @@ class FlightSchdInvariantTest { val meta = MetaFields("AODB", "SCHD", "DNLD", 3L, 20260907040000L) val decoded = DecodedMessage(meta, MsgKind.Schd(MsgKind.SchdSubtype.DNLD), "") - val newFlights = listOf("REG_NEW" to """{"FLID":"REG_NEW","type":"REG"}""") + val newFlights = listOf("REG_NEW" to mapOf("FLID" to "REG_NEW", "FLTY" to "REG")) SnapshotFlow.StageResult.parser = { SnapshotFlow.StageResult.Ok(day, newFlights) } val snapshotFlow = SnapshotFlow(procState, flightSchd, msgEvents, reqTrack, procFailure, txManager, inbox) @@ -298,7 +298,7 @@ class FlightSchdInvariantTest { // ③ 增量 ADFT_888 航班由于 FDAY=NULL 天然存活、绝不被误删! val adftRecord = flightSchd.findByFlid("ADFT_888") assertNotNull(adftRecord) - assertTrue(adftRecord!!.contains("ADFT_888")) + assertEquals("ADFT_888", adftRecord!!["FLID"]) } @Test @@ -308,11 +308,11 @@ class FlightSchdInvariantTest { val dayNew = "2026-09-08" // 初始属旧代 - flightSchd.upsertSnapshotBatch(dayOld, listOf("FL_MIG" to """{"FLID":"FL_MIG","day":"07"}""")) + flightSchd.upsertSnapshotBatch(dayOld, listOf("FL_MIG" to mapOf("FLID" to "FL_MIG", "REMC" to "07"))) flightSchd.putGenIfVersion(dayOld, 0L, FlightSchdRepository.GenMeta(dayOld, 1L, setOf("FL_MIG"))) // 随后 09-08 快照写入将 FDAY 更新为 2026-09-08 - flightSchd.upsertSnapshotBatch(dayNew, listOf("FL_MIG" to """{"FLID":"FL_MIG","day":"08"}""")) + flightSchd.upsertSnapshotBatch(dayNew, listOf("FL_MIG" to mapOf("FLID" to "FL_MIG", "REMC" to "08"))) // 此时 09-07 再次执行差删(差集中包含 FL_MIG) val deleted = flightSchd.deleteDiffByDay(dayOld, listOf("FL_MIG")) @@ -334,7 +334,7 @@ class FlightSchdInvariantTest { TimeZone.setDefault(TimeZone.getTimeZone(ZoneId.of("Asia/Tokyo"))) val now = Instant.parse("2026-09-07T08:00:00.123456Z") - flightSchd.upsertSnapshotBatch("2026-09-07", listOf("TZ_01" to """{"test":true}"""), now) + flightSchd.upsertSnapshotBatch("2026-09-07", listOf("TZ_01" to mapOf("FLID" to "TZ_01")), now) flightSchd.putGenIfVersion("2026-09-07", 0L, FlightSchdRepository.GenMeta("2026-09-07", 1L, setOf("TZ_01")), now) val gen = flightSchd.getGen("2026-09-07") diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/MessageProcessorTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/MessageProcessorTest.kt index 561c491..653034d 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/MessageProcessorTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/MessageProcessorTest.kt @@ -18,6 +18,7 @@ import com.gzzn.omms.msgexchange.domain.ProcStatus import com.gzzn.omms.msgexchange.domain.SchdPush import com.gzzn.omms.msgexchange.domain.Targets import com.gzzn.omms.msgexchange.infra.persistence.CminmsgInboxRepository +import com.gzzn.omms.msgexchange.infra.persistence.FlightFields import com.gzzn.omms.msgexchange.infra.persistence.MsgEventRepository import com.gzzn.omms.msgexchange.infra.persistence.ProcStateRepository import com.gzzn.omms.msgexchange.infra.persistence.PumpJobRepository @@ -148,23 +149,23 @@ class MessageProcessorTest { } private class FakeFlightSchd : FlightSchdRepository { - val flights = mutableMapOf() + val flights = mutableMapOf() val gens = mutableMapOf() val incrementalChanges = mutableListOf() - override fun upsertSnapshotBatch(day: String, flights: List>, now: Instant) { + override fun upsertSnapshotBatch(day: String, flights: List>, now: Instant) { this.flights.putAll(flights) } override fun upsertIncremental(changes: List, now: Instant) { incrementalChanges.addAll(changes) - changes.forEach { flights[it.flid] = it.payloadJson } + changes.forEach { flights[it.flid] = it.fields } } override fun deleteDiffByDay(day: String, delFlids: Collection): Int = 0 - override fun findByFlid(flid: String): String? = flights[flid] - override fun findByFlids(flids: Collection): Map = + override fun findByFlid(flid: String): FlightFields? = flights[flid] + override fun findByFlids(flids: Collection): Map = flids.mapNotNull { f -> flights[f]?.let { f to it } }.toMap() - override fun findByDay(day: String): List> = flights.map { it.key to it.value } - override fun findAll(): Map = flights.toMap() + override fun findByDay(day: String): List> = flights.map { it.key to it.value } + override fun findAll(): Map = flights.toMap() override fun deleteByFlids(flids: Set): Int = 0 override fun getGen(day: String): FlightSchdRepository.GenMeta? = gens[day] override fun putGenIfVersion(day: String, expected: Long, newGen: FlightSchdRepository.GenMeta, now: Instant): Boolean { @@ -193,9 +194,9 @@ class MessageProcessorTest { private fun deliveHandler() = object : Handler { override val kind: MsgKind = MsgKind.Flop("DELY") - override fun decide(flightView: Map, msg: DecodedMessage): Decision = Decision( + override fun decide(flightView: Map>, msg: DecodedMessage): Decision = Decision( msgNotifies = listOf(NotifyPayload("""{"n":1}""")), - schdPush = listOf(SchdPush(flid = "F1", fltrJson = """{"FLID":"F1","v":2}""")), + schdPush = listOf(SchdPush(flid = "F1", payloadJson = """{"FLID":"F1","v":2}""")), ) } diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/tools/FlightStoreDiffToolTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/tools/FlightStoreDiffToolTest.kt index 44b90a5..a19a79a 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/tools/FlightStoreDiffToolTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/tools/FlightStoreDiffToolTest.kt @@ -10,9 +10,14 @@ class FlightStoreDiffToolTest { private val tool = FlightStoreDiffTool() @Test - fun `AST comparison ignores key ordering differences`() { + fun `Field map comparison ignores legacy JSON key ordering differences`() { val pg = mapOf( - "F1" to """{"flid":"F1","airline":"CA","flightNo":"123","status":"SCHD"}""", + "F1" to mapOf( + "flid" to "F1", + "airline" to "CA", + "flightNo" to "123", + "status" to "SCHD", + ), ) val legacy = mapOf( "F1" to """{"status":"SCHD","flightNo":"123","airline":"CA","flid":"F1"}""", @@ -25,12 +30,16 @@ class FlightStoreDiffToolTest { } @Test - fun `AST comparison normalizes numeric precision and null equivalents`() { + fun `Field map comparison normalizes numeric precision and null equivalents`() { val pg = mapOf( - "F1" to """{"flid":"F1","weight":100.0,"delay":0,"extra":null}""", + "F1" to mapOf( + "flid" to "F1", + "weight" to "100.0", + "delay" to "0", + ), ) val legacy = mapOf( - "F1" to """{"flid":"F1","weight":100,"delay":0.0}""", + "F1" to """{"flid":"F1","weight":100,"delay":0.0,"extra":null}""", ) val report = tool.diff(pg, legacy) @@ -43,8 +52,8 @@ class FlightStoreDiffToolTest { fun `Recognizes known deviation 1 - Cross-day migrated flights protected in PG`() { // PG 中包含跨日迁移未被误删的航班,Legacy Redis 中已被误删 val pg = mapOf( - "F_STABLE" to """{"flid":"F_STABLE"}""", - "F_MIGRATED" to """{"flid":"F_MIGRATED","day":"2026-09-08"}""", + "F_STABLE" to mapOf("flid" to "F_STABLE"), + "F_MIGRATED" to mapOf("flid" to "F_MIGRATED", "FDAY" to "2026-09-08"), ) val legacy = mapOf( "F_STABLE" to """{"flid":"F_STABLE"}""", @@ -62,8 +71,8 @@ class FlightStoreDiffToolTest { @Test fun `Detects real field mismatch and missing flight as unexpected deviation`() { val pg = mapOf( - "F1" to """{"flid":"F1","status":"BOARDING"}""", - "F_EXTRA" to """{"flid":"F_EXTRA"}""", + "F1" to mapOf("flid" to "F1", "status" to "BOARDING"), + "F_EXTRA" to mapOf("flid" to "F_EXTRA"), ) val legacy = mapOf( "F1" to """{"flid":"F1","status":"DEPARTED"}""", @@ -84,4 +93,23 @@ class FlightStoreDiffToolTest { assertTrue(summary.contains("RED (未通过)")) assertTrue(summary.contains("F1")) } + + @Test + fun `Nested collection value is compared as parsed JSON tree`() { + val pg = mapOf( + "F1" to mapOf( + "flid" to "F1", + // 1:N 明细集合列:序列化文本(JSON 数组),按解析树与 legacy 递归比对 + "GTDT" to """[{"GATE":"A1","PGOT":"07SEP260800"}]""", + ), + ) + val legacy = mapOf( + "F1" to """{"flid":"F1","GTDT":[{"PGOT":"07SEP260800","GATE":"A1"}]}""", + ) + + val report = tool.diff(pg, legacy) + assertTrue(report.isGreen) + assertEquals(1, report.matchedCount) + assertEquals(0, report.unexpectedDeviations.size) + } }