fix(acm2-96): schd 单条上限超限分轮,文档与实现对齐 C-9

This commit is contained in:
windyboy
2026-09-23 16:59:09 +08:00
parent 4fbb974307
commit 71698b240b
8 changed files with 41 additions and 11 deletions
@@ -23,7 +23,7 @@ interface DeliveryPort {
*/
fun sendKafkaSchd(topic: String, payloadJson: String)
/** 发一条删除通知:key 是 FLID、value 为空;下游按"整态里这个键没了"理解成删除。 */
/** 遗留的空 value 删除通知(tombstone);当前删除只走 `KAFKA:msg` 的 UPSERT`C-9`),生产路径不再调用。 */
fun sendKafkaNull(topic: String, key: String)
/** 给健康检查用的连通性探测;默认返回 true,真实 Kafka 实现要覆写成向 broker 拉一次 metadata 来判断。 */
@@ -156,8 +156,8 @@ class Dispatcher(
}
/**
* 发 `KAFKA:schd`:本 tick 窗口内待发 UPSERT 聚成一条 `SCHD.FLTR` JSON 数组、不设 key`C-9`
* 空窗口不发。schd 侧 TOMBSTONE 不发送(删航班只走 msg)。
* 发 `KAFKA:schd`:本待发 UPSERT 取上限(`PARAM:msgx.schd.flush-limit`聚成一条 `SCHD.FLTR` JSON 数组、不设 key`C-9`
* 超限剩余行留待后续轮次发出。空窗口不发。schd 侧 TOMBSTONE 不发送(删航班只走 msg)。
*/
internal fun flushSchd() {
val batch = try {
@@ -2,7 +2,7 @@ package com.gzzn.omms.msgexchange.domain
/**
* 事件形态。生产路径只写 [UPSERT]:删除通知也走 `KAFKA:msg` 的 UPSERT 变更(`C-9`)。
* [TOMBSTONE] 仍保留在枚举与投递分支里,当前无生产写入方(待 `Q5` 定稿清理或保留)
* [TOMBSTONE] 仅枚举与写入守卫保留兼容语义,当前无生产写入方、不发送
*/
enum class EventType { UPSERT, TOMBSTONE }
@@ -2,7 +2,8 @@ package com.gzzn.omms.msgexchange.domain
/**
* 事件往哪儿投。KAFKA_SCHD 按 FLID(航班实例 ID)发航班的完整状态,KAFKA_MSG 只发"变化了"
* 的通知;两个主题之间不保证顺序,删除用 value 为空的 tombstone 消息表示。
* 的通知;两个主题之间不保证顺序。删航班只走 `KAFKA:msg` 的 UPSERT 变更(`C-9`),
* 不发空 value 的 tombstone。
*
* 常量值就是 MSG_EVENT.TARGET 落库的值,真正的 topic 名由投递适配层映射
* KAFKA:msg → "msg"KAFKA:schd → "schd")。ES 投影还没做,先不登记目标。
@@ -95,6 +95,35 @@ class DispatcherTickTest {
assertEquals(0, repo.rows.values.count { it.state.name == "PENDING" })
}
@Test
fun `schd overflow beyond flush-limit is sent in follow-up rounds`() {
val repo = StubMsgEvents()
val port = StubDeliveryPort()
val props = PipelineProps().apply { schd.flushLimit = 2 }
repo.insertAll(
listOf(
ev(1, Targets.KAFKA_SCHD, "F1", """{"v":"f1"}""", 1),
ev(2, Targets.KAFKA_SCHD, "F2", """{"v":"f2"}""", 1),
ev(3, Targets.KAFKA_SCHD, "F3", """{"v":"f3"}""", 1),
),
)
val d = dispatcher(repo, port, props)
d.flushSchd()
val first = port.sent.filter { it.topic == "schd" }
assertEquals(1, first.size)
assertNull(first.single().key)
assertTrue(first.single().payload!!.contains("f1"))
assertTrue(first.single().payload!!.contains("f2"))
assertEquals(1, repo.rows.values.count { it.state == EventStatus.PENDING })
d.flushSchd()
val all = port.sent.filter { it.topic == "schd" }
assertEquals(2, all.size)
assertTrue(all.last().payload!!.contains("f3"))
assertEquals(0, repo.rows.values.count { it.state == EventStatus.PENDING })
}
@Test
fun `schd retry is not selected before its backoff expires`() {
val repo = StubMsgEvents()