refactor(flight-schd): 全宽表零子表收敛 16 个集合并定案 PIPELINE_LOCK 行锁 (ACM2-29)

依据 ACM2-29 新定案【全宽表·零子表】,废除「*_TXT 过渡 + 按需升独立子表」
原规划(对 XSD maxOccurs=99 的过度设计),结合报文样例与成都现场地服实际
规律(登机门 1~2、值机柜台 1~3、转盘 1~2、延误单有效、靠撤桥/轮挡各 1 次)
将 16 个明细集合全部收敛为 FLIGHT_SCHD 宽表标量列或紧凑 VARCHAR 字串。

行为变更:
- V1.2.0 迁移:新增集合平铺槽位列(GTDT×2/CKDT×3/CLDT×2/PSDT×2/CHDT×2,
  跨集合同名属性按 B 前缀/CH 前缀消解)、里程碑标量列(DELY_*、ABTM_A/D、
  CHOT_ON/OFF)、异常前缀标量列(FDIV/FRET/FLAB)、紧凑航路字串
  (ROUT_PATH/ERUT_PATH)与无界集合 JSON 字串列(SRVT/VIPF/MAFL_TEXT);
  删除全部 16 个 *_TXT 文本列;不建任何子表、零 CLOB;
- 仓储层:删除子表替换/回查机制,写侧集合键平铺为列(序号属性 0 = 显式
  删除标记,分舱复用条目抹平去重,超界按定案丢弃),读侧由平铺列重建
  16 个集合键,视图与 legacy flightInfo hash 保持同构(KAFKA_SCHD 线格式
  与 FS7 Diff 逐字段比对不受存储形态影响);
- I5 单写者锁修正:PIPELINE_LOCK 单行 SELECT ... FOR UPDATE 取代
  PG advisory lock,PG/Oracle 11g 同构,无 DBMS_LOCK DBA 特权依赖;
- FlightFieldsJson 对契约结构字段解析 JSON tree 输出原生数组/对象,
  杜绝集合被再次编码成字符串的双重转义。

不变量:消息严格 FIFO、单写者互斥、增量=字段级合并(集合键出现=单资源
集合级全量快照替换)、快照=整体替换、按代域化差删均不变。

迁移影响:本地开发库为一次性测试数据,已重置并由 Flyway 全新应用
V1.0.0→V1.2.0(同文件名内容变更,沿用旧库会触发校验和不匹配)。

验证:./gradlew test 全绿(66 个用例,含本地 PG 真实方言集成:集合槽位
替换/清除、里程碑与航路平铺回读、PIPELINE_LOCK NOWAIT 互斥与提交释放)。
This commit is contained in:
windyboy
2026-09-08 08:48:24 +08:00
parent ff6cec08c8
commit cd5de56937
8 changed files with 651 additions and 45 deletions
+1 -1
View File
@@ -74,7 +74,7 @@ CIIMS / AODB 等上游
## 5. 必须保持的约束
- **消息严格 FIFO**:队头失败并退避时,后续消息仍不能越过它。只有队头完成或按失败策略进入终态后,队列才继续推进。收报重扫和水位设计必须防止较小 ID 漏入队而被后续消息越过。
- **动态状态单写者**`FLIGHT_SCHD` 运营航班表和快照 `SCHD_GEN` 只由主泵单线程写入。生产单实例通过 PG advisory lock 保护,不能通过增加实例或处理线程直接扩容。
- **动态状态单写者**`FLIGHT_SCHD` 运营航班表和快照 `SCHD_GEN` 只由主泵单线程写入。事务通过 `PIPELINE_LOCK``FLIGHT_SCHD_WRITER` 行执行 `SELECT ... FOR UPDATE` 互斥;该方案在 PG/Oracle 11g 均无需数据库扩展或额外 DBA 特权。不能通过增加实例或处理线程直接扩容。
- **身份去重**:同一业务身份只能绑定一条有效处理记录,重复报文不应再次产生业务副作用。具体身份组成和重放规则见设计文档。
- **快照可恢复**:快照覆盖、旧数据清理和版本推进需要原子性与重放保护;不能在恢复时把旧代数据重新写回。
- **作业不与消息混排**`PUMP_JOB` 是独立队列,只在没有消息队头或队头处于退避窗口时执行。执行窗口与饥饿边界需要明确验证。
+35
View File
@@ -117,3 +117,38 @@ A 是最小代价回退位;B 不建议。
`ON CONFLICT` upsert 需改 `MERGE`)、I5 advisory lock 的 11g 替代(`DBMS_LOCK`)、Flyway/驱动
对 11.2 的支持矩阵。若「自有 PG → 现场 11g」成为确定部署形态,需按 ACM2-28 体例开独立决策记录
重开「自有库」平台定案,评估点在方言与并发原语,不在本修订已解决的存储形态。
## 9. ACM2-29 修正落盘:全宽表·零子表(Zero-Subtable Wide Table
R2 的「`*_TXT` 过渡 + 按需升独立子表」经 `docs/legacy/SIS_AODB_RMS-V0.1.md` 报文样例核验与
成都现场(双流/天府)地服实际规律复核,确认为对 XSD 理论上限(maxOccurs="99")的过度设计,
正式废除。16 个明细集合全部收敛为 `FLIGHT_SCHD` 宽表的标准 VARCHAR 标量列或紧凑格式化字串,
零子表、零 CLOB(V1.2.0 迁移落地):
- **平铺标量列**(槽位号即报文序号属性,可按需直接追加普通 B-tree):登机门 `GATE1/2+PGOT/PGCT/GOTM/GCTM/GTYP`
×2 组;值机柜台 `CHKC1..3` 全属性(同柜台分舱复用条目写侧抹平去重);转盘 `BELT1/2` 全属性
(计划时刻 `BPCOT/BPCCT` 专属前缀避免与值机同名);计划机位 `PSST1/2`;离港通道 `CHUT1/2`
(等级/类型 `CHCLS/CHTYP` 专属前缀)。槽位数按现场规律定界(2/3/2/2/2),超界条目按定案丢弃。
- **单值/里程碑标量列**:延误 `DELY_CODE/DELY_STRT/DELY_DURA/DELY_REMC`(业务上任意时刻仅 1 个
有效延误,覆盖语义);靠撤桥 `ABTM_A/ABTM_D`(桥号沿用派生列 ABDG);轮挡 `CHOT_ON/CHOT_OFF`
(机位沿用派生列 STND);异常 `FDIV/FRET/FLAB` 前缀标量列(规范 1:0..1 单值异常)。
- **紧凑 VARCHAR 字串列**(无界/航路型集合,11g VARCHAR2 内联存储,非 CLOB):
`ROUT_PATH/ERUT_PATH`(≤7 站紧凑航路 `"APCD/SCAT/SCDT,..."`);`SRVT_TEXT/VIPF_TEXT/MAFL_TEXT`
(服务/VIP/共享列表条数无 XSD 上限,存 JSON 数组;超长由列约束 fail fast,不静默截断)。
- **消息语义**:运营资源 = 单资源集合级全量快照替换(集合键出现即整集合覆盖,缺槽位置 NULL,
序号属性 "0" 为显式删除标记);外层航班字段仍为字段级增量合并(仅新增/覆盖)。
- **读侧视图**:由平铺列重建 16 个集合键,与 legacy `flightInfo` hash 同构——KAFKA_SCHD 线格式
FLTR JSON 数组)与 FS7 Diff 逐字段比对均依赖该同构视图,存储形态变化不出仓储边界。
已知保真边界:`ABTM/CHOT` 只保留里程碑时刻(资源号沿用 ABDG/STND 派生列);航路空串属性
与缺失属性在紧凑字串中合并为空槽位;`FDIV/FRET/FLAB` 自由文本键名待阶段 2 Handler 钉死。
- **出站序列化**`FlightFieldsJson` 对契约结构字段解析 JSON tree 输出原生数组/对象,
杜绝集合被再次编码成字符串的双重转义。
- **I5 单写者锁修正**:不再依赖 PG advisory lock/11g `DBMS_LOCK`。V1.2.0 创建 `PIPELINE_LOCK`
单行锁表,处理事务对 `FLIGHT_SCHD_WRITER` 行执行 `SELECT ... FOR UPDATE`;连接断开、提交或
回滚时由数据库释放,PG 与 11g 同构,无 DBA 特权依赖。
16 个 `*_TXT` 文本列由 V1.2.0 全部删除;运营查询索引(如 `GATE1/CHKC1/DELY_CODE`)按需以普通
B-tree 逐列追加,11g/PG 零方言成本。
Oracle 11g 的 `MERGE INTO` 方言、Flyway 11.2 支持版本和 ojdbc 认证组合仍必须在目标数据库环境
中完成平台决策与集成验收;不能用 PostgreSQL 测试结果替代该外部环境证据。
@@ -1,7 +1,6 @@
package com.gzzn.omms.msgexchange.infra.persistence
import com.fasterxml.jackson.databind.ObjectMapper
import com.fasterxml.jackson.databind.node.ObjectNode
/**
* 航班字段集(FLIGHT_SCHD 宽表字段)→ KAFKA_SCHD 事件载荷序列化。
@@ -11,9 +10,26 @@ import com.fasterxml.jackson.databind.node.ObjectNode
object FlightFieldsJson {
private val mapper = ObjectMapper()
/** Fields whose legacy wire representation is a JSON array/object rather than a JSON string. */
private val STRUCTURED_FIELDS = setOf(
"ROUT", "ERUT", "CHDT", "GTDT", "PSDT", "CKDT", "CLDT", "DELY",
"CHOT", "ABTM", "SRVT", "VIPF", "MAFL", "FDIV", "FRET", "FLAB",
)
fun toJson(fields: Map<String, String>): String {
val node = mapper.createObjectNode()
fields.toSortedMap().forEach { (k, v) -> node.put(k, v) }
fields.toSortedMap().forEach { (key, value) ->
if (key in STRUCTURED_FIELDS) {
val parsed = runCatching { mapper.readTree(value) }.getOrNull()
if (parsed != null && (parsed.isArray || parsed.isObject || parsed.isNull)) {
node.set<com.fasterxml.jackson.databind.JsonNode>(key, parsed)
} else {
node.put(key, value)
}
} else {
node.put(key, value)
}
}
return node.toString()
}
}
@@ -70,9 +70,12 @@ interface MsgEventRepository {
}
/**
* 阶段 A 运营航班权威与日计划代(ACM2-28 定案 + 11g 修订:废除 FLTR_JSON JSONB 整文档存储):
* 阶段 A 运营航班权威与日计划代(ACM2-28 定案 + ACM2-29 全宽表·零子表定案):
* - 表 FLIGHT_SCHD:一行一航班的运营航班宽表(FLID 主键 + FDAY 可空所属代 + SCHD.FLTR 标量字段列
* + 1:N 明细集合序列化文本列),字段即列、天然可索引可直查,PG/Oracle 11g 方言一致;
* + 集合平铺标量列/紧凑 VARCHAR 字串列),零子表、零 CLOB,字段即列、天然可索引可直查,
* PG/Oracle 11g 方言一致;读侧视图由平铺列重建集合键,与 legacy flightInfo hash 同构;
* - 运营资源消息语义 = 单资源集合级全量快照替换;外层航班增量仍为字段级合并(仅新增/覆盖);
* - 表 PIPELINE_LOCK:单写者行级互斥(SELECT ... FOR UPDATEPG/11g 同构,无 DBMS_LOCK 特权依赖);
* - 表 SCHD_GEN:各日代版本(SQL CAS);表 SCHD_GEN_FLID:当前代有效航班 FLID 集合(差删依据);
* - Redis 退出动态权威与全部写路径;事务 2 与快照发布全在自有库内单事务原子提交;
* - 增量更新(FLOP/ADFT):新插 FDAY=NULL,已有行保留原 FDAY,按字段列更新
@@ -1,5 +1,8 @@
package com.gzzn.omms.msgexchange.infra.persistence.jdbc
import com.fasterxml.jackson.databind.JsonNode
import com.fasterxml.jackson.databind.ObjectMapper
import com.fasterxml.jackson.databind.node.ObjectNode
import com.gzzn.omms.msgexchange.domain.ErrorClass
import com.gzzn.omms.msgexchange.domain.EventStatus
import com.gzzn.omms.msgexchange.domain.MsgEvent
@@ -296,7 +299,16 @@ class JdbcPumpJobRepository(
class JdbcPipelineTransactionManager(
private val ds: DataSource,
) : PipelineTransactionManager {
override fun <T> inTransaction(block: () -> T): T = ds.withTransaction(block)
override fun <T> inTransaction(block: () -> T): T = ds.withTransaction {
// Standard row lock works on PostgreSQL and Oracle 11g without DBMS_LOCK privileges.
// The lock is connection-scoped and is released automatically on commit, rollback, or crash.
ds.queryOne(
"SELECT lock_key FROM pipeline_lock WHERE lock_key = ? FOR UPDATE",
{ ps -> ps.setString(1, "FLIGHT_SCHD_WRITER") },
) { rs -> rs.getString("lock_key") }
?: error("missing FLIGHT_SCHD_WRITER row in PIPELINE_LOCK")
block()
}
}
@Singleton
@@ -317,33 +329,212 @@ class JdbcFlightSchdRepository(
"ABDG", "LPSDT", "ABN",
)
/** 1:N 明细集合键(库列名 = 键 + _TXT 后缀,序列化文本)。 */
private val COLLECTION_KEYS: List<String> = listOf(
"ROUT", "ERUT", "CHDT", "GTDT", "PSDT", "CKDT", "CLDT", "DELY",
"CHOT", "ABTM", "SRVT", "VIPF", "MAFL", "FDIV", "FRET", "FLAB",
/** 异常明细前缀标量列(FDIV/FRET/FLAB:规范 1:0..1 单值异常)。 */
private val EXCEPTION_COLUMNS: List<String> = listOf(
"FDIV_DDES", "FDIV_DDIR", "FDIV_REMC", "FRET_REID", "FRET_RSN",
"FLAB_ARES", "FLAB_RSN",
)
/** 库列名全集(标量原样;集合列带 _TXT 后缀)。 */
private val ALL_COLUMNS: List<String> = SCALAR_COLUMNS + COLLECTION_KEYS.map { "${it}_TXT" }
/**
* 平铺出现槽位集合(ACM2-29 零子表定案):槽位序号即报文序号属性(GTNO/CKNO/...),
* 列名 = legacy 属性名 + 槽位号;跨集合同名属性按 V1.2.0 注释加专属前缀
* (转盘时刻 BPCOT/BPCCT、通道等级/类型 CHCLS/CHTYP)。槽位数按现场地服
* 规律定界(登机门 2、值机柜台 3、转盘 2、计划机位 2、离港通道 2),
* 超界条目按定案丢弃,不再为 XSD maxOccurs=99 的理论上限设计。
*/
private data class SlotCollection(
val key: String,
val seqAttr: String,
val dedupeAttr: String,
val slots: List<Map<String, String>>, // 报文属性 → 库列名
)
/** 字段键 → 库列名白名单;未知键拒绝写入(fail fast,经处理边界落 FAILED(INFRA))。
* FLID 即主键列:写侧忽略该键(主键已承载),读侧合成返回,视图与 legacy hash 同构。 */
private val FIELD_TO_COLUMN: Map<String, String> =
SCALAR_COLUMNS.associateWith { it } + COLLECTION_KEYS.associateWith { "${it}_TXT" }
private fun slotCollection(
key: String,
seqAttr: String,
dedupeAttr: String,
attrs: List<String>,
count: Int,
renames: Map<String, String> = emptyMap(),
): SlotCollection = SlotCollection(
key = key,
seqAttr = seqAttr,
dedupeAttr = dedupeAttr,
slots = (1..count).map { i -> attrs.associateWith { attr -> (renames[attr] ?: attr) + i } },
)
private fun columnToFieldKey(column: String): String =
if (column.endsWith("_TXT")) column.removeSuffix("_TXT") else column
private val OCCURRENCE_COLLECTIONS: List<SlotCollection> = listOf(
slotCollection("GTDT", "GTNO", "GATE", listOf("GATE", "PGOT", "PGCT", "GOTM", "GCTM", "GTYP"), 2),
slotCollection(
"CKDT", "CKNO", "CHKC",
listOf("CHKC", "CCLS", "PCOT", "PCCT", "COTM", "CCTM", "CTYP"), 3,
),
slotCollection(
"CLDT", "CLNO", "BELT",
listOf("BELT", "BCLS", "PCOT", "PCCT", "FBAG", "LBAG", "BTYP"), 2,
renames = mapOf("PCOT" to "BPCOT", "PCCT" to "BPCCT"),
),
slotCollection("PSDT", "PSNO", "PSST", listOf("PSST", "STST", "STET"), 2),
slotCollection(
"CHDT", "CHNO", "CHUT",
listOf("CHUT", "CCLS", "PCBT", "PCET", "CBTM", "CETM", "CTYP"), 2,
renames = mapOf("CCLS" to "CHCLS", "CTYP" to "CHTYP"),
),
)
/** 里程碑/单值集合列(延误单有效覆盖语义、靠撤桥/轮挡时刻)。 */
private val MILESTONE_COLUMNS: List<String> = listOf(
"DELY_CODE", "DELY_STRT", "DELY_DURA", "DELY_REMC",
"ABTM_A", "ABTM_D", "CHOT_ON", "CHOT_OFF",
)
/** 紧凑航路字串列("APCD/SCAT/SCDT,...",≤7 站)。 */
private val PATH_COLUMNS: List<String> = listOf("ROUT_PATH", "ERUT_PATH")
/** 无界集合紧凑 JSON 数组字串列(键名去 _TEXT 后缀即 legacy 视图键,直存直取)。 */
private val TEXT_COLUMNS: List<String> = listOf("SRVT_TEXT", "VIPF_TEXT", "MAFL_TEXT")
/** 库列名全集(读侧 SELECT 与视图重建的稳定顺序,与 V1.2.0 迁移一致)。 */
private val ALL_COLUMNS: List<String> = SCALAR_COLUMNS + EXCEPTION_COLUMNS +
OCCURRENCE_COLLECTIONS.flatMap { c -> c.slots.flatMap(Map<String, String>::values) } +
MILESTONE_COLUMNS + PATH_COLUMNS + TEXT_COLUMNS
private val SCALAR_KEY_SET: Set<String> = SCALAR_COLUMNS.toSet()
/** 16 个集合视图键:写侧消费平铺、读侧重建,视图与 legacy flightInfo hash 同构。 */
private val COLLECTION_KEYS: Set<String> = OCCURRENCE_COLLECTIONS.map { it.key }.toSet() +
setOf("DELY", "ABTM", "CHOT", "ROUT", "ERUT", "SRVT", "VIPF", "MAFL", "FDIV", "FRET", "FLAB")
}
private val mapper = ObjectMapper()
private fun toSqlDate(day: String): java.sql.Date =
java.sql.Date.valueOf(day.trim().take(10))
private fun requireColumns(fields: FlightFields): List<Pair<String, String>> =
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
/** 集合键平铺 + 标量键校验 → (列名, 值);未知键拒绝写入(fail fast,经处理边界落 FAILED(INFRA))。
* FLID 即主键列:写侧忽略该键(主键已承载),读侧合成返回。 */
private fun storagePairs(fields: FlightFields): List<Pair<String, String?>> =
flattenForStorage(fields).mapNotNull { (key, value) ->
if (key == "FLID") null else key to value
}
/** 字段集 → 列名键值(标量键同名直取;16 个集合键消费为平铺列/紧凑字串列;
* 未知键 fail fast)。值为 null 表示显式清空(单资源集合级全量快照替换语义)。 */
private fun flattenForStorage(fields: FlightFields): Map<String, String?> {
val out = mutableMapOf<String, String?>()
fields.forEach { (key, value) ->
when {
key in SCALAR_KEY_SET -> out[key] = value
key == "FLID" || key in COLLECTION_KEYS -> Unit // 主键跳过;集合键由下方平铺消费
else -> throw IllegalArgumentException("unknown flight field: $key")
}
}
flattenExceptions(fields, out)
OCCURRENCE_COLLECTIONS.forEach { c -> flattenSlots(c, fields[c.key], out) }
flattenDelay(fields["DELY"], out)
flattenAirbridge(fields["ABTM"], out)
flattenChocks(fields["CHOT"], out)
fields["ROUT"]?.let { out["ROUT_PATH"] = routeToPath(it) }
fields["ERUT"]?.let { out["ERUT_PATH"] = routeToPath(it) }
fields["SRVT"]?.let { out["SRVT_TEXT"] = it }
fields["VIPF"]?.let { out["VIPF_TEXT"] = it }
fields["MAFL"]?.let { out["MAFL_TEXT"] = it }
return out
}
/** FDIV/FRET/FLAB1:0..1 单值异常对象 → 前缀标量列(自由文本取首个命中的文本键)。 */
private fun flattenExceptions(fields: FlightFields, out: MutableMap<String, String?>) {
fun flatten(key: String, mappings: Map<String, String>, textTarget: String) {
val raw = fields[key] ?: return
val node = runCatching { mapper.readTree(raw) }.getOrNull() ?: return
if (!node.isObject) return
mappings.forEach { (source, target) ->
node[source]?.takeUnless(JsonNode::isNull)?.asText()?.let { out[target] = it }
}
listOf("REMC", "RSN", "value", "#text").asSequence()
.mapNotNull { name -> node[name]?.takeUnless(JsonNode::isNull)?.asText() }
.firstOrNull()
?.let { out[textTarget] = it }
}
flatten("FDIV", mapOf("DDES" to "FDIV_DDES", "DDIR" to "FDIV_DDIR"), "FDIV_REMC")
flatten("FRET", mapOf("REID" to "FRET_REID"), "FRET_RSN")
flatten("FLAB", mapOf("ARES" to "FLAB_ARES"), "FLAB_RSN")
}
/** 出现槽位集合 → 槽位序号列。序号属性 "0" = 显式删除标记;按资源号抹平复用冗余;
* 超出现场定界槽位数的条目按定案丢弃;集合键出现但条目为空 = 全槽位清空。 */
private fun flattenSlots(c: SlotCollection, raw: String?, out: MutableMap<String, String?>) {
if (raw == null) return
val seen = mutableSetOf<String>()
val items = collectionItems(raw)
.filterNot { it[c.seqAttr]?.takeUnless(JsonNode::isNull)?.asText() == "0" }
.filter { item ->
val dedupeKey = item[c.dedupeAttr]?.takeUnless(JsonNode::isNull)?.asText()
dedupeKey == null || seen.add(dedupeKey)
}
.take(c.slots.size)
c.slots.forEachIndexed { index, slot ->
val item = items.getOrNull(index)
slot.forEach { (attr, column) ->
out[column] = item?.get(attr)?.takeUnless(JsonNode::isNull)?.asText()
}
}
}
/** 延误:业务上任意时刻仅 1 个有效延误(覆盖语义)→ 首条延误平铺。 */
private fun flattenDelay(raw: String?, out: MutableMap<String, String?>) {
if (raw == null) return
val first = collectionItems(raw).firstOrNull()
out["DELY_CODE"] = first?.get("CODE")?.takeUnless(JsonNode::isNull)?.asText()
out["DELY_STRT"] = first?.get("STRT")?.takeUnless(JsonNode::isNull)?.asText()
out["DELY_DURA"] = first?.get("DURA")?.takeUnless(JsonNode::isNull)?.asText()
out["DELY_REMC"] = first?.let { f ->
listOf("REMC", "value", "#text").asSequence()
.mapNotNull { name -> f[name]?.takeUnless(JsonNode::isNull)?.asText() }
.firstOrNull()
}
}
/** 靠撤桥:按 ABOP 取靠桥(A)/撤桥(D) 时刻;桥号沿用派生列 ABDG。 */
private fun flattenAirbridge(raw: String?, out: MutableMap<String, String?>) {
if (raw == null) return
val items = collectionItems(raw)
fun timeOf(op: String): String? = items.firstOrNull { it["ABOP"]?.asText() == op }
?.get("AOTM")?.takeUnless(JsonNode::isNull)?.asText()
out["ABTM_A"] = timeOf("A")
out["ABTM_D"] = timeOf("D")
}
/** 轮挡:按 CHID 取上轮挡(ON)/下轮挡(OFF) 时刻;机位沿用派生列 STND。 */
private fun flattenChocks(raw: String?, out: MutableMap<String, String?>) {
if (raw == null) return
val items = collectionItems(raw)
fun timeOf(id: String): String? = items.firstOrNull { it["CHID"]?.asText() == id }
?.get("CHTM")?.takeUnless(JsonNode::isNull)?.asText()
out["CHOT_ON"] = timeOf("ON")
out["CHOT_OFF"] = timeOf("OFF")
}
/** 航路集合 → 紧凑字串 "APCD/SCAT/SCDT,..."(空属性保留空槽位;空集合 → NULL = 显式清空)。 */
private fun routeToPath(raw: String): String? {
val items = collectionItems(raw)
if (items.isEmpty()) return null
return items.joinToString(",") { item ->
listOf("APCD", "SCAT", "SCDT").joinToString("/") { attr ->
item[attr]?.takeUnless(JsonNode::isNull)?.asText() ?: ""
}
}
}
/** 集合值统一形状:JSON 数组(多值)/对象(单值)/null(清空);其余形状 fail fast。 */
private fun collectionItems(raw: String): List<JsonNode> {
val node = mapper.readTree(raw)
return when {
node.isNull -> emptyList()
node.isArray -> node.toList()
node.isObject -> listOf(node)
else -> throw IllegalArgumentException("collection value must be a JSON array or object")
}
}
override fun upsertSnapshotBatch(day: String, flights: List<Pair<String, FlightFields>>, now: Instant) {
@@ -364,11 +555,12 @@ class JdbcFlightSchdRepository(
conn.prepareStatement(sql).use { ps ->
var count = 0
for ((flid, fields) in flights) {
val stored = flattenForStorage(fields)
var i = 1
ps.setString(i++, flid)
ps.setDate(i++, sqlDate)
for (column in ALL_COLUMNS) {
ps.setString(i++, fields[columnToFieldKey(column)])
ps.setString(i++, stored[column])
}
ps.setTimestamp(i++, sqlTimestamp)
ps.setTimestamp(i++, sqlTimestamp)
@@ -412,10 +604,11 @@ class JdbcFlightSchdRepository(
ps.executeBatch()
}
}
// 字段列更新(legacy flightInfo hash 同语义:仅新增/覆盖,不删除缺失字段
// 字段列更新(legacy flightInfo hash 同语义:仅新增/覆盖,不删除缺失字段
// 集合键出现 = 单资源集合级全量快照替换,缺槽位置 NULL)
for (c in changes) {
val columns = requireColumns(c.fields)
if (columns.isEmpty()) continue
val columns = storagePairs(c.fields)
if (columns.isNotEmpty()) {
val sql = "UPDATE flight_schd SET " +
columns.joinToString(", ") { "${it.first} = ?" } +
", updated_at = ? WHERE flid = ?"
@@ -426,6 +619,7 @@ class JdbcFlightSchdRepository(
ps.setString(i, c.flid)
}
}
}
} finally {
conn.releaseIfNotInTransaction()
}
@@ -446,25 +640,105 @@ class JdbcFlightSchdRepository(
return totalDeleted
}
/** 行 → 字段视图:标量列直读,16 个集合键由平铺列重建,与 legacy flightInfo hash 同构
* KAFKA_SCHD 线格式与 FS7 对拍逐字段比对均依赖该同构视图)。 */
private fun mapFlightRow(rs: java.sql.ResultSet): Pair<String, FlightFields> {
val flid = rs.getString("flid")
val fields = linkedMapOf<String, String>()
fields["FLID"] = flid // 视图与 legacy flightInfo hash 同构:FLID 字段常在
val row = linkedMapOf<String, String?>()
for (column in ALL_COLUMNS) {
val value = rs.getString(column) ?: continue
fields[columnToFieldKey(column)] = value
row[column] = rs.getString(column)
}
SCALAR_COLUMNS.forEach { column -> row[column]?.let { fields[column] = it } }
TEXT_COLUMNS.forEach { column -> row[column]?.let { fields[column.removeSuffix("_TEXT")] = it } }
exceptionView("FDIV", listOf("DDES" to "FDIV_DDES", "DDIR" to "FDIV_DDIR", "REMC" to "FDIV_REMC"), row)
?.let { fields["FDIV"] = it }
exceptionView("FRET", listOf("REID" to "FRET_REID", "RSN" to "FRET_RSN"), row)
?.let { fields["FRET"] = it }
exceptionView("FLAB", listOf("ARES" to "FLAB_ARES", "RSN" to "FLAB_RSN"), row)
?.let { fields["FLAB"] = it }
delayView(row)?.let { fields["DELY"] = it }
airbridgeView(row)?.let { fields["ABTM"] = it }
chocksView(row)?.let { fields["CHOT"] = it }
row["ROUT_PATH"]?.let { fields["ROUT"] = routeView(it) }
row["ERUT_PATH"]?.let { fields["ERUT"] = routeView(it) }
OCCURRENCE_COLLECTIONS.forEach { c -> slotView(c, row)?.let { fields[c.key] = it } }
return flid to fields
}
override fun findByFlid(flid: String): FlightFields? {
val row = ds.queryOne(
/** 前缀标量列 → 单值异常对象视图(缺失属性省略;全空 → 无键)。 */
private fun exceptionView(key: String, attrs: List<Pair<String, String>>, row: Map<String, String?>): String? {
val node = mapper.createObjectNode()
attrs.forEach { (attr, column) -> row[column]?.let { node.put(attr, it) } }
return if (node.size() == 0) null else mapper.writeValueAsString(node)
}
/** 延误列 → 0..1 条延误对象数组视图(legacy DELY 为可重复标签的数组形状)。 */
private fun delayView(row: Map<String, String?>): String? {
val node = mapper.createObjectNode()
listOf("CODE" to "DELY_CODE", "STRT" to "DELY_STRT", "DURA" to "DELY_DURA", "REMC" to "DELY_REMC")
.forEach { (attr, column) -> row[column]?.let { node.put(attr, it) } }
return if (node.size() == 0) null else mapper.writeValueAsString(arrayOf(node))
}
/** 靠撤桥时刻 → 保障里程碑数组视图(ASNO/ABOP 为操作序合成;桥号见派生列 ABDG)。 */
private fun airbridgeView(row: Map<String, String?>): String? {
val items = mutableListOf<ObjectNode>()
row["ABTM_A"]?.let {
items += mapper.createObjectNode().apply { put("ASNO", "1"); put("ABOP", "A"); put("AOTM", it) }
}
row["ABTM_D"]?.let {
items += mapper.createObjectNode().apply { put("ASNO", "2"); put("ABOP", "D"); put("AOTM", it) }
}
return if (items.isEmpty()) null else mapper.writeValueAsString(items)
}
/** 轮挡时刻 → 保障里程碑数组视图(CHID 标识上/下轮挡;机位见派生列 STND)。 */
private fun chocksView(row: Map<String, String?>): String? {
val items = mutableListOf<ObjectNode>()
row["CHOT_ON"]?.let {
items += mapper.createObjectNode().apply { put("CSNO", "1"); put("CHID", "ON"); put("CHTM", it) }
}
row["CHOT_OFF"]?.let {
items += mapper.createObjectNode().apply { put("CSNO", "1"); put("CHID", "OFF"); put("CHTM", it) }
}
return if (items.isEmpty()) null else mapper.writeValueAsString(items)
}
/** 紧凑航路字串 → legacy ROUTEDAILY 数组视图(RTNO 按段序合成;空槽位属性省略)。 */
private fun routeView(path: String): String {
val items = path.split(",").mapIndexed { index, segment ->
val parts = segment.split("/", limit = 3)
mapper.createObjectNode().apply {
put("RTNO", (index + 1).toString())
listOf("APCD", "SCAT", "SCDT").forEachIndexed { i, attr ->
if (i < parts.size && parts[i].isNotEmpty()) put(attr, parts[i])
}
}
}
return mapper.writeValueAsString(items)
}
/** 出现槽位列 → legacy 集合数组视图(空槽位跳过;序号属性按槽位合成)。 */
private fun slotView(c: SlotCollection, row: Map<String, String?>): String? {
val items = c.slots.mapIndexed { index, slot ->
mapper.createObjectNode().apply {
put(c.seqAttr, (index + 1).toString())
slot.forEach { (attr, column) -> row[column]?.let { put(attr, it) } }
}
}.filterNot { it.size() <= 1 }
return if (items.isEmpty()) null else mapper.writeValueAsString(items)
}
override fun findByFlid(flid: String): FlightFields? =
ds.queryOne(
"SELECT flid, ${ALL_COLUMNS.joinToString(", ")} FROM flight_schd WHERE flid = ?",
{ ps -> ps.setString(1, flid) },
::mapFlightRow,
)
return row?.second
}
)?.second
override fun findByFlids(flids: Collection<String>): Map<String, FlightFields> {
if (flids.isEmpty()) return emptyMap()
@@ -0,0 +1,118 @@
-- =====================================================================
-- ACM2-29 · 全宽表·零子表(Zero-Subtable Wide Table)定案落地
-- ---------------------------------------------------------------------
-- 废除「16 个明细集合保留 *_TXT 并逐个升独立子表」的原规划(对 XSD
-- maxOccurs=99 的过度设计)。结合 docs/legacy/SIS_AODB_RMS-V0.1.md 报文
-- 样例与成都现场(双流/天府)地服实际规律——登机门 1~2 个、值机柜台专属
-- 1~3 个、转盘 1~2 个、延误单有效、靠撤桥/轮挡各 1 次——全部集合收敛为
-- FLIGHT_SCHD 宽表的标准 VARCHAR 标量列或紧凑格式化字串,零子表、零 CLOB:
-- · 平铺标量列(属性名 + 出现槽位序号,可按需直接追加普通 B-tree,
-- 11g/PG 零方言成本):
-- GTDT → GATE/PGOT/PGCT/GOTM/GCTM/GTYP ×2 组(登机门 1~2 个)
-- CKDT → CHKC/CCLS/PCOT/PCCT/COTM/CCTM/CTYP ×3 组
-- (值机柜台 1~3 个;同柜台分舱复用条目在写侧抹平去重)
-- CLDT → BELT/BCLS/BPCOT/BPCCT/FBAG/LBAG/BTYP ×2 组(转盘 1~2 个;
-- B 前缀时刻 = BELT 专属,避免与值机 PCOT/PCCT 同名)
-- PSDT → PSST/STST/STET ×2 组(计划机位 1~2 个)
-- CHDT → CHUT/CHCLS/PCBT/PCET/CBTM/CETM/CHTYP ×2 组(离港通道 1~2 个;
-- CH 前缀 = CHUTE 专属,避免与值机 CCLS/CTYP 同名)
-- DELY → DELY_CODE/DELY_STRT/DELY_DURA/DELY_REMC(单有效覆盖语义)
-- ABTM → ABTM_A/ABTM_D(靠桥/撤桥时刻;桥号沿用派生列 ABDG)
-- CHOT → CHOT_ON/CHOT_OFF(上/下轮挡时刻;机位沿用派生列 STND)
-- FDIV/FRET/FLAB → 前缀标量列(规范 1:0..1 单值异常)
-- 超出现场定界的条目按定案丢弃(不再为 XSD 理论上限设计)。
-- · 紧凑 VARCHAR 字串列(无界/航路型集合,11g VARCHAR2 内联存储,非 CLOB):
-- ROUT/ERUT → ROUT_PATH/ERUT_PATH(≤7 站紧凑航路 "APCD/SCAT/SCDT,..."
-- SRVT/VIPF/MAFL → *_TEXT(服务/VIP/共享列表条数无 XSD 上限,存 JSON 数组;
-- 超长写侧 fail fast,不允许静默截断)
-- · 值语义保持 legacy 字符串原样;读侧视图由标量列重建集合键,
-- 与 legacy flightInfo hash 同构(KAFKA_SCHD 线格式与 FS7 对拍不变)。
-- · 运营资源消息语义 = 单资源集合级全量快照替换;外层航班字段仍为
-- 字段级增量合并(集合键出现即整集合覆盖,缺槽位置 NULL)。
-- ---------------------------------------------------------------------
-- 交叉实例单写者互斥(I5 修正定案):PIPELINE_LOCK 行级锁。主泵写事务
-- 持有 FLIGHT_SCHD_WRITER 行 SELECT ... FOR UPDATE;连接断开、提交或回滚
-- 时由数据库释放,PG 与 Oracle 11g 同构,无 DBMS_LOCK DBA 特权依赖。
-- =====================================================================
ALTER TABLE FLIGHT_SCHD
-- 登机门 GTDT(现场 1~2 个)
ADD COLUMN GATE1 VARCHAR(64), ADD COLUMN PGOT1 VARCHAR(64), ADD COLUMN PGCT1 VARCHAR(64),
ADD COLUMN GOTM1 VARCHAR(64), ADD COLUMN GCTM1 VARCHAR(64), ADD COLUMN GTYP1 VARCHAR(64),
ADD COLUMN GATE2 VARCHAR(64), ADD COLUMN PGOT2 VARCHAR(64), ADD COLUMN PGCT2 VARCHAR(64),
ADD COLUMN GOTM2 VARCHAR(64), ADD COLUMN GCTM2 VARCHAR(64), ADD COLUMN GTYP2 VARCHAR(64),
-- 值机柜台 CKDT(现场专属 1~3 个,分舱复用条目写侧抹平)
ADD COLUMN CHKC1 VARCHAR(64), ADD COLUMN CCLS1 VARCHAR(64),
ADD COLUMN PCOT1 VARCHAR(64), ADD COLUMN PCCT1 VARCHAR(64),
ADD COLUMN COTM1 VARCHAR(64), ADD COLUMN CCTM1 VARCHAR(64), ADD COLUMN CTYP1 VARCHAR(64),
ADD COLUMN CHKC2 VARCHAR(64), ADD COLUMN CCLS2 VARCHAR(64),
ADD COLUMN PCOT2 VARCHAR(64), ADD COLUMN PCCT2 VARCHAR(64),
ADD COLUMN COTM2 VARCHAR(64), ADD COLUMN CCTM2 VARCHAR(64), ADD COLUMN CTYP2 VARCHAR(64),
ADD COLUMN CHKC3 VARCHAR(64), ADD COLUMN CCLS3 VARCHAR(64),
ADD COLUMN PCOT3 VARCHAR(64), ADD COLUMN PCCT3 VARCHAR(64),
ADD COLUMN COTM3 VARCHAR(64), ADD COLUMN CCTM3 VARCHAR(64), ADD COLUMN CTYP3 VARCHAR(64),
-- 行李转盘 CLDT(现场 1~2 个;B 前缀时刻 = BELT 专属)
ADD COLUMN BELT1 VARCHAR(64), ADD COLUMN BCLS1 VARCHAR(64),
ADD COLUMN BPCOT1 VARCHAR(64), ADD COLUMN BPCCT1 VARCHAR(64),
ADD COLUMN FBAG1 VARCHAR(64), ADD COLUMN LBAG1 VARCHAR(64), ADD COLUMN BTYP1 VARCHAR(64),
ADD COLUMN BELT2 VARCHAR(64), ADD COLUMN BCLS2 VARCHAR(64),
ADD COLUMN BPCOT2 VARCHAR(64), ADD COLUMN BPCCT2 VARCHAR(64),
ADD COLUMN FBAG2 VARCHAR(64), ADD COLUMN LBAG2 VARCHAR(64), ADD COLUMN BTYP2 VARCHAR(64),
-- 计划机位 PSDT(现场 1~2 个)
ADD COLUMN PSST1 VARCHAR(64), ADD COLUMN STST1 VARCHAR(64), ADD COLUMN STET1 VARCHAR(64),
ADD COLUMN PSST2 VARCHAR(64), ADD COLUMN STST2 VARCHAR(64), ADD COLUMN STET2 VARCHAR(64),
-- 离港行李通道 CHDT(现场 1~2 个;CH 前缀 = CHUTE 专属)
ADD COLUMN CHUT1 VARCHAR(64), ADD COLUMN CHCLS1 VARCHAR(64),
ADD COLUMN PCBT1 VARCHAR(64), ADD COLUMN PCET1 VARCHAR(64),
ADD COLUMN CBTM1 VARCHAR(64), ADD COLUMN CETM1 VARCHAR(64), ADD COLUMN CHTYP1 VARCHAR(64),
ADD COLUMN CHUT2 VARCHAR(64), ADD COLUMN CHCLS2 VARCHAR(64),
ADD COLUMN PCBT2 VARCHAR(64), ADD COLUMN PCET2 VARCHAR(64),
ADD COLUMN CBTM2 VARCHAR(64), ADD COLUMN CETM2 VARCHAR(64), ADD COLUMN CHTYP2 VARCHAR(64),
-- 延误 DELY(单有效覆盖语义)
ADD COLUMN DELY_CODE VARCHAR(64),
ADD COLUMN DELY_STRT VARCHAR(64),
ADD COLUMN DELY_DURA VARCHAR(64),
ADD COLUMN DELY_REMC VARCHAR(80),
-- 靠撤桥 ABTM / 轮挡 CHOT(保障里程碑时刻)
ADD COLUMN ABTM_A VARCHAR(64),
ADD COLUMN ABTM_D VARCHAR(64),
ADD COLUMN CHOT_ON VARCHAR(64),
ADD COLUMN CHOT_OFF VARCHAR(64),
-- 异常明细 FDIV/FRET/FLAB(规范 1:0..1 单值异常)
ADD COLUMN FDIV_DDES VARCHAR(64),
ADD COLUMN FDIV_DDIR VARCHAR(64),
ADD COLUMN FDIV_REMC VARCHAR(80),
ADD COLUMN FRET_REID VARCHAR(64),
ADD COLUMN FRET_RSN VARCHAR(80),
ADD COLUMN FLAB_ARES VARCHAR(64),
ADD COLUMN FLAB_RSN VARCHAR(80),
-- 紧凑航路字串(≤7 站,免除多表关联)
ADD COLUMN ROUT_PATH VARCHAR(256),
ADD COLUMN ERUT_PATH VARCHAR(256),
-- 无界集合紧凑 JSON 数组字串(11g VARCHAR2 内联,非 CLOB;超长写侧 fail fast
ADD COLUMN SRVT_TEXT VARCHAR(2000),
ADD COLUMN VIPF_TEXT VARCHAR(512),
ADD COLUMN MAFL_TEXT VARCHAR(512);
-- 16 个集合的过渡 *_TXT 文本列全清场
ALTER TABLE FLIGHT_SCHD
DROP COLUMN ROUT_TXT, DROP COLUMN ERUT_TXT, DROP COLUMN CHDT_TXT,
DROP COLUMN GTDT_TXT, DROP COLUMN PSDT_TXT, DROP COLUMN CKDT_TXT,
DROP COLUMN CLDT_TXT, DROP COLUMN DELY_TXT, DROP COLUMN CHOT_TXT,
DROP COLUMN ABTM_TXT, DROP COLUMN SRVT_TXT, DROP COLUMN VIPF_TXT,
DROP COLUMN MAFL_TXT, DROP COLUMN FDIV_TXT, DROP COLUMN FRET_TXT,
DROP COLUMN FLAB_TXT;
-- 运营查询索引按需逐列追加示例(11g/PG 均为普通 B-tree,零方言成本):
-- CREATE INDEX idx_flight_schd_gate1 ON FLIGHT_SCHD (GATE1);
-- CREATE INDEX idx_flight_schd_chkc1 ON FLIGHT_SCHD (CHKC1);
-- CREATE INDEX idx_flight_schd_dely_code ON FLIGHT_SCHD (DELY_CODE);
-- 交叉实例单写者互斥行(I5SELECT ... FOR UPDATE 行级锁)
CREATE TABLE PIPELINE_LOCK (
LOCK_KEY VARCHAR(32) NOT NULL PRIMARY KEY,
OWNER_INFO VARCHAR(128),
UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL
);
INSERT INTO PIPELINE_LOCK (LOCK_KEY, OWNER_INFO, UPDATED_AT)
VALUES ('FLIGHT_SCHD_WRITER', NULL, CURRENT_TIMESTAMP);
@@ -0,0 +1,28 @@
package com.gzzn.omms.msgexchange.infra.persistence
import com.fasterxml.jackson.databind.ObjectMapper
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
class FlightFieldsJsonTest {
private val mapper = ObjectMapper()
@Test
fun `collection JSON is emitted as a native array`() {
val json = FlightFieldsJson.toJson(
mapOf("FLID" to "F1", "GTDT" to """[{"GTNO":"1","GATE":"A1"}]"""),
)
val root = mapper.readTree(json)
assertTrue(root["GTDT"].isArray)
assertEquals("A1", root["GTDT"][0]["GATE"].asText())
}
@Test
fun `malformed collection text stays a string instead of breaking delivery`() {
val root = mapper.readTree(FlightFieldsJson.toJson(mapOf("GTDT" to "legacy-value")))
assertTrue(root["GTDT"].isTextual)
assertEquals("legacy-value", root["GTDT"].asText())
}
}
@@ -18,14 +18,17 @@ import java.time.ZoneId
import java.util.TimeZone
/**
* ACM2-28 FS6PostgreSQL 真实方言集成测试(JdbcFlightSchdRepository
* ACM2-28 FS6 + ACM2-29PostgreSQL 真实方言集成测试(JdbcFlightSchdRepository
* 验证:
* 1. DNLD 快照批处理写入(JDBC batch200 批次);
* 2. 增量 Upsert FDAY 保留策略(ON CONFLICT 保留原代,新插置 NULL);
* 3. 域化差删(按 FDAY 严格隔离);
* 4. SCHD_GEN SQL CAS(防并发断言);
* 5. 单事务原子回滚(崩溃无残留);
* 6. ES 历史清场 deleteByFlids 幂等删除
* 6. ES 历史清场 deleteByFlids 幂等删除
* 7. ACM2-29 零子表:集合平铺槽位列写读(集合级全量替换/序号 0 清除/里程碑与航路紧凑列),
* 读侧视图重建 legacy 同构集合键;
* 8. I5PIPELINE_LOCK 行级单写者互斥与提交释放。
*/
class FlightSchdJdbcPgTest {
@@ -137,6 +140,135 @@ class FlightSchdJdbcPgTest {
assertNotNull(repo.findByFlid("TEST_ADFT_01"))
}
@Test
fun `resource slots replace the whole collection and zero sequence clears it`() {
val day = "2026-09-07"
val mapper = com.fasterxml.jackson.databind.ObjectMapper()
repo.upsertSnapshotBatch(
day,
listOf(
"TEST_RESOURCE" to mapOf(
"GTDT" to """[{"GTNO":"1","GATE":"A1"},{"GTNO":"2","GATE":"A2"}]""",
),
),
)
// 零子表:两个登机门落平铺槽位列
val slots = ds.queryOne(
"SELECT gate1, pgot1, gate2, pgot2 FROM flight_schd WHERE flid = 'TEST_RESOURCE'",
{},
) { rs -> listOf(rs.getString("gate1"), rs.getString("pgot1"), rs.getString("gate2"), rs.getString("pgot2")) }
assertEquals(listOf("A1", null, "A2", null), slots)
// 读侧视图重建 legacy 同构 GTDT 数组(槽位序号合成 GTNO)
val view = repo.findByFlid("TEST_RESOURCE")!!["GTDT"]?.let { mapper.readTree(it) }
assertEquals(2, view!!.size())
assertEquals("A2", view[1]["GATE"].asText())
assertEquals("2", view[1]["GTNO"].asText())
// 增量 = 单资源集合级全量快照替换:1 个登机门 → 槽位 1 覆盖、槽位 2 清空
repo.upsertIncremental(
listOf(FlightChange("TEST_RESOURCE", mapOf("GTDT" to """[{"GTNO":"3","GATE":"B1"}]"""))),
)
val replaced = mapper.readTree(repo.findByFlid("TEST_RESOURCE")!!["GTDT"])
assertEquals(1, replaced.size())
assertEquals("B1", replaced[0]["GATE"].asText())
// 序号属性 "0" = 显式删除标记 → 全槽位清空,视图无 GTDT 键
repo.upsertIncremental(
listOf(FlightChange("TEST_RESOURCE", mapOf("GTDT" to """[{"GTNO":"0"}]"""))),
)
assertNull(repo.findByFlid("TEST_RESOURCE")!!["GTDT"])
val cleared = ds.queryOne(
"SELECT gate1, gate2 FROM flight_schd WHERE flid = 'TEST_RESOURCE'",
{},
) { rs -> listOf(rs.getString("gate1"), rs.getString("gate2")) }
assertEquals(listOf(null, null), cleared)
}
@Test
fun `milestone delay route and exception collections flatten and rebuild the legacy view`() {
val day = "2026-09-07"
val mapper = com.fasterxml.jackson.databind.ObjectMapper()
repo.upsertSnapshotBatch(
day,
listOf(
"TEST_SCALAR" to mapOf(
"DELY" to """[{"CODE":"YY","STRT":"07SEP261605","DURA":"0200","REMC":"Flight Delayed"}]""",
"ROUT" to """[{"RTNO":"1","APCD":"ORD","SCAT":"07SEP261125","SCDT":"07SEP261315"},{"RTNO":"2","APCD":"MSP"}]""",
"ABTM" to """[{"ASNO":"1","ABOP":"A","AOTM":"07SEP261700"}]""",
"FDIV" to """{"DDES":"PEK","DDIR":"TO","REMC":"weather"}""",
),
),
)
// 库内为零子表平铺列:延误覆盖、紧凑航路字串、里程碑时刻、异常前缀标量
val stored = ds.queryOne(
"SELECT dely_code, dely_strt, rout_path, abtm_a, abtm_d, fdiv_ddes, fdiv_remc FROM flight_schd WHERE flid = 'TEST_SCALAR'",
{},
) { rs ->
listOf(
rs.getString("dely_code"), rs.getString("dely_strt"), rs.getString("rout_path"),
rs.getString("abtm_a"), rs.getString("abtm_d"), rs.getString("fdiv_ddes"), rs.getString("fdiv_remc"),
)
}
assertEquals(
listOf("YY", "07SEP261605", "ORD/07SEP261125/07SEP261315,MSP//", "07SEP261700", null, "PEK", "weather"),
stored,
)
// 读侧视图重建 legacy 同构集合键
val fields = repo.findByFlid("TEST_SCALAR")!!
assertEquals("YY", mapper.readTree(fields["DELY"])[0]["CODE"].asText())
val route = mapper.readTree(fields["ROUT"])
assertEquals(2, route.size())
assertEquals("MSP", route[1]["APCD"].asText())
assertFalse(route[1].has("SCAT"))
assertEquals("A", mapper.readTree(fields["ABTM"])[0]["ABOP"].asText())
assertEquals("PEK", mapper.readTree(fields["FDIV"])["DDES"].asText())
// 增量:空延误数组 = 覆盖清空;航路整集合替换
repo.upsertIncremental(
listOf(
FlightChange("TEST_SCALAR", mapOf("DELY" to "[]")),
FlightChange("TEST_SCALAR", mapOf("ROUT" to """[{"RTNO":"1","APCD":"CTU"}]""")),
),
)
val updated = ds.queryOne(
"SELECT dely_code, rout_path FROM flight_schd WHERE flid = 'TEST_SCALAR'",
{},
) { rs -> listOf(rs.getString("dely_code"), rs.getString("rout_path")) }
assertEquals(listOf(null, "CTU//"), updated)
assertNull(repo.findByFlid("TEST_SCALAR")!!["DELY"])
}
@Test
fun `PIPELINE_LOCK row lock guards the writer transaction and releases on commit`() {
val txManager = JdbcPipelineTransactionManager(ds)
// 持锁事务内:另一连接 NOWAIT 获取同行锁必须立即失败(互斥生效)
assertEquals("held", txManager.inTransaction {
DriverManager.getConnection("jdbc:postgresql://localhost:5432/msgx", "msgx_dev", "msgx_dev_pass").use { other ->
other.autoCommit = false
try {
other.prepareStatement(
"SELECT lock_key FROM pipeline_lock WHERE lock_key = 'FLIGHT_SCHD_WRITER' FOR UPDATE NOWAIT",
).use { ps -> ps.executeQuery() }
throw AssertionError("expected PIPELINE_LOCK contention")
} catch (_: java.sql.SQLException) {
// 期望:行锁已被写者事务持有
} finally {
other.rollback()
}
}
"held"
})
// 提交后锁释放:可再次进入,且种子行完整
assertEquals("again", txManager.inTransaction { "again" })
val owner = ds.queryOne(
"SELECT owner_info FROM pipeline_lock WHERE lock_key = 'FLIGHT_SCHD_WRITER'",
{},
) { rs -> rs.getString("owner_info") }
assertNull(owner)
}
@Test
fun `PG dialect - SCHD_GEN SQL CAS version advance and conflict rejection`() {
val day = "2026-09-07"