fix(delivery): claimBatch 排除已暂停 FLID,避免投递饿死(ACM2-95)
同 FLID 事件占满 delivery-batch 时后续轮次不再重复领到同一批;补饥饿回归用例。 Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -104,8 +104,9 @@ class Dispatcher(
|
|||||||
val paused = HashSet<String>() // 本轮投递中被暂停的 FLID:队头失败或退避未到期
|
val paused = HashSet<String>() // 本轮投递中被暂停的 FLID:队头失败或退避未到期
|
||||||
while (round < maxRounds) {
|
while (round < maxRounds) {
|
||||||
round++
|
round++
|
||||||
|
// 领批时排除已暂停 FLID,避免单航班占满 delivery-batch 饿死其他航班(ACM2-95)
|
||||||
val batch = try {
|
val batch = try {
|
||||||
msgEvents.claimBatch(Targets.KAFKA_MSG, batchSize)
|
msgEvents.claimBatch(Targets.KAFKA_MSG, batchSize, paused)
|
||||||
} catch (e: Exception) {
|
} catch (e: Exception) {
|
||||||
log.error("claimBatch(KAFKA_MSG) failed", e)
|
log.error("claimBatch(KAFKA_MSG) failed", e)
|
||||||
return sent
|
return sent
|
||||||
|
|||||||
@@ -172,7 +172,11 @@ data class Backlog(
|
|||||||
interface MsgEventRepository {
|
interface MsgEventRepository {
|
||||||
fun insertAll(events: List<MsgEvent>): List<Long>
|
fun insertAll(events: List<MsgEvent>): List<Long>
|
||||||
|
|
||||||
fun claimBatch(target: String, limit: Int): List<MsgEvent>
|
/**
|
||||||
|
* 领取待发事件;[excludePartitionKeys] 为本轮已暂停的 FLID,排除后避免单航班占满整批
|
||||||
|
* 导致其他航班饿死(ACM2-95)。
|
||||||
|
*/
|
||||||
|
fun claimBatch(target: String, limit: Int, excludePartitionKeys: Set<String> = emptySet()): List<MsgEvent>
|
||||||
|
|
||||||
/** 领取待发的整态事件:写入侧已保证每个 `FLID` 只有一行,读端不再去重合并。 */
|
/** 领取待发的整态事件:写入侧已保证每个 `FLID` 只有一行,读端不再去重合并。 */
|
||||||
fun mergePendingSchd(now: Instant, limit: Int): List<MsgEvent>
|
fun mergePendingSchd(now: Instant, limit: Int): List<MsgEvent>
|
||||||
|
|||||||
+21
-4
@@ -512,12 +512,29 @@ class JdbcMsgEventRepository(
|
|||||||
ps.setTimestamp(9, e.createdAt.toSqlTimestamp())
|
ps.setTimestamp(9, e.createdAt.toSqlTimestamp())
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun claimBatch(target: String, limit: Int): List<MsgEvent> =
|
override fun claimBatch(target: String, limit: Int, excludePartitionKeys: Set<String>): List<MsgEvent> {
|
||||||
ds.query(
|
if (excludePartitionKeys.isEmpty()) {
|
||||||
"SELECT * FROM msg_event WHERE target = ? AND state = 'PENDING' ORDER BY event_id ASC LIMIT ?",
|
return ds.query(
|
||||||
{ ps -> ps.setString(1, target); ps.setInt(2, limit) },
|
"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,
|
::mapEvent,
|
||||||
)
|
)
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 领取待发的整态事件。写入侧已保证每个 `FLID` 只有一行(部分唯一索引 `uq_schd_event`),
|
* 领取待发的整态事件。写入侧已保证每个 `FLID` 只有一行(部分唯一索引 `uq_schd_event`),
|
||||||
|
|||||||
@@ -280,8 +280,12 @@ class StubMsgEvents : MsgEventRepository {
|
|||||||
return id
|
return id
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun claimBatch(target: String, limit: Int): List<MsgEvent> =
|
override fun claimBatch(target: String, limit: Int, excludePartitionKeys: Set<String>): List<MsgEvent> =
|
||||||
rows.values.filter { it.target == target && it.state == EventStatus.PENDING }.sortedBy { it.eventId!! }.take(limit)
|
rows.values
|
||||||
|
.filter { it.target == target && it.state == EventStatus.PENDING }
|
||||||
|
.filter { it.partitionKey !in excludePartitionKeys }
|
||||||
|
.sortedBy { it.eventId!! }
|
||||||
|
.take(limit)
|
||||||
|
|
||||||
/** 写入侧已保证每个 FLID 单行,按退避与事件先后输出。 */
|
/** 写入侧已保证每个 FLID 单行,按退避与事件先后输出。 */
|
||||||
override fun mergePendingSchd(now: Instant, limit: Int): List<MsgEvent> =
|
override fun mergePendingSchd(now: Instant, limit: Int): List<MsgEvent> =
|
||||||
|
|||||||
@@ -247,6 +247,34 @@ class DispatcherTickTest {
|
|||||||
assertTrue(f1Rows.all { it.state == EventStatus.PENDING })
|
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
|
@Test
|
||||||
fun `msg head backoff not yet due pauses only the same FLID`() {
|
fun `msg head backoff not yet due pauses only the same FLID`() {
|
||||||
val repo = StubMsgEvents()
|
val repo = StubMsgEvents()
|
||||||
|
|||||||
Reference in New Issue
Block a user