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
@@ -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"