diff --git a/README.md b/README.md index 0640617..7ff1c37 100644 --- a/README.md +++ b/README.md @@ -27,7 +27,8 @@ 自 legacy 仓库复制以自包含;codec 实装依据,ACM2-2/ACM2-3)。 - `db/migration/V2.0.0__aux_tables.sql`:六表 DDL(PROC_STATE / MSG_EVENT / REF_DATA / REQ_TRACK / PUMP_JOB / FLIGHT_STATE),与 ACMA-8 v4 数据模型节逐字一致;影子实例在 - 独立 schema 执行同一脚本。 + 独立 schema 执行同一脚本。**ACM2-11 拆分**:21 类静态主数据迁独立 PostgreSQL 参考库 + (`datasources.reference`,见 design.md §2);`REF_DATA` 留守 SCHD_GEN 行(流程 4 gen CAS)。 - `lua/snapshot_replace.lua`:同一 hash 原子“覆盖新代 + 按代差删”(流程 4,I4/I5)。 - `lua/batch_delete.lua`:3:30 清场批量删除(仅 ES 写成功集,I4)。 - `application.yml`:口令全部环境变量外置(零入库);`msgx.phase` 为阶段 A/B 总开关; @@ -42,6 +43,8 @@ 「Micronaut×现网 Eureka 互操作冒烟 + logstash + ES REST」通过后固化。 4. **仓储实装**:`infra/persistence/Repositories.kt` 目前是接口(Micronaut Data JDBC 实装属阶段 1 后续),主泵/调度循环以接口驱动,纯逻辑已抽离可单测。 +5. **21 类静态主数据参考库**(ACM2-11):独立 PostgreSQL(`datasources.reference` / + `StaticRefRepository`)已占位(enabled=false),表结构/迁移/同步实装属阶段 6。 ## 数据库初始化(U04/R01,务必先读) @@ -55,6 +58,10 @@ REQ_TRACK / PUMP_JOB / FLIGHT_STATE),并假定 `CMINMSGS`(及其历史表 收报第一句 SQL 即报「表不存在」。影子实例在独立 schema 执行同一脚本时同样先建旧表。 - 为什么没有 CMINMSGS 的 V1 迁移:六表之外的旧 schema 归 legacy 仓库维护(冻结期), 本仓库不重复声明;若未来要求空库一键初始化,再补 V1 基线快照(见 ACM2-10 U04)。 +- 21 类静态主数据(航空公司/航线/机位等)走**独立 PostgreSQL 参考库**(ACM2-11): + `application.yml` 的 `datasources.reference` 默认 `enabled=false`,阶段 6 实装时置 true + 并设 `MSGX_REF_DB_URL`(参考库迁移目录 `db/ref-migration`);与业务 MySQL 库(六表+legacy) + 和 Redis 航班动态互不干扰。 ## 构建 diff --git a/build.gradle.kts b/build.gradle.kts index ebfdd43..b562270 100644 --- a/build.gradle.kts +++ b/build.gradle.kts @@ -48,6 +48,8 @@ dependencies { implementation(libs.logstash.logback.encoder) runtimeOnly(libs.mysql.connector.j) + // ACM2-11:datasources.reference(21 类静态主数据独立 PG 库)驱动 org.postgresql:postgresql + // 随阶段 6 实装加入(当前 reference enabled=false 不加载;版本由 micronaut-platform BOM 钉 ~42.7) runtimeOnly(libs.snakeyaml) // U02:显式版本入 catalog(原无版本号依赖 BOM 覆盖,snakeyaml 不在 micronaut BOM 内) testImplementation(libs.junit.jupiter) diff --git a/docs/architecture.md b/docs/architecture.md index 84d0bf7..0fab526 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -24,7 +24,7 @@ |---|---|---| | 语言/运行时 | Kotlin 2.3 + JDK 25 | JDK 21 不可行(Micronaut 5.1 系要求 JVM 25+,ACMA-9 实测) | | 框架 | Micronaut 5.1.3 | 编译期 DI:KSP(`kotlin-ksp` + `micronaut-inject-kotlin`)生成 `*$Definition` | -| 持久化 | MySQL + Flyway | 仓储现为接口(Micronaut Data JDBC 实装属 U05,阶段 1 后续) | +| 持久化 | 业务库:MySQL + Flyway;静态参考库:**PostgreSQL**(ACM2-11) | 仓储现为接口(Micronaut Data JDBC 实装属 U05,阶段 1 后续;`datasources.reference` enabled=false 待阶段 6) | | 权威存储 | Redis(阶段 A) | flightInfo hash;仅主泵线程写(I5);Lua 脚本原子覆盖 | | 投递 | Kafka(acks=all + 幂等) | outbox 模式,经 MSG_EVENT 表中转 | | 投影(阶段 B) | Elasticsearch + Redis 投影 + FLIGHT_STATE | 仅阶段 B 启用(`msgx.phase`) | @@ -104,6 +104,10 @@ - **本仓库 Flyway 只建六张辅助表**:`PROC_STATE` / `MSG_EVENT` / `REF_DATA` / `REQ_TRACK` / `PUMP_JOB` / `FLIGHT_STATE`(`V2.0.0__aux_tables.sql`)。 +- **21 类静态主数据 → 独立 PostgreSQL 参考库**(ACM2-11 定案;`datasources.reference` / + `StaticRefRepository`):航空公司/航线/机位/登机桥等静态数据迁出 `REF_DATA`; + **`REF_DATA` 仅留守 SCHD_GEN 行**(流程 4 gen CAS 与 SUCCEEDED 同事务,见 design.md §2 注); + 航班动态权威 Redis(阶段 A flightInfo)**不变**。参考库表结构/迁移/21 类同步实装属阶段 6。 - `CMINMSGS` / `CMINMSGS_HST` / `COUTMSGS` 等 legacy 旧表归 legacy 仓库维护(冻结期), 本仓库不重复声明;**全新空库需先建 legacy schema**,否则收报首句 SQL 报表不存在 (README「数据库初始化」节)。 diff --git a/docs/design.md b/docs/design.md index ea29b53..53c47f1 100644 --- a/docs/design.md +++ b/docs/design.md @@ -46,7 +46,7 @@ ErrorClass(两侧共用): schdPush / outboundIntents / refUpserts),不触碰 Redis/Kafka——副作用全部由泵边界执行。 - Handler 实装:0/32(骨架),翻译属阶段 2/3,逐条对照 ACM2-4 行为基线与 KEEP/FIX 矩阵。 -## 2. 数据模型(六辅助表) +## 2. 数据模型(六辅助表 + 静态参考库) `db/migration/V2.0.0__aux_tables.sql`(legacy 旧表不在本仓库声明,见 architecture.md §6): @@ -54,7 +54,7 @@ ErrorClass(两侧共用): |---|---|---| | PROC_STATE | 消息处理伴生状态(不动 CMINMSGS 旧列) | `UK_PROC_IDENTITY(IDENTITY_KEY)` 唯一约束=I3 依据;`IDX_PROC_HEAD(STATE, CMINMSGS_ID)`=队头查询 | | MSG_EVENT | 统一投递 outbox | `EVENT_ID` 自增=全序;`IDX_EVT_HEAD(TARGET, STATE, EVENT_ID)`=每 target 队头 | -| REF_DATA | 21 类参考数据 + SCHD_GEN | `VERSION` 列支撑流程 4 CAS;`SOURCE` 区分 ADMINAPI/AODB/PIPELINE | +| REF_DATA | **SCHD_GEN(流程 4 gen 协议,留守业务库)**;原 21 类静态数据已拆分迁独立 PG 参考库(ACM2-11) | `VERSION` 列支撑流程 4 CAS(gen 行) | | REQ_TRACK | 15 类请求状态机 | REGISTERED/SENT/WAITING/DONE/EXPIRED | | PUMP_JOB | 泵作业队列 | kind:ARCHIVE/HISTORY_SWEEP/PROJECTION_REBUILD | | FLIGHT_STATE | 阶段 B 权威 | `replaceDay` 单事务删差集+写新代+版本提升 | @@ -63,6 +63,17 @@ ErrorClass(两侧共用): 主键应含 FDAY;时间列 TIMESTAMP(秒级+会话时区)应 DATETIME(6)/显式 UTC,否则退避/毒丸 判定存在系统性偏移风险。 +**库边界(ACM2-11,2026-09-07 定案)**: +- 业务事务库(MySQL `cdairport`):PROC_STATE / MSG_EVENT / REQ_TRACK / PUMP_JOB / + FLIGHT_STATE + `REF_DATA` 的 **SCHD_GEN 行**(gen CAS 与 SUCCEEDED 同事务,快照协议依赖,不能随迁); +- **21 类静态主数据 → 独立 PostgreSQL 参考库**(`datasources.reference`,`StaticRefRepository`, + SOURCE=ADMINAPI/AODB/PIPELINE):写少读多、外部来源、弱事务;Redis 只作只读热点投影 + (legacy orms_stand 语义延续,阶段 6); +- 航班动态权威 Redis(阶段 A flightInfo)**不变**。 +- 实现状态:Repository 已拆(`RefDataRepository`=gen-only / `StaticRefRepository`=21 类); + reference datasource 已占位(enabled=false);表结构草案、迁移 SQL、ReferenceService 21 类 + 同步与 codec 应答接线属阶段 6(U05 数据层批次后)。 + ## 3. 核心流程设计 ### 3.1 流程 1:收报(`InboxService.accept`) @@ -173,9 +184,10 @@ CAS 冲突 → FAILED(INFRA)+退避(串行泵下不应发生→告警语义) | `identity.include-day-boundary` | false | 幂等键日边界(CONFIRM 前禁开) | | `consistency-check.on-startup` / `daily-sample-ratio` | true / 0.01 | 一致性哨兵(实装属 U25) | -基础设施键位口径(Micronaut 5.1,U03):`datasources.default.*`、 -`flyway.datasources.default.*`、`kafka.producers.default.*`、`eureka.client.*`; -logback 独立于本文件,环境变量前缀 `MSGX_LOGSTASH_*`。 +基础设施键位口径(Micronaut 5.1,U03):`datasources.default.*`(业务 MySQL)、 +`datasources.reference.*`(21 类静态 PG,ACM2-11,enabled=false 待阶段 6)、 +`flyway.datasources.default.*` / `flyway.datasources.reference.*`、`kafka.producers.default.*`、 +`eureka.client.*`;logback 独立于本文件,环境变量前缀 `MSGX_LOGSTASH_*`。 ## 7. 测试策略 @@ -215,3 +227,4 @@ logback 独立于本文件,环境变量前缀 `MSGX_LOGSTASH_*`。 | U22–U24 | 载荷类型收敛、eventSeq 未接线、每报文全量读语义定案 | WP3 | | U25/U28/U30 | 一致性哨兵实装、README 安全节/入口、索引与杂项 | WP3/4 | | U27 | 专有材料(SIS md 703KB / XSD 版权头)治理决策 | WP4(ACL 核验先行) | +| ACM2-11 | 21 类静态主数据独立 PG 参考库(datasources.reference):Repository 接口已拆(RefDataRepository=gen-only / StaticRefRepository),表结构/迁移/SOURCE 审计/ReferenceService 同步与应答接线未实装 | 阶段 6(U05 批次后);生产 Oracle 仅可能性,触发条件见 ACM2-11 | diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index 3070eac..75d76b4 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -15,6 +15,8 @@ micronaut = "5.1.3" ksp = "2.3.0" junit = "5.11.4" snakeyaml = "2.4" +# ACM2-11:21 类静态主数据独立库选型 PostgreSQL(datasources.reference,阶段 6 实装) +postgresql = "42.7.4" [plugins] kotlin-jvm = { id = "org.jetbrains.kotlin.jvm", version.ref = "kotlin" } @@ -37,6 +39,7 @@ jackson-dataformat-xml = { module = "com.fasterxml.jackson.dataformat:jackson-da jackson-module-kotlin = { module = "com.fasterxml.jackson.module:jackson-module-kotlin", version = "2.18.2" } logstash-logback-encoder = { module = "net.logstash.logback:logstash-logback-encoder", version = "8.0" } mysql-connector-j = { module = "com.mysql:mysql-connector-j", version = "9.1.0" } +postgresql = { module = "org.postgresql:postgresql", version.ref = "postgresql" } snakeyaml = { module = "org.yaml:snakeyaml", version.ref = "snakeyaml" } junit-jupiter = { module = "org.junit.jupiter:junit-jupiter", version.ref = "junit" } h2 = { module = "com.h2database:h2", version = "2.3.232" } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/domain/Decision.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/domain/Decision.kt index 0736e61..a8e051e 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/domain/Decision.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/domain/Decision.kt @@ -9,7 +9,7 @@ data class Decision( val msgNotifies: List = emptyList(), // → MSG_EVENT(KAFKA:msg) val schdPush: List = emptyList(), // → MSG_EVENT(KAFKA:schd),PARTITION_KEY=FLID val outboundIntents: List = emptyList(), // → COUTMSGS(沿用既有列语义) - val refUpserts: List = emptyList(), // → REF_DATA + val refUpserts: List = emptyList(), // → 静态主数据(独立 PG reference 库,ACM2-11) ) /** 航班状态变更(阶段 A 由主泵线程 redisApply;阶段 B 落 FLIGHT_STATE 同事务)。 */ diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/persistence/Repositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/persistence/Repositories.kt index fd27d08..21cc2a3 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/persistence/Repositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/persistence/Repositories.kt @@ -5,6 +5,7 @@ import com.gzzn.omms.msgexchange.nextgen.domain.EventStatus import com.gzzn.omms.msgexchange.nextgen.domain.MsgEvent import com.gzzn.omms.msgexchange.nextgen.domain.ProcState import com.gzzn.omms.msgexchange.nextgen.domain.ProcStatus +import com.gzzn.omms.msgexchange.nextgen.domain.RefUpsert import java.time.Instant /** @@ -58,6 +59,11 @@ interface MsgEventRepository { fun insertSync(events: List) } +/** + * 快照 generation(SCHD_GEN)协议——只留 gen(ACM2-11 拆分后)。 + * 位置:业务事务库(MySQL,与 PROC_STATE/SUCCEEDED 同事务,流程 4 重放幂等依据); + * 21 类静态主数据已迁独立 PostgreSQL 参考库(见 [StaticRefRepository],阶段 6 实装)。 + */ interface RefDataRepository { data class GenMeta(val flids: List, val version: Long) @@ -65,8 +71,18 @@ interface RefDataRepository { /** 流程 4:版本 CAS(expected 未变才写,重放 no-op,不二次自增)。 */ fun putGenIfVersion(day: String, expected: Long, new: GenMeta): Boolean +} - fun upsertAll(rows: List>) // (rtype, rkey, payloadJson) +/** + * 21 类静态主数据(航空公司/航线/机位/登机桥等)——独立 PostgreSQL 参考库 + * (`datasources.reference`,ACM2-11 定案;SOURCE=ADMINAPI/AODB/PIPELINE,N19 对齐)。 + * 与主处理链路弱事务耦合:写入者为 ReferenceService(21 类同步)与请求应答路径, + * 不参与主泵事务 2;Redis 只作只读热点投影(legacy orms_stand 语义延续,阶段 6)。 + */ +interface StaticRefRepository { + fun upsertAll(refs: List) + + fun findByType(type: String): List } interface ReqTrackRepository { diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/stub/StubRepositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/stub/StubRepositories.kt index 8043328..89da81f 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/stub/StubRepositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/stub/StubRepositories.kt @@ -5,6 +5,7 @@ import com.gzzn.omms.msgexchange.nextgen.domain.EventStatus import com.gzzn.omms.msgexchange.nextgen.domain.MsgEvent import com.gzzn.omms.msgexchange.nextgen.domain.ProcState import com.gzzn.omms.msgexchange.nextgen.domain.ProcStatus +import com.gzzn.omms.msgexchange.nextgen.domain.RefUpsert import com.gzzn.omms.msgexchange.nextgen.infra.persistence.CminmsgInboxRepository import com.gzzn.omms.msgexchange.nextgen.infra.persistence.FlightStateRepository import com.gzzn.omms.msgexchange.nextgen.infra.persistence.MsgEventRepository @@ -12,6 +13,7 @@ import com.gzzn.omms.msgexchange.nextgen.infra.persistence.ProcStateRepository import com.gzzn.omms.msgexchange.nextgen.infra.persistence.PumpJobRepository import com.gzzn.omms.msgexchange.nextgen.infra.persistence.RefDataRepository import com.gzzn.omms.msgexchange.nextgen.infra.persistence.ReqTrackRepository +import com.gzzn.omms.msgexchange.nextgen.infra.persistence.StaticRefRepository import io.micronaut.context.annotation.Requires import jakarta.inject.Singleton import java.time.Instant @@ -166,13 +168,13 @@ class StubPumpJobs : PumpJobRepository { override fun markFailed(jobId: Long, lastError: String) { rows[jobId]?.state = "FAILED"; rows[jobId]?.lastError = lastError } } +/** 快照 gen 协议 stub(业务库 REF_DATA 的 SCHD_GEN 行;ACM2-11 拆分后 21 类不在本实现)。 */ @Requires(property = "msgx.stubs", value = "true") @Singleton class StubRefData : RefDataRepository { private val gens = mutableMapOf() - private val flat = mutableMapOf, String>() - fun clear() { gens.clear(); flat.clear() } + fun clear() { gens.clear() } override fun getGen(day: String): RefDataRepository.GenMeta? = gens[day] @@ -182,10 +184,22 @@ class StubRefData : RefDataRepository { gens[day] = new return true } +} - override fun upsertAll(rows: List>) { - rows.forEach { flat[it.first to it.second] = it.third } +/** 21 类静态主数据 stub(独立 PG reference 库;内存实现,source 保留供审计断言)。 */ +@Requires(property = "msgx.stubs", value = "true") +@Singleton +class StubStaticRef : StaticRefRepository { + private val rows = mutableMapOf, RefUpsert>() + + fun clear() { rows.clear() } + + override fun upsertAll(refs: List) { + refs.forEach { rows[it.rtype to it.rkey] = it } } + + override fun findByType(type: String): List = + rows.values.filter { it.rtype == type } } @Requires(property = "msgx.stubs", value = "true") diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/processing/Pump.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/processing/Pump.kt index 9a492ee..b38e044 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/processing/Pump.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/processing/Pump.kt @@ -198,7 +198,8 @@ class MessageProcessor( if (props.phase == PipelineProps.Phase.A) { // 阶段 A:主泵线程先写 Redis(幂等),先于事件创建(I2);TODO: redisApply(flightChanges) } - // 事务 2:事件 + 出站意图 + REF_DATA + CMINMSGS 回填 + SUCCEEDED(实装后 @Transactional) + // 事务 2:事件 + CMINMSGS 回填 + SUCCEEDED(实装后 @Transactional); + // 静态主数据(refUpserts)落独立 PG reference 库,弱事务不入本事务(ACM2-11) 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)) } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/reference/ReferenceService.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/reference/ReferenceService.kt index fba3bb7..1f3f76d 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/reference/ReferenceService.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/reference/ReferenceService.kt @@ -1,18 +1,19 @@ package com.gzzn.omms.msgexchange.nextgen.reference -import com.gzzn.omms.msgexchange.nextgen.infra.persistence.RefDataRepository +import com.gzzn.omms.msgexchange.nextgen.infra.persistence.StaticRefRepository import jakarta.inject.Singleton /** * ACMA-8 决策 6:21 类 admin-api 基础数据同步(阶段 6)。 - * 只产 REF_DATA 行(权威写唯一入口);影子实例不运行本模块(v4 影子二分)。 + * 落点 = 21 类静态主数据独立 PostgreSQL 参考库(StaticRefRepository,ACM2-11); + * 只产静态主数据行(权威写唯一入口);影子实例不运行本模块(v4 影子二分)。 */ @Singleton class ReferenceService( - private val refData: RefDataRepository, + private val staticRef: StaticRefRepository, // TODO(阶段6): @Client(id="ADMINAPI") 声明式客户端 + 21 类端点配置化清单 ) { fun refresh(type: String) { - // TODO(阶段6): 拉取 → upsertAll(type, key, payload);失败旧数据可用(REF_DATA 保留旧值) + // TODO(阶段6): 拉取 → staticRef.upsertAll(rows, source=ADMINAPI);失败旧数据可用(参考库保留旧值) } } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/reference/RequestCoordinator.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/reference/RequestCoordinator.kt index dea38d5..a4ef1ce 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/reference/RequestCoordinator.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/reference/RequestCoordinator.kt @@ -1,7 +1,8 @@ package com.gzzn.omms.msgexchange.nextgen.reference -import com.gzzn.omms.msgexchange.nextgen.infra.persistence.RefDataRepository +import com.gzzn.omms.msgexchange.nextgen.domain.RefUpsert import com.gzzn.omms.msgexchange.nextgen.infra.persistence.ReqTrackRepository +import com.gzzn.omms.msgexchange.nextgen.infra.persistence.StaticRefRepository import jakarta.inject.Singleton import java.time.Instant @@ -9,11 +10,13 @@ import java.time.Instant * ACMA-8 流程 6:请求式查询状态机(15 类,决策树版 v4)。 * 同类并发=1(注册新请求先强制旧请求 EXPIRED,终态不再被匹配); * 应答匹配:回显字段 CONFIRM(矩阵 #13)→ 精确匹配;否则退化模式(DTTM ≥ SENT 才应用 + 审计)。 + * 应答落库目标 = 21 类静态主数据(独立 PG 参考库,StaticRefRepository,ACM2-11); + * 注:N04 量纲缺陷(dttm vs epochMillis)与 N16 死分支未修,接线阶段 6 前处理(U19–U21)。 */ @Singleton class RequestCoordinator( private val reqTrack: ReqTrackRepository, - private val refData: RefDataRepository, + private val staticRef: StaticRefRepository, ) { fun request(kind: String, rangeJson: String): Long { reqTrack.forceExpireOpenOf(kind) // 防护① @@ -29,7 +32,7 @@ class RequestCoordinator( kind: String, dttm: Long, echoSeqn: Long?, - records: List>, // (rtype, rkey, payloadJson) + records: List, // 应答载荷 → 静态主数据(source=AODB) ): String { val req = reqTrack.findOpenByKind(kind) ?: return audit("late/unknown resp kind=$kind") val echoConfirmed = false // CONFIRM(矩阵 #13)待机场方结论 @@ -44,8 +47,8 @@ class RequestCoordinator( } } - private fun applyResp(reqId: Long, records: List>) { - refData.upsertAll(records) // 大应答复用流程 4 staging 路径(阶段 6) + private fun applyResp(reqId: Long, records: List) { + staticRef.upsertAll(records) // 独立 PG 参考库(阶段 6 接线) reqTrack.markDone(reqId) } diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index 818359b..b3b4989 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -38,6 +38,15 @@ datasources: username: ${MSGX_DB_USER} password: ${MSGX_DB_PASSWORD} driver-class-name: com.mysql.cj.jdbc.Driver + # ACM2-11:21 类静态主数据独立 PostgreSQL 参考库(reference datasource)。 + # 数据层实装(阶段 6 / U05 批次)前默认 disabled——不建连不解析占位符; + # 接入时置 enabled=true 并设 MSGX_REF_DB_URL(或 host/port/name)等环境变量。 + reference: + enabled: false + url: ${MSGX_REF_DB_URL} + username: ${MSGX_REF_DB_USER:reference} + password: ${MSGX_REF_DB_PASSWORD:} + driver-class-name: org.postgresql.Driver flyway: datasources: @@ -46,6 +55,10 @@ flyway: locations: classpath:db/migration # U04(R01):V2.0.0 只建六张辅助表,假定 CMINMSGS/CMINMSGS_HST 等 legacy 旧表已存在; # 全新库请按 README「数据库初始化」先落 legacy schema(或 baseline),否则收报首句 SQL 会报表不存在。 + reference: + # ACM2-11:21 类静态主数据 PG 库迁移(类表 DDL 草案见 ACM2-11 Checks,阶段 6 落库) + enabled: false + locations: classpath:db/ref-migration redis: uri: ${MSGX_REDIS_URI} # 阶段 A 权威存储;仅主泵线程写(I5) diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/processing/MessageProcessorTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/processing/MessageProcessorTest.kt index a324203..e73cc49 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/processing/MessageProcessorTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/processing/MessageProcessorTest.kt @@ -144,7 +144,6 @@ class MessageProcessorTest { private class FakeRefData : RefDataRepository { override fun getGen(day: String): RefDataRepository.GenMeta? = null override fun putGenIfVersion(day: String, expected: Long, new: RefDataRepository.GenMeta): Boolean = true - override fun upsertAll(rows: List>) = Unit } // ---------- helpers ----------