From d298fcc48f244847643ae00ee4653e412e803bc4 Mon Sep 17 00:00:00 2001 From: windyboy Date: Mon, 21 Sep 2026 14:11:53 +0800 Subject: [PATCH] =?UTF-8?q?fix(delivery):=20claimBatch=20=E6=8E=92?= =?UTF-8?q?=E9=99=A4=E5=B7=B2=E6=9A=82=E5=81=9C=20FLID=EF=BC=8C=E9=81=BF?= =?UTF-8?q?=E5=85=8D=E6=8A=95=E9=80=92=E9=A5=BF=E6=AD=BB=EF=BC=88ACM2-95?= =?UTF-8?q?=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 同 FLID 事件占满 delivery-batch 时后续轮次不再重复领到同一批;补饥饿回归用例。 Co-authored-by: Cursor --- .../omms/msgexchange/delivery/Dispatcher.kt | 3 +- .../infra/persistence/Repositories.kt | 6 +++- .../persistence/jdbc/JdbcPgRepositories.kt | 25 ++++++++++++++--- .../infra/stub/StubRepositories.kt | 8 ++++-- .../delivery/DispatcherTickTest.kt | 28 +++++++++++++++++++ 5 files changed, 62 insertions(+), 8 deletions(-) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/Dispatcher.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/Dispatcher.kt index 712724a..239c66d 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/Dispatcher.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/Dispatcher.kt @@ -104,8 +104,9 @@ class Dispatcher( val paused = HashSet() // 本轮投递中被暂停的 FLID:队头失败或退避未到期 while (round < maxRounds) { round++ + // 领批时排除已暂停 FLID,避免单航班占满 delivery-batch 饿死其他航班(ACM2-95) val batch = try { - msgEvents.claimBatch(Targets.KAFKA_MSG, batchSize) + msgEvents.claimBatch(Targets.KAFKA_MSG, batchSize, paused) } catch (e: Exception) { log.error("claimBatch(KAFKA_MSG) failed", e) return sent diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt index e5d3990..49b634d 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt @@ -172,7 +172,11 @@ data class Backlog( interface MsgEventRepository { fun insertAll(events: List): List - fun claimBatch(target: String, limit: Int): List + /** + * 领取待发事件;[excludePartitionKeys] 为本轮已暂停的 FLID,排除后避免单航班占满整批 + * 导致其他航班饿死(ACM2-95)。 + */ + fun claimBatch(target: String, limit: Int, excludePartitionKeys: Set = emptySet()): List /** 领取待发的整态事件:写入侧已保证每个 `FLID` 只有一行,读端不再去重合并。 */ fun mergePendingSchd(now: Instant, limit: Int): List diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt index 6ba2e63..19c0a74 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt @@ -512,12 +512,29 @@ class JdbcMsgEventRepository( ps.setTimestamp(9, e.createdAt.toSqlTimestamp()) } - override fun claimBatch(target: String, limit: Int): List = - ds.query( - "SELECT * FROM msg_event WHERE target = ? AND state = 'PENDING' ORDER BY event_id ASC LIMIT ?", - { ps -> ps.setString(1, target); ps.setInt(2, limit) }, + override fun claimBatch(target: String, limit: Int, excludePartitionKeys: Set): List { + if (excludePartitionKeys.isEmpty()) { + return ds.query( + "SELECT * FROM msg_event WHERE target = ? AND state = 'PENDING' ORDER BY event_id ASC LIMIT ?", + { ps -> ps.setString(1, target); ps.setInt(2, limit) }, + ::mapEvent, + ) + } + val placeholders = excludePartitionKeys.joinToString(",") { "?" } + return ds.query( + """ + SELECT * FROM msg_event + WHERE target = ? AND state = 'PENDING' AND partition_key NOT IN ($placeholders) + ORDER BY event_id ASC LIMIT ? + """.trimIndent(), + { ps -> + ps.setString(1, target) + excludePartitionKeys.forEachIndexed { i, key -> ps.setString(i + 2, key) } + ps.setInt(excludePartitionKeys.size + 2, limit) + }, ::mapEvent, ) + } /** * 领取待发的整态事件。写入侧已保证每个 `FLID` 只有一行(部分唯一索引 `uq_schd_event`), diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt index 8de5332..5820c77 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt @@ -280,8 +280,12 @@ class StubMsgEvents : MsgEventRepository { return id } - override fun claimBatch(target: String, limit: Int): List = - rows.values.filter { it.target == target && it.state == EventStatus.PENDING }.sortedBy { it.eventId!! }.take(limit) + override fun claimBatch(target: String, limit: Int, excludePartitionKeys: Set): List = + rows.values + .filter { it.target == target && it.state == EventStatus.PENDING } + .filter { it.partitionKey !in excludePartitionKeys } + .sortedBy { it.eventId!! } + .take(limit) /** 写入侧已保证每个 FLID 单行,按退避与事件先后输出。 */ override fun mergePendingSchd(now: Instant, limit: Int): List = diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/delivery/DispatcherTickTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/delivery/DispatcherTickTest.kt index e086c00..d02f67e 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/delivery/DispatcherTickTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/delivery/DispatcherTickTest.kt @@ -247,6 +247,34 @@ class DispatcherTickTest { assertTrue(f1Rows.all { it.state == EventStatus.PENDING }) } + /** 暂停 FLID 事件数 ≥ delivery-batch 时,领批排除已暂停键,其他航班同轮仍可发(ACM2-95)。 */ + @Test + fun `paused FLID filling the claim batch does not starve other flights`() { + val repo = StubMsgEvents() + val port = FailingMsgKeyPort("F1") + val props = PipelineProps().apply { + pipeline.deliveryBatch = 3 + pipeline.deliveryDrainRounds = 5 + schd.flushPeriod = java.time.Duration.ofHours(1) + } + val d = dispatcher(repo, port, props) + // F1 三条占满 batch=3;F2 一条排在后面——旧实现会整轮领到同一批 F1 并 sent=0 + repo.insertAll( + listOf( + ev(1, Targets.KAFKA_MSG, "F1", """{"flid":"F1","v":1}"""), + ev(2, Targets.KAFKA_MSG, "F1", """{"flid":"F1","v":2}"""), + ev(3, Targets.KAFKA_MSG, "F1", """{"flid":"F1","v":3}"""), + ev(4, Targets.KAFKA_MSG, "F2", """{"flid":"F2","v":1}"""), + ), + ) + + d.tick() + + assertEquals(1, port.sent.count { it.topic == "msg" && it.key == "F2" }) + assertEquals(0, port.sent.count { it.topic == "msg" && it.key == "F1" }) + assertTrue(repo.rows.values.filter { it.partitionKey == "F1" }.all { it.state == EventStatus.PENDING }) + } + @Test fun `msg head backoff not yet due pauses only the same FLID`() { val repo = StubMsgEvents()