feat(jobs): expose last job failure time

JobActivity 记录最近失败时刻并进入快照;JobRunner 失败时上报该时刻;新增指标 msgx.pipeline.job.last_failure_age_seconds 与健康明细 lastFailureAgeSeconds(对齐 reference 作业健康口径)。

验证:./gradlew test 145 tests / 0 fail / 0 skipped。
This commit is contained in:
windyboy
2026-09-12 20:56:41 +08:00
parent 8a380c7646
commit 267eb20219
6 changed files with 22 additions and 6 deletions
@@ -45,6 +45,7 @@ internal fun jobRunnerHealth(snapshot: JobActivity.Snapshot, now: Instant): Heal
"started" to snapshot.started,
"lastTickAgeSeconds" to (ageSeconds ?: -1L),
"lastTickDurationMs" to snapshot.lastTickMillis,
"lastFailureAgeSeconds" to (snapshot.lastFailureAt?.let { Duration.between(it, now).seconds } ?: -1L),
"ticks" to snapshot.ticks,
"failures" to snapshot.failures,
"lastSweepSelected" to snapshot.lastSweepSelected,
@@ -22,6 +22,9 @@ class JobActivity {
@Volatile
private var lastTickMillis = 0L
@Volatile
private var lastFailureAt: Instant? = null
/** 上一轮回填扫描选中的待办条数(扫描积压的即时值)。 */
@Volatile
var lastSweepSelected = 0
@@ -46,18 +49,20 @@ class JobActivity {
ticks.incrementAndGet()
}
/** 一次作业 tick 抛错:只累计失败,**不更新心跳时刻**——停摆由心跳年龄暴露。 */
fun tickFailed() {
/** 一次作业 tick 抛错:记下最近失败时刻,**不更新心跳时刻**——停摆由心跳年龄暴露。 */
fun tickFailed(now: Instant) {
lastFailureAt = now
failures.incrementAndGet()
}
fun snapshot(): Snapshot =
Snapshot(started, lastTickAt, lastTickMillis, lastSweepSelected, ticks.get(), failures.get())
Snapshot(started, lastTickAt, lastTickMillis, lastFailureAt, lastSweepSelected, ticks.get(), failures.get())
data class Snapshot(
val started: Boolean,
val lastTickAt: Instant?,
val lastTickMillis: Long,
val lastFailureAt: Instant?,
val lastSweepSelected: Int,
val ticks: Long,
val failures: Long,
@@ -24,6 +24,7 @@ import java.time.Duration
* - `msgx.pipeline.watermark.lag`:水位落后信箱最新 ID 的距离
* - `msgx.pipeline.hole.aged_out.total`:永久空洞放行次数
* - `msgx.pipeline.job.heartbeat_age_seconds`:距上一次作业 tick 完成的秒数(未跑过为 -1)
* - `msgx.pipeline.job.last_failure_age_seconds`:距最近一次作业 tick 失败的秒数(从未失败为 -1)
* - `msgx.pipeline.job.ticks.total` / `msgx.pipeline.job.failures.total`:作业 tick 完成/抛错次数
* - `msgx.pipeline.job.last_sweep_selected`:上一轮回填扫描选中的待办条数(扫描积压)
*
@@ -84,6 +85,10 @@ class PipelineMetrics(
a.snapshot().lastTickAt?.let { Duration.between(it, clock.instant()).seconds.toDouble() } ?: -1.0
}.strongReference(true).register(registry)
Gauge.builder("msgx.pipeline.job.last_failure_age_seconds", activity) { a ->
a.snapshot().lastFailureAt?.let { Duration.between(it, clock.instant()).seconds.toDouble() } ?: -1.0
}.strongReference(true).register(registry)
Gauge.builder("msgx.pipeline.job.ticks.total", activity) { it.snapshot().ticks.toDouble() }
.strongReference(true)
.register(registry)
@@ -79,7 +79,7 @@ class JobRunner(
} catch (e: InterruptedException) {
throw e
} catch (e: Exception) {
activity.tickFailed()
activity.tickFailed(startedAt)
log.warn("job tick failed: {}", e.message ?: e.javaClass.simpleName)
}
}
@@ -59,12 +59,16 @@ class JobRunnerHealthIndicatorTest {
fun `a failed tick counts up without refreshing the heartbeat`() {
val activity = JobActivity().apply { started() }
activity.tickFailed()
activity.tickFailed()
activity.tickFailed(now.minusSeconds(5))
activity.tickFailed(now)
val snapshot = activity.snapshot()
assertEquals(2L, snapshot.failures)
assertEquals(0L, snapshot.ticks)
assertEquals(null, snapshot.lastTickAt)
val details = details(jobRunnerHealth(snapshot, now))
assertEquals(0L, details["lastFailureAgeSeconds"])
assertEquals(2L, details["failures"])
}
}
@@ -69,6 +69,7 @@ class JobRunnerTest {
val snapshot = activity.snapshot()
assertNull(snapshot.lastTickAt)
assertEquals(MutableClock.BASE, snapshot.lastFailureAt)
assertEquals(0L, snapshot.ticks)
assertEquals(1L, snapshot.failures)
}