refactor(storage): ACM2-12 落地——自有 PostgreSQL 全内部状态,共享 MySQL 仅信箱
按 ACM2-12 定案把仓库存储口径从"MySQL 六辅助表"推进到最终形态: - 迁移:删除 MySQL V2.0.0 六表脚本,新增自有 PG V1.0.0 (PROC_STATE/MSG_EVENT/PUMP_JOB/REQ_TRACK/REF_MASTER,PG 方言,REQ_TRACK.COUTMSGS_ID 按 U18 修为 BIGINT;FLIGHT_STATE 阶段 B 缓做不建表)。 - 配置:datasources.default = 自有 PostgreSQL(enabled=false 待 U05); 移除 datasources.reference/flyway.reference;新增 mailbox.shared-mysql(共享信箱, 仅 DML,不建表);test profile 显式启用 H2 内存 datasource。 - 接口/注释:Repositories KDoc 按 ACM2-12 归属(自有 PG / 信箱封装 / gen→Redis 占位 / FLIGHT_STATE 缓做);InboxService、Pump 事务模型注释对正 (信箱外部副作用 + PG 本地事务 + 回填最终一致)。 - 文档:architecture §1/§2/§3/§4/§5(D1/D2)/§6/§7/§8、design §1/§2/§3.1/§3.4/§3.5/§6/§9、 README 全部按 ACM2-12 收口(自有 PG + 共享信箱 + Redis 动态/gen + 阶段 B 缓做)。 验证:37 测试全绿;dev stub 冒烟仍可启动、/health UP、收报 200。 信箱适配层(CminmsgMailbox/OutboxMailbox)、gen→Redis Lua、作业窗口语义、影子重设计 属 ACM2-12 Checks ②③④⑤(U05/U09 批次)。
This commit is contained in:
+33
-13
@@ -9,7 +9,11 @@ import com.gzzn.omms.msgexchange.nextgen.domain.RefUpsert
|
||||
import java.time.Instant
|
||||
|
||||
/**
|
||||
* ACMA-8 v4 仓储接口(接口驱动,主泵/调度循环可单测;Micronaut Data JDBC 实装属阶段 1 后续)。
|
||||
* ACM2-12 仓储接口(接口驱动,主泵/调度循环可单测;Micronaut Data JDBC 实装属 U05 批次)。
|
||||
* 存储边界:自有 PostgreSQL(datasources.default)= 本文件除 CminmsgInboxRepository 外
|
||||
* 的全部接口(消息管道 PROC_STATE/MSG_EVENT、PUMP_JOB、REQ_TRACK、21 类 REF_MASTER);
|
||||
* 共享 MySQL 信箱(CMINMSGS 收 / COUTMSGS 出)经信箱封装访问,仅 DML、不建表;
|
||||
* Redis = 航班动态 + 快照 gen(RefDataRepository 目标实现);FLIGHT_STATE 缓做。
|
||||
*/
|
||||
interface ProcStateRepository {
|
||||
fun insert(cminmsgsId: Long, state: ProcStatus = ProcStatus.PENDING)
|
||||
@@ -60,9 +64,10 @@ interface MsgEventRepository {
|
||||
}
|
||||
|
||||
/**
|
||||
* 快照 generation(SCHD_GEN)协议——只留 gen(ACM2-11 拆分后)。
|
||||
* 位置:业务事务库(MySQL,与 PROC_STATE/SUCCEEDED 同事务,流程 4 重放幂等依据);
|
||||
* 21 类静态主数据已迁独立 PostgreSQL 参考库(见 [StaticRefRepository],阶段 6 实装)。
|
||||
* 快照 generation(SCHD_GEN)协议——只留 gen。
|
||||
* ACM2-12:gen 迁 Redis(与 flightInfo 同源,Lua 内原子「覆盖+按代差删+版本推进」,
|
||||
* DB 仅写 SUCCEEDED;重放幂等由 Lua 承接,协议重设计属 U09)。本接口为过渡占位,
|
||||
* 目标实现为 Redis gen store(script 化),非关系表。
|
||||
*/
|
||||
interface RefDataRepository {
|
||||
data class GenMeta(val flids: List<String>, val version: Long)
|
||||
@@ -74,10 +79,10 @@ interface RefDataRepository {
|
||||
}
|
||||
|
||||
/**
|
||||
* 21 类静态主数据(航空公司/航线/机位/登机桥等)——独立 PostgreSQL 参考库
|
||||
* (`datasources.reference`,ACM2-11 定案;SOURCE=ADMINAPI/AODB/PIPELINE,N19 对齐)。
|
||||
* 与主处理链路弱事务耦合:写入者为 ReferenceService(21 类同步)与请求应答路径,
|
||||
* 不参与主泵事务 2;Redis 只作只读热点投影(legacy orms_stand 语义延续,阶段 6)。
|
||||
* 21 类静态主数据(航空公司/航线/机位/登机桥等)——自有 PostgreSQL `REF_MASTER` 表
|
||||
* (ACM2-12:与消息管道同自有库;SOURCE=ADMINAPI/AODB/PIPELINE,N19 对齐)。
|
||||
* 与主链弱事务耦合:写入者为 ReferenceService(21 类同步)与请求应答路径;
|
||||
* Redis 只作只读热点投影(legacy orms_stand 语义延续,阶段 6)。
|
||||
*/
|
||||
interface StaticRefRepository {
|
||||
fun upsertAll(refs: List<RefUpsert>)
|
||||
@@ -85,6 +90,10 @@ interface StaticRefRepository {
|
||||
fun findByType(type: String): List<RefUpsert>
|
||||
}
|
||||
|
||||
/**
|
||||
* 15 类请求状态机——自有 PG `REQ_TRACK`(ACM2-12;Reference & Query 域)。
|
||||
* 与共享库 COUTMSGS(出站信箱)跨库:先 COUTMSGS 落库成功 → 再 markSent,补偿重扫,最终一致。
|
||||
*/
|
||||
interface ReqTrackRepository {
|
||||
data class Req(
|
||||
val reqId: Long,
|
||||
@@ -108,6 +117,11 @@ interface ReqTrackRepository {
|
||||
fun markDone(reqId: Long)
|
||||
}
|
||||
|
||||
/**
|
||||
* 泵作业调度记录——自有 PG `PUMP_JOB`(ACM2-12)。
|
||||
* 决策 1 修订:作业不插队,仅在消息队头空闲/退避窗口由主泵执行(跨库/异队列无全序);
|
||||
* 作业动作本身(归档写共享库 CMINMSGS_HST、清场删 Redis/ES 等)仍在各自目标存储。
|
||||
*/
|
||||
interface PumpJobRepository {
|
||||
data class Job(val jobId: Long, val kind: String) // ARCHIVE/HISTORY_SWEEP/PROJECTION_REBUILD
|
||||
|
||||
@@ -122,6 +136,10 @@ interface PumpJobRepository {
|
||||
fun markFailed(jobId: Long, lastError: String)
|
||||
}
|
||||
|
||||
/**
|
||||
* 阶段 B 权威(ACM2-12:缓做,不落表——航班动态权威保持 Redis;阶段 B 重新
|
||||
* 评估后再定是否引入事务化权威)。本接口仅供占位与测试,勿据此建表。
|
||||
*/
|
||||
interface FlightStateRepository {
|
||||
/** 阶段 B 权威;replaceDay = 单事务删差集+写新代+版本提升。 */
|
||||
fun replaceDay(day: String, flights: List<Pair<String, String>>)
|
||||
@@ -129,15 +147,17 @@ interface FlightStateRepository {
|
||||
fun findByDay(day: String): List<Pair<String, String>>
|
||||
}
|
||||
|
||||
/**
|
||||
* 收报信箱(共享 MySQL CMINMSGS,仅 DML——ACM2-12:库属他人系统,本系统不建表)。
|
||||
* 事务模型:insertRaw = 外部副作用「信箱落信」(成功即返回其主键 CMINMSGS_ID)→
|
||||
* 随后自有 PG 建 PROC_STATE(PENDING) 入队;PG 建行失败以共享库
|
||||
* DATE_PROCESSED IS NULL 重扫补建。backfillOnSuccess 为处理成功后的外部回填
|
||||
* (DATE_PROCESSED/STATUS,最终一致、失败重试+告警)。
|
||||
*/
|
||||
interface CminmsgInboxRepository {
|
||||
/** 事务 1:原文落库(沿用 SUBSYSTEM_* 列;仅 CLOB/DATE_RECEIVED,I3:无接收唯一约束)。 */
|
||||
fun insertRaw(rawXml: String): Long
|
||||
|
||||
fun rawOf(cminmsgsId: Long): String?
|
||||
|
||||
/**
|
||||
* v4 回滚兼容关键:SUCCEEDED 时回填 SUBSYSTEM_*(FIX)+ DATE_PROCESSED/STATUS,
|
||||
* 旧系统按自身 processed 语义可无缝接管(Runbook 第 7 步)。
|
||||
*/
|
||||
fun backfillOnSuccess(cminmsgsId: Long, sndr: String, type: String, styp: String, seqn: Long)
|
||||
}
|
||||
|
||||
@@ -5,23 +5,24 @@ import com.gzzn.omms.msgexchange.nextgen.infra.persistence.ProcStateRepository
|
||||
import jakarta.inject.Singleton
|
||||
import java.time.Instant
|
||||
|
||||
/** ACMA-8 流程 1:入站接收(事务 1)。响应语义 = “已持久化”(与现役一致);不解析报文(I3)。 */
|
||||
/** ACMA-8 流程 1:入站接收。响应语义 = “已持久化”(与现役一致);不解析报文(I3)。 */
|
||||
@Singleton
|
||||
class InboxService(
|
||||
private val inbox: CminmsgInboxRepository,
|
||||
private val procState: ProcStateRepository,
|
||||
// TODO(阶段1后续): 事务边界(@Transactional)随 Micronaut Data 实装补齐;
|
||||
// pump 唤醒仅加速,崩溃后主泵 1s 轮询兜底。
|
||||
// TODO(U05 批次): 事务模型(ACM2-12)——insertRaw 为共享信箱外部写(成功即返回 ID),
|
||||
// 随后自有 PG 建 PROC_STATE(PENDING);两写跨库,PG 建行失败以共享库
|
||||
// DATE_PROCESSED IS NULL 重扫补建;pump 唤醒仅加速,崩溃后主泵 1s 轮询兜底。
|
||||
) {
|
||||
private val log = org.slf4j.LoggerFactory.getLogger(InboxService::class.java)
|
||||
|
||||
data class Receipt(val cminmsgsId: Long, val receivedAt: Instant)
|
||||
|
||||
fun accept(rawXml: String): Receipt {
|
||||
val id = inbox.insertRaw(rawXml)
|
||||
procState.insert(id) // 同事务(实装后);接收层无唯一约束(I3)
|
||||
val id = inbox.insertRaw(rawXml) // 信箱外部写(ACM2-12,与 PG 跨库)
|
||||
procState.insert(id) // 自有 PG 入队;接收层无唯一约束(I3)
|
||||
wakePump()
|
||||
log.info("accepted cminmsgsId={} (tx1)", id)
|
||||
log.info("accepted cminmsgsId={}", id)
|
||||
return Receipt(id, Instant.now())
|
||||
}
|
||||
|
||||
|
||||
@@ -198,8 +198,9 @@ class MessageProcessor(
|
||||
if (props.phase == PipelineProps.Phase.A) {
|
||||
// 阶段 A:主泵线程先写 Redis(幂等),先于事件创建(I2);TODO: redisApply(flightChanges)
|
||||
}
|
||||
// 事务 2:事件 + CMINMSGS 回填 + SUCCEEDED(实装后 @Transactional);
|
||||
// 静态主数据(refUpserts)落独立 PG reference 库,弱事务不入本事务(ACM2-11)
|
||||
// 事务 2(自有 PG 内原子):事件 + SUCCEEDED(实装后 @Transactional);
|
||||
// CMINMSGS 回填 = 共享信箱外部副作用(最终一致,ACM2-12);
|
||||
// 静态主数据(refUpserts)同自有 PG 但弱事务独立提交(21 类 REF_MASTER)。
|
||||
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)) }
|
||||
|
||||
Reference in New Issue
Block a user