refactor(flight-schd): 按现场 11g 约束将运营航班存储改为宽表 (ACM2-28 修订)
现场环境仅提供 Oracle 11g(无任何 JSON 能力),废除 FLTR_JSON JSONB 整文档存储: - V1.1.0 迁移重写:FLIGHT_SCHD 改为一行一航班宽表(SCHD.FLTR 标量字段列 + 1:N 明细集合序列化文本列),SCHD_GEN.FLIDS_JSON 展开为 SCHD_GEN_FLID(FDAY, FLID); 已应用过旧版迁移的环境需重建 schema 重放 - 契约层:FlightChange/Handler 在线视图改为字段集映射(与 legacy flightInfo hash 同构),仓储白名单拒绝未知字段;增量写=字段级合并(hmset 同语义), 快照=整体替换 - 事件载荷:快照投递与投影重建经 FlightFieldsJson 确定性序列化(键序稳定) - 影子对拍:FlightStoreDiffTool 改为 PG 宽表列值 vs legacy FLTR JSON 逐字段 归一化比对,新增嵌套集合比对用例 - 测试:61 用例全绿(不变量门槛 1/2/3、UTC 方言、真实 PG 方言集成、清场五场景) - 文档:decision-flight-state.md 追加 §8 修订记录(不改写定案历史),design.md 表说明同步 11g 方言移植(ON CONFLICT→MERGE、advisory lock 替代、Flyway/驱动矩阵)、 集合列升子表、类型化列提升等未尽事项另立 Plane issue 跟踪。
This commit is contained in:
@@ -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, String>): String {
|
||||
val node = mapper.createObjectNode()
|
||||
fields.toSortedMap().forEach { (k, v) -> node.put(k, v) }
|
||||
return node.toString()
|
||||
}
|
||||
}
|
||||
@@ -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<String, String>
|
||||
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<Pair<String, String>>, now: Instant = Instant.now())
|
||||
/** 快照全量写入(DNLD):强行声明/更新 FDAY 归属,批处理写入,航班字段集整体替换。 */
|
||||
fun upsertSnapshotBatch(day: String, flights: List<Pair<String, FlightFields>>, now: Instant = Instant.now())
|
||||
|
||||
/** 增量更新(FLOP/ADFT):新插 FDAY=NULL,已有行保留原 FDAY。 */
|
||||
/** 增量更新(FLOP/ADFT):新插 FDAY=NULL,已有行保留原 FDAY;字段级合并。 */
|
||||
fun upsertIncremental(changes: List<com.gzzn.omms.msgexchange.domain.FlightChange>, now: Instant = Instant.now())
|
||||
|
||||
/** 按代差删域化:仅删除 FDAY = day 且在 delFlids 中的记录(ADFT 与跨代已迁移行受保护)。 */
|
||||
fun deleteDiffByDay(day: String, delFlids: Collection<String>): Int
|
||||
|
||||
/** 点查单航班 FLTR_JSON。 */
|
||||
fun findByFlid(flid: String): String?
|
||||
/** 点查单航班字段集;无字段行(含航班不存在)返回 null。 */
|
||||
fun findByFlid(flid: String): FlightFields?
|
||||
|
||||
/** 点查多航班 FLTR_JSON。 */
|
||||
fun findByFlids(flids: Collection<String>): Map<String, String>
|
||||
/** 点查多航班字段集:FLID → 字段集,仅含实际存在的航班。 */
|
||||
fun findByFlids(flids: Collection<String>): Map<String, FlightFields>
|
||||
|
||||
/** 按计划日查询当前有效航班。 */
|
||||
fun findByDay(day: String): List<Pair<String, String>>
|
||||
fun findByDay(day: String): List<Pair<String, FlightFields>>
|
||||
|
||||
/** 全量查询(供影子对拍 / 一致性对账)。 */
|
||||
fun findAll(): Map<String, String>
|
||||
fun findAll(): Map<String, FlightFields>
|
||||
|
||||
/**
|
||||
* 历史清场删除:仅删除已确认归档至 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<Pair<String, String>>)
|
||||
fun replaceDay(day: String, flights: List<Pair<String, FlightFields>>)
|
||||
|
||||
fun findByDay(day: String): List<Pair<String, String>>
|
||||
fun findByDay(day: String): List<Pair<String, FlightFields>>
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+167
-77
@@ -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<Set<String>>() {}
|
||||
private val mapper = com.fasterxml.jackson.databind.ObjectMapper()
|
||||
|
||||
private fun serializeFlids(flids: Collection<String>): String = mapper.writeValueAsString(flids)
|
||||
companion object {
|
||||
/** SCHD.FLTR 标量列 + legacy 派生列(与 V1.1.0 迁移一致,全部可空 VARCHAR)。 */
|
||||
private val SCALAR_COLUMNS: List<String> = 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<String> {
|
||||
if (json.isBlank()) return emptySet()
|
||||
return try {
|
||||
mapper.readValue(json, flidsTypeRef)
|
||||
} catch (_: Exception) {
|
||||
emptySet()
|
||||
}
|
||||
/** 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",
|
||||
)
|
||||
|
||||
/** 库列名全集(标量原样;集合列带 _TXT 后缀)。 */
|
||||
private val ALL_COLUMNS: List<String> = SCALAR_COLUMNS + COLLECTION_KEYS.map { "${it}_TXT" }
|
||||
|
||||
/** 字段键 → 库列名白名单;未知键拒绝写入(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 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<Pair<String, String>>, now: Instant) {
|
||||
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
|
||||
}
|
||||
|
||||
override fun upsertSnapshotBatch(day: String, flights: List<Pair<String, FlightFields>>, 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<com.gzzn.omms.msgexchange.domain.FlightChange>, 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<String, FlightFields> {
|
||||
val flid = rs.getString("flid")
|
||||
val fields = linkedMapOf<String, String>()
|
||||
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<String>): Map<String, String> {
|
||||
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<String>): Map<String, FlightFields> {
|
||||
if (flids.isEmpty()) return emptyMap()
|
||||
val result = mutableMapOf<String, String>()
|
||||
val result = linkedMapOf<String, FlightFields>()
|
||||
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<Pair<String, String>> =
|
||||
override fun findByDay(day: String): List<Pair<String, FlightFields>> =
|
||||
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<String, String> =
|
||||
override fun findAll(): Map<String, FlightFields> =
|
||||
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<String>): 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<String>) {
|
||||
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<Pair<String, String>>) =
|
||||
override fun replaceDay(day: String, flights: List<Pair<String, FlightFields>>) =
|
||||
flightSchd.upsertSnapshotBatch(day, flights)
|
||||
|
||||
override fun findByDay(day: String): List<Pair<String, String>> =
|
||||
override fun findByDay(day: String): List<Pair<String, FlightFields>> =
|
||||
flightSchd.findByDay(day)
|
||||
}
|
||||
|
||||
@@ -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<String, Record>()
|
||||
private val records = linkedMapOf<String, Record>()
|
||||
private val gens = mutableMapOf<String, FlightSchdRepository.GenMeta>()
|
||||
|
||||
fun clear() {
|
||||
@@ -214,13 +215,13 @@ class StubFlightSchd : FlightSchdRepository {
|
||||
gens.clear()
|
||||
}
|
||||
|
||||
override fun upsertSnapshotBatch(day: String, flights: List<Pair<String, String>>, now: Instant) {
|
||||
for ((flid, json) in flights) {
|
||||
override fun upsertSnapshotBatch(day: String, flights: List<Pair<String, FlightFields>>, 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<String>): Map<String, String> =
|
||||
flids.mapNotNull { flid -> records[flid]?.let { flid to it.fltrJson } }.toMap()
|
||||
override fun findByFlids(flids: Collection<String>): Map<String, FlightFields> =
|
||||
flids.mapNotNull { flid -> records[flid]?.let { flid to it.fields } }.toMap()
|
||||
|
||||
override fun findByDay(day: String): List<Pair<String, String>> =
|
||||
override fun findByDay(day: String): List<Pair<String, FlightFields>> =
|
||||
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<String, String> =
|
||||
records.mapValues { it.value.fltrJson }
|
||||
override fun findAll(): Map<String, FlightFields> =
|
||||
records.mapValues { it.value.fields }
|
||||
|
||||
override fun deleteByFlids(flids: Set<String>): 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<Pair<String, String>>) =
|
||||
override fun replaceDay(day: String, flights: List<Pair<String, FlightFields>>) =
|
||||
flightSchd.upsertSnapshotBatch(day, flights)
|
||||
|
||||
override fun findByDay(day: String): List<Pair<String, String>> =
|
||||
override fun findByDay(day: String): List<Pair<String, FlightFields>> =
|
||||
flightSchd.findByDay(day)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user