2026-09-07 15:11:33 +08:00
|
|
|
|
package com.gzzn.omms.msgexchange.delivery
|
2026-09-06 16:13:50 +08:00
|
|
|
|
|
2026-09-07 15:11:33 +08:00
|
|
|
|
import com.gzzn.omms.msgexchange.config.PipelineProps
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.domain.ErrorClass
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.domain.EventStatus
|
2026-09-09 17:53:08 +08:00
|
|
|
|
import com.gzzn.omms.msgexchange.domain.EventType
|
2026-09-07 15:11:33 +08:00
|
|
|
|
import com.gzzn.omms.msgexchange.domain.MsgEvent
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.domain.Targets
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.infra.persistence.MsgEventRepository
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.infra.retry.FailureScheduler
|
2026-09-06 16:13:50 +08:00
|
|
|
|
import jakarta.inject.Singleton
|
|
|
|
|
|
import java.time.Duration
|
|
|
|
|
|
import java.time.Instant
|
2026-09-11 08:08:35 +08:00
|
|
|
|
import java.util.concurrent.atomic.AtomicLong
|
2026-09-06 16:13:50 +08:00
|
|
|
|
|
2026-09-10 11:14:11 +08:00
|
|
|
|
/** 对外投递的出口:同步等 Kafka 确认,语义是至少发一次(可能重复,但不会丢)。以后 ES 投影也从这里加。 */
|
2026-09-06 16:13:50 +08:00
|
|
|
|
interface DeliveryPort {
|
2026-09-10 11:14:11 +08:00
|
|
|
|
/** 发一条变化通知,消息 key 是 FLID(航班实例 ID)。 */
|
2026-09-09 17:53:08 +08:00
|
|
|
|
fun sendKafka(topic: String, key: String, payloadJson: String)
|
2026-09-06 16:13:50 +08:00
|
|
|
|
|
2026-09-10 11:14:11 +08:00
|
|
|
|
/** 发一条完整状态,消息 key 是 FLID。适配层负责把 target 映射成 topic:KAFKA:msg 对应 "msg",KAFKA:schd 对应 "schd"。 */
|
2026-09-09 17:53:08 +08:00
|
|
|
|
fun sendKafkaSchd(topic: String, key: String, payloadJson: String)
|
2026-09-06 16:13:50 +08:00
|
|
|
|
|
2026-09-10 11:14:11 +08:00
|
|
|
|
/** 发一条删除通知:key 是 FLID、value 为空;下游按"整态里这个键没了"理解成删除。 */
|
2026-09-09 17:53:08 +08:00
|
|
|
|
fun sendKafkaNull(topic: String, key: String)
|
2026-09-06 22:12:23 +08:00
|
|
|
|
|
2026-09-10 11:14:11 +08:00
|
|
|
|
/** 给健康检查用的连通性探测;默认返回 true,真实 Kafka 实现要覆写成向 broker 拉一次 metadata 来判断。 */
|
2026-09-06 22:12:23 +08:00
|
|
|
|
fun ping(): Boolean = true
|
2026-09-06 16:13:50 +08:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
/**
|
2026-09-10 11:14:11 +08:00
|
|
|
|
* 投递调度:把 outbox(待发事件表)里的事件发给下游。
|
|
|
|
|
|
*
|
|
|
|
|
|
* KAFKA_MSG 一条一条按登记顺序发,不插队。KAFKA_SCHD 走 flushSchd 批量发:同一个 FLID 攒了
|
|
|
|
|
|
* 多条未发事件时只发版本号最新的那条,旧的自然作废;删除通知发 value 为空的 tombstone。
|
|
|
|
|
|
* 两个主题之间不保证先后顺序。
|
|
|
|
|
|
*
|
|
|
|
|
|
* 失败处理:队首的重试时间没到就不取;一批里有发送失败,整批重试次数加一并推后退避,
|
|
|
|
|
|
* 次数用尽整批转 DEAD 当死信。见 docs/flight-state.md §5。
|
2026-09-06 16:13:50 +08:00
|
|
|
|
*/
|
|
|
|
|
|
@Singleton
|
|
|
|
|
|
class Dispatcher(
|
|
|
|
|
|
private val msgEvents: MsgEventRepository,
|
|
|
|
|
|
private val port: DeliveryPort,
|
|
|
|
|
|
private val props: PipelineProps,
|
2026-09-06 19:08:02 +08:00
|
|
|
|
private val scheduler: FailureScheduler,
|
2026-09-06 16:13:50 +08:00
|
|
|
|
) {
|
2026-09-06 21:02:27 +08:00
|
|
|
|
private val log = org.slf4j.LoggerFactory.getLogger(Dispatcher::class.java)
|
|
|
|
|
|
|
2026-09-11 08:08:35 +08:00
|
|
|
|
/** 连续失败计数:仅用于日志/排障。 */
|
|
|
|
|
|
private val tickFailures = AtomicLong(0)
|
|
|
|
|
|
|
2026-09-06 16:13:50 +08:00
|
|
|
|
@Volatile
|
|
|
|
|
|
private var running = true
|
|
|
|
|
|
|
2026-09-09 17:53:08 +08:00
|
|
|
|
private var lastFlush: Instant? = null
|
2026-09-06 16:13:50 +08:00
|
|
|
|
|
2026-09-10 11:14:11 +08:00
|
|
|
|
/** 请求停机:置位后 loop 走完当前一轮就退出;真正中断线程由 PipelineLifecycle 负责。 */
|
2026-09-06 20:53:14 +08:00
|
|
|
|
fun stop() {
|
|
|
|
|
|
running = false
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-09-06 16:13:50 +08:00
|
|
|
|
fun loop() {
|
|
|
|
|
|
while (running) {
|
2026-09-06 18:19:04 +08:00
|
|
|
|
try {
|
|
|
|
|
|
tick()
|
2026-09-11 08:08:35 +08:00
|
|
|
|
tickFailures.set(0)
|
2026-09-06 19:08:02 +08:00
|
|
|
|
} catch (e: InterruptedException) {
|
|
|
|
|
|
Thread.currentThread().interrupt()
|
|
|
|
|
|
return
|
2026-09-06 18:19:04 +08:00
|
|
|
|
} catch (e: Exception) {
|
2026-09-11 08:08:35 +08:00
|
|
|
|
// tick 整体失败(DB 断连等)必须在别处留痕,否则只剩"投递不动"没有原因。
|
|
|
|
|
|
val failures = tickFailures.incrementAndGet()
|
|
|
|
|
|
log.error("dispatcher tick failed (failure #{})", failures, e)
|
2026-09-06 18:19:04 +08:00
|
|
|
|
sleepQuietly(props.pipeline.pollInterval)
|
|
|
|
|
|
}
|
2026-09-06 16:13:50 +08:00
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
internal fun tick() {
|
2026-09-11 08:08:35 +08:00
|
|
|
|
val sent = drainKafkaMsg()
|
2026-09-09 17:53:08 +08:00
|
|
|
|
if (flushDue()) flushSchd()
|
2026-09-11 08:08:35 +08:00
|
|
|
|
// 有活就不睡:pollInterval 只用于"本轮无事可做或队头在退避",不能当投递节流用。
|
|
|
|
|
|
if (sent == 0 && running) sleepQuietly(props.pipeline.pollInterval)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
|
* 按 `EVENT_ID` 顺序批量投递 `KAFKA:msg`。
|
|
|
|
|
|
*
|
|
|
|
|
|
* 保序规则不变:**同一目标内不越序**——队头失败或退避未到期时立即停止本轮,
|
|
|
|
|
|
* 不跳过它去投后面的。批量只用来省掉"每条一次 DB 往返 + 一轮一次 sleep"。
|
|
|
|
|
|
*
|
|
|
|
|
|
* @return 本轮成功发出的条数
|
|
|
|
|
|
*/
|
|
|
|
|
|
private fun drainKafkaMsg(): Int {
|
|
|
|
|
|
val batchSize = props.pipeline.deliveryBatch.coerceAtLeast(1)
|
|
|
|
|
|
val maxRounds = props.pipeline.deliveryDrainRounds.coerceAtLeast(1)
|
|
|
|
|
|
var sent = 0
|
|
|
|
|
|
var round = 0
|
|
|
|
|
|
while (round < maxRounds) {
|
|
|
|
|
|
round++
|
|
|
|
|
|
val batch = try {
|
|
|
|
|
|
msgEvents.claimBatch(Targets.KAFKA_MSG, batchSize)
|
|
|
|
|
|
} catch (e: Exception) {
|
|
|
|
|
|
log.error("claimBatch(KAFKA_MSG) failed", e)
|
|
|
|
|
|
return sent
|
|
|
|
|
|
}
|
|
|
|
|
|
if (batch.isEmpty()) return sent
|
|
|
|
|
|
for (e in batch) {
|
|
|
|
|
|
val next = e.nextAttemptAt
|
|
|
|
|
|
if (next != null && next > scheduler.now()) return sent // 队头退避未到期:停止推进
|
|
|
|
|
|
if (!deliver(Targets.KAFKA_MSG, e)) return sent // 失败已记账:停止以保序
|
|
|
|
|
|
sent++
|
|
|
|
|
|
}
|
|
|
|
|
|
if (batch.size < batchSize) return sent
|
|
|
|
|
|
}
|
|
|
|
|
|
log.warn("kafka msg drain hit round cap ({} rounds, sent={}); yielding to schd flush", maxRounds, sent)
|
|
|
|
|
|
return sent
|
2026-09-06 16:13:50 +08:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-09-06 18:19:04 +08:00
|
|
|
|
private fun flushDue(): Boolean =
|
2026-09-10 20:39:22 +08:00
|
|
|
|
lastFlush?.let { Duration.between(it, scheduler.now()) >= props.schd.flushPeriod } ?: true
|
2026-09-06 19:08:02 +08:00
|
|
|
|
|
2026-09-11 08:08:35 +08:00
|
|
|
|
/** @return 是否已确认发出(false 表示已按退避/死信记账,调用方应停止本轮以保序) */
|
|
|
|
|
|
private fun deliver(target: String, e: MsgEvent): Boolean {
|
|
|
|
|
|
val eventId = e.eventId ?: return false
|
|
|
|
|
|
return try {
|
2026-09-09 17:53:08 +08:00
|
|
|
|
when (target) {
|
|
|
|
|
|
Targets.KAFKA_MSG -> port.sendKafka("msg", e.partitionKey, e.payloadJson)
|
2026-09-11 08:08:35 +08:00
|
|
|
|
// 未知目标不能在"未发送"的情况下被标记为已发;记账后停止本轮。
|
|
|
|
|
|
else -> error("unknown delivery target: $target")
|
2026-09-09 17:53:08 +08:00
|
|
|
|
}
|
2026-09-11 08:08:35 +08:00
|
|
|
|
msgEvents.markSent(eventId)
|
|
|
|
|
|
true
|
2026-09-09 17:53:08 +08:00
|
|
|
|
} catch (ex: Exception) {
|
|
|
|
|
|
retryOrDead(e, ex.message ?: ex.javaClass.simpleName)
|
2026-09-11 08:08:35 +08:00
|
|
|
|
false
|
2026-09-06 19:08:02 +08:00
|
|
|
|
}
|
|
|
|
|
|
}
|
2026-09-06 18:19:04 +08:00
|
|
|
|
|
2026-09-10 11:14:11 +08:00
|
|
|
|
/** 批量发 KAFKA_SCHD:每个 FLID 只发版本号最新的那条未发事件,删除通知发空 value。 */
|
2026-09-06 16:13:50 +08:00
|
|
|
|
internal fun flushSchd() {
|
2026-09-09 17:53:08 +08:00
|
|
|
|
val batch = try {
|
2026-09-10 20:39:22 +08:00
|
|
|
|
msgEvents.mergePendingSchd(scheduler.now(), props.schd.flushLimit)
|
2026-09-09 17:53:08 +08:00
|
|
|
|
} catch (e: Exception) {
|
|
|
|
|
|
log.warn("mergePendingSchd failed: {}", e.message)
|
|
|
|
|
|
return
|
|
|
|
|
|
}
|
|
|
|
|
|
if (batch.isEmpty()) {
|
2026-09-06 19:08:02 +08:00
|
|
|
|
lastFlush = scheduler.now()
|
2026-09-06 18:19:04 +08:00
|
|
|
|
return
|
|
|
|
|
|
}
|
2026-09-09 17:53:08 +08:00
|
|
|
|
val failures = mutableListOf<MsgEvent>()
|
|
|
|
|
|
for (e in batch) {
|
|
|
|
|
|
try {
|
|
|
|
|
|
when (e.eventType) {
|
|
|
|
|
|
EventType.TOMBSTONE -> port.sendKafkaNull("schd", e.partitionKey)
|
|
|
|
|
|
EventType.UPSERT -> port.sendKafkaSchd("schd", e.partitionKey, e.payloadJson)
|
|
|
|
|
|
}
|
|
|
|
|
|
} catch (ex: Exception) {
|
|
|
|
|
|
failures.add(e)
|
|
|
|
|
|
}
|
2026-09-06 19:08:02 +08:00
|
|
|
|
}
|
2026-09-09 17:53:08 +08:00
|
|
|
|
val sentIds = batch.mapNotNull { it.eventId }.toSet() - failures.mapNotNull { it.eventId }.toSet()
|
|
|
|
|
|
if (sentIds.isNotEmpty()) msgEvents.markAllSent(sentIds.toList())
|
|
|
|
|
|
val sentVersions = batch.associate { it.partitionKey to it.stateVersion }
|
2026-09-11 08:08:35 +08:00
|
|
|
|
// 被更新版本压掉的旧事件也要标成已发,否则它们会一直留在队里:同一个 FLID 只按最新版本输出一次。
|
|
|
|
|
|
// 这里必须**有界**:若 markAllSent 因任何原因没有生效(例如行缺 eventId),
|
|
|
|
|
|
// 旧的无界 while(true) 会原地空转。
|
2026-09-09 17:53:08 +08:00
|
|
|
|
runCatching {
|
2026-09-11 08:08:35 +08:00
|
|
|
|
var rounds = 0
|
|
|
|
|
|
while (rounds < SUPERSEDED_CLEANUP_MAX_ROUNDS) {
|
|
|
|
|
|
rounds++
|
2026-09-10 20:39:22 +08:00
|
|
|
|
val superseded = msgEvents.mergePendingSchd(scheduler.now(), props.schd.flushLimit)
|
2026-09-09 17:53:08 +08:00
|
|
|
|
.filter { sentVersions[it.partitionKey]?.let { v -> it.stateVersion < v } == true }
|
|
|
|
|
|
if (superseded.isEmpty()) break
|
2026-09-11 08:08:35 +08:00
|
|
|
|
val ids = superseded.mapNotNull { it.eventId }
|
|
|
|
|
|
if (ids.isEmpty()) {
|
|
|
|
|
|
log.warn("superseded cleanup: {} rows without event_id, stop to avoid a spin", superseded.size)
|
|
|
|
|
|
break
|
|
|
|
|
|
}
|
|
|
|
|
|
msgEvents.markAllSent(ids)
|
|
|
|
|
|
if (rounds == SUPERSEDED_CLEANUP_MAX_ROUNDS) {
|
|
|
|
|
|
log.warn("superseded cleanup hit round cap ({}); remaining rows will be handled next flush", rounds)
|
|
|
|
|
|
}
|
2026-09-09 17:53:08 +08:00
|
|
|
|
}
|
|
|
|
|
|
}.onFailure { log.warn("superseded cleanup failed: {}", it.message) }
|
|
|
|
|
|
failures.forEach { retryOrDead(it, it.lastError ?: "send-failed") }
|
2026-09-06 19:08:02 +08:00
|
|
|
|
lastFlush = scheduler.now()
|
2026-09-06 18:19:04 +08:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-09-10 11:14:11 +08:00
|
|
|
|
/** 一条事件发失败之后怎么走:重试次数加一,到上限就标成 DEAD(EXHAUSTED) 留作死信(次数落库便于追查),否则按退避推到下次再发。 */
|
2026-09-09 17:53:08 +08:00
|
|
|
|
private fun retryOrDead(e: MsgEvent, lastError: String) {
|
|
|
|
|
|
val eventId = e.eventId ?: return
|
|
|
|
|
|
val attempts = e.attempts + 1
|
|
|
|
|
|
if (scheduler.exhausted(attempts)) {
|
|
|
|
|
|
msgEvents.markDead(eventId, ErrorClass.EXHAUSTED, lastError, attempts)
|
|
|
|
|
|
} else {
|
|
|
|
|
|
msgEvents.scheduleRetry(eventId, scheduler.nextAttemptAt(attempts), attempts)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-09-06 18:19:04 +08:00
|
|
|
|
private fun sleepQuietly(d: Duration) {
|
|
|
|
|
|
if (!d.isNegative && !d.isZero) Thread.sleep(d.toMillis().coerceAtLeast(1))
|
2026-09-06 16:13:50 +08:00
|
|
|
|
}
|
2026-09-11 08:08:35 +08:00
|
|
|
|
|
|
|
|
|
|
private companion object {
|
|
|
|
|
|
/** superseded 清理的轮数上限:只用于防止"标不掉又不报错"时的原地空转。 */
|
|
|
|
|
|
const val SUPERSEDED_CLEANUP_MAX_ROUNDS = 100
|
|
|
|
|
|
}
|
2026-09-06 16:13:50 +08:00
|
|
|
|
}
|