7.5 配套代码:韧性客户端与超时审计
对应小节:7.5 背压与韧性 三部分:① 完整的韧性客户端;② 超时预算审计;③ 受限重试的实现。
一、完整的韧性客户端
// src/main/kotlin/resilience/ResilientClient.kt
package resilience
import io.micrometer.core.instrument.Counter
import io.micrometer.core.instrument.MeterRegistry
import kotlinx.coroutines.TimeoutCancellationException
import kotlinx.coroutines.delay
import kotlinx.coroutines.sync.Semaphore
import kotlinx.coroutines.sync.withPermit
import kotlinx.coroutines.withTimeout
import java.util.concurrent.atomic.AtomicInteger
import java.util.concurrent.atomic.AtomicLong
import kotlin.math.min
import kotlin.random.Random
/**
* 五种韧性手段的组合实现:
* ① 限流(RateLimiter)—— 防止自己被打垮
* ② 熔断(CircuitBreaker)—— 防止被下游拖垮
* ③ 隔板(Semaphore)—— 防止一个依赖拖垮全部
* ④ 超时(withTimeout)—— 及时放弃,释放资源
* ⑤ 重试(带退避与上限)—— 应对偶发失败
*/
class ResilientClient(
private val name: String,
private val call: suspend (String) -> String,
// ① 限流:每秒最多 N 个请求
private val maxQps: Int = 1000,
// ③ 隔板:最多 N 个并发
private val maxConcurrent: Int = 32,
// ④ 超时:单次调用的超时
private val timeoutMs: Long = 100,
// ⑤ 重试
private val maxAttempts: Int = 3,
private val initialBackoffMs: Long = 50,
// ② 熔断
private val failureThreshold: Double = 0.5, // 错误率阈值
private val minRequestsForCircuit: Int = 20, // 样本不足不触发
private val openDurationMs: Long = 20_000, // 打开时长
private val halfOpenProbes: Int = 3, // 半开探测数
registry: MeterRegistry,
) {
// ── 限流(令牌桶的简化实现)──────────────────────────────
private val tokens = AtomicInteger(maxQps)
private val lastRefill = AtomicLong(System.currentTimeMillis())
// ── 隔板 ────────────────────────────────────────────────
private val semaphore = Semaphore(maxConcurrent)
// ── 熔断状态 ────────────────────────────────────────────
private enum class State { CLOSED, OPEN, HALF_OPEN }
@Volatile private var circuitState = State.CLOSED
@Volatile private var openedAt = 0L
private val totalRequests = AtomicLong(0)
private val failedRequests = AtomicLong(0)
private val halfOpenSuccess = AtomicInteger(0)
// ── 指标 ────────────────────────────────────────────────
private val callsTotal = Counter.builder("client.calls.total").tag("downstream", name).register(registry)
private val callsFailed = Counter.builder("client.calls.failed").tag("downstream", name).register(registry)
private val callsRejected = Counter.builder("client.calls.rejected").tag("downstream", name).register(registry)
private val callsTimedOut = Counter.builder("client.calls.timeout").tag("downstream", name).register(registry)
private val callsRetried = Counter.builder("client.calls.retried").tag("downstream", name).register(registry)
private val circuitOpened = Counter.builder("client.circuit.opened").tag("downstream", name).register(registry)
private fun tryAcquireToken(): Boolean {
val now = System.currentTimeMillis()
val elapsed = now - lastRefill.get()
if (elapsed >= 1000) { // 每秒补充
tokens.set(maxQps)
lastRefill.set(now)
}
return tokens.getAndUpdate { if (it > 0) it - 1 else 0 } > 0
}
suspend fun fetch(id: String): Result<String> {
// ① 限流
if (!tryAcquireToken()) {
callsRejected.increment()
return Result.failure(RejectedException("$name 超过速率限制"))
}
// ② 熔断检查
if (isCircuitOpen()) {
callsRejected.increment()
return Result.failure(CircuitOpenException("$name 熔断中"))
}
callsTotal.increment()
return try {
val value = withRetry {
// ③ 隔板 + ④ 超时
semaphore.withPermit {
withTimeout(timeoutMs) { call(id) }
}
}
recordSuccess()
Result.success(value)
} catch (e: TimeoutCancellationException) {
callsTimedOut.increment()
recordFailure()
Result.failure(TimeoutException("$name 超时(${timeoutMs}ms)", e))
} catch (e: Exception) {
callsFailed.increment()
recordFailure()
Result.failure(e)
}
}
/** ⑤ 重试:指数退避 + 抖动 + 次数上限 */
private suspend fun <T> withRetry(block: suspend () -> T): T {
var lastError: Throwable? = null
repeat(maxAttempts) { attempt ->
try {
return block()
} catch (e: TimeoutCancellationException) {
lastError = e
if (attempt < maxAttempts - 1) {
callsRetried.increment()
val base = min(initialBackoffMs shl attempt, 2_000)
delay(base + Random.nextLong(0, base / 2)) // 抖动:避免惊群
}
} catch (e: NonRetryableException) {
throw e // 不重试不可重试的异常
} catch (e: Exception) {
lastError = e
if (attempt < maxAttempts - 1) {
callsRetried.increment()
val base = min(initialBackoffMs shl attempt, 2_000)
delay(base + Random.nextLong(0, base / 2))
}
}
}
throw lastError ?: IllegalStateException("unreachable")
}
// ── 熔断状态机 ──────────────────────────────────────────
private fun isCircuitOpen(): Boolean = when (circuitState) {
State.CLOSED -> false
State.OPEN -> {
if (System.currentTimeMillis() - openedAt > openDurationMs) {
circuitState = State.HALF_OPEN
halfOpenSuccess.set(0)
false
} else true
}
State.HALF_OPEN -> false
}
private fun recordSuccess() {
if (circuitState == State.HALF_OPEN) {
if (halfOpenSuccess.incrementAndGet() >= halfOpenProbes) {
circuitState = State.CLOSED
totalRequests.set(0)
failedRequests.set(0)
}
} else {
totalRequests.incrementAndGet()
}
}
private fun recordFailure() {
if (circuitState == State.HALF_OPEN) {
circuitState = State.OPEN
openedAt = System.currentTimeMillis()
circuitOpened.increment()
return
}
val total = totalRequests.incrementAndGet()
val failed = failedRequests.incrementAndGet()
if (total >= minRequestsForCircuit && failed.toDouble() / total > failureThreshold) {
circuitState = State.OPEN
openedAt = System.currentTimeMillis()
circuitOpened.increment()
}
}
}
class RejectedException(message: String) : RuntimeException(message)
class CircuitOpenException(message: String) : RuntimeException(message)
class TimeoutException(message: String, cause: Throwable?) : RuntimeException(message, cause)
class NonRetryableException(message: String) : RuntimeException(message)
二、超时预算审计
// src/main/kotlin/resilience/TimeoutAudit.kt
package resilience
/**
* 超时链审计:找出「本层超时 > 上游给的预算」的配置。
*
* 这是决定「系统崩溃后能否自愈」的关键(第 3 章 3.6 节):
* 如果应用给数据库的超时(30s)远大于网关给应用的超时(1s),
* 那么网关超时返回后,应用线程仍然被占住 30 秒 → 连接池持续满 → 无法自愈。
*/
data class TimeoutLayer(
val name: String,
val upstreamBudgetMs: Long, // 上游给这一层的时间
val selfTimeoutMs: Long, // 这一层设置的超时
val note: String = "",
)
data class AuditResult(
val problems: List<String>,
val warnings: List<String>,
val chain: List<Pair<String, Long>>, // (层名, 留给下游的时间)
)
fun auditTimeouts(layers: List<TimeoutLayer>): AuditResult {
val problems = mutableListOf<String>()
val warnings = mutableListOf<String>()
val chain = mutableListOf<Pair<String, Long>>()
var remainingBudget = layers.firstOrNull()?.upstreamBudgetMs ?: Long.MAX_VALUE
layers.forEach { layer ->
// ① 本层超时不应超过上游给的预算
if (layer.selfTimeoutMs > layer.upstreamBudgetMs) {
problems += buildString {
append("❌ ${layer.name}: 本层超时 ${layer.selfTimeoutMs}ms > 上游预算 ${layer.upstreamBudgetMs}ms")
append("\n 后果:上游已超时返回,本层的调用仍在继续,资源被占住 → 可能无法自愈")
if (layer.note.isNotEmpty()) append("\n 备注:${layer.note}")
}
}
// ② 上游预算不应超过更上游给的预算(逐层递减)
if (layer.upstreamBudgetMs > remainingBudget) {
problems += "❌ ${layer.name}: 上游给了 ${layer.upstreamBudgetMs}ms," +
"但更上游只剩 ${remainingBudget}ms(超时链未递减)"
}
// ③ 检查余量是否合理
val margin = layer.upstreamBudgetMs - layer.selfTimeoutMs
if (margin in 1..5) {
warnings += "⚠️ ${layer.name}: 余量只有 ${margin}ms,可能不足以处理上游的开销"
}
if (layer.selfTimeoutMs <= 0) {
problems += "❌ ${layer.name}: 超时为 ${layer.selfTimeoutMs}ms(无效值)"
}
remainingBudget = layer.selfTimeoutMs
chain += layer.name to layer.selfTimeoutMs
}
// ④ 检查是否所有层都配置了超时
if (layers.any { it.selfTimeoutMs <= 0 }) {
warnings += "⚠️ 存在未配置超时的层(默认超时可能非常长,如 30 秒)"
}
return AuditResult(problems, warnings, chain)
}
fun main() {
println("═══ 反例:会导致无法自愈的配置 ═══")
println()
val bad = auditTimeouts(listOf(
TimeoutLayer("客户端 → 网关", upstreamBudgetMs = 5_000, selfTimeoutMs = 3_000),
TimeoutLayer("网关 → 应用", upstreamBudgetMs = 3_000, selfTimeoutMs = 1_000),
TimeoutLayer("应用 → 数据库", upstreamBudgetMs = 30_000, selfTimeoutMs = 30_000,
note = "默认值,没人改过 —— 这是最常见的坑"),
TimeoutLayer("应用 → Redis", upstreamBudgetMs = 30_000, selfTimeoutMs = 5_000),
))
bad.problems.forEach { println(it); println() }
bad.warnings.forEach { println(it) }
println()
println("超时链:")
bad.chain.forEachIndexed { i, (name, ms) ->
println(" ${" ".repeat(i)}└─ $name: ${ms}ms")
}
println()
println("═".repeat(70))
println()
println("═══ 正例:逐层递减,可自愈 ═══")
println()
val good = auditTimeouts(listOf(
TimeoutLayer("客户端 → 网关", upstreamBudgetMs = 5_000, selfTimeoutMs = 500),
TimeoutLayer("网关 → 应用", upstreamBudgetMs = 500, selfTimeoutMs = 300),
TimeoutLayer("应用 → 数据库", upstreamBudgetMs = 300, selfTimeoutMs = 60),
TimeoutLayer("应用 → Redis", upstreamBudgetMs = 300, selfTimeoutMs = 10),
TimeoutLayer("应用 → 下游", upstreamBudgetMs = 300, selfTimeoutMs = 100),
))
if (good.problems.isEmpty()) {
println("✅ 超时链自洽,逐层递减")
} else {
good.problems.forEach { println(it) }
}
good.warnings.forEach { println(it) }
println()
println("超时链:")
good.chain.forEachIndexed { i, (name, ms) ->
println(" ${" ".repeat(i)}└─ $name: ${ms}ms")
}
println()
println("核心规则:本层超时必须小于上游预算,且逐层递减。")
println("超时不是「容忍时间」,而是「放弃的时机」—— 放弃得早,系统才能自愈。")
}
预期输出(反例部分):
❌ 应用 → 数据库: 本层超时 30000ms > 上游预算 30000ms
后果:上游已超时返回,本层的调用仍在继续,资源被占住 → 可能无法自愈
备注:默认值,没人改过 —— 这是最常见的坑
❌ 应用 → Redis: 上游给了 30000ms,但更上游只剩 30000ms(超时链未递减)
三、受限重试的实现
// src/main/kotlin/resilience/Retry.kt
package resilience
import kotlinx.coroutines.delay
import java.util.concurrent.atomic.AtomicLong
import kotlin.math.min
import kotlin.random.Random
/**
* 带四条纪律的重试:
* ① 有次数上限
* ② 指数退避 + 抖动
* ③ 只重试可重试的异常
* ④ 有全局重试预算(防止重试风暴)
*/
class RetryPolicy(
private val maxAttempts: Int = 3,
private val initialDelayMs: Long = 100,
private val maxDelayMs: Long = 2_000,
private val retryBudgetRatio: Double = 0.1, // 重试请求不超过正常请求的 10%
) {
private val totalRequests = AtomicLong(0)
private val retriedRequests = AtomicLong(0)
/** 检查是否还有重试预算 */
private fun hasBudget(): Boolean {
val total = totalRequests.get()
if (total < 100) return true // 样本太少时不限制
return retriedRequests.get().toDouble() / total < retryBudgetRatio
}
suspend fun <T> execute(
operation: String,
isRetryable: (Throwable) -> Boolean = { it !is NonRetryableException },
block: suspend (attempt: Int) -> T,
): T {
totalRequests.incrementAndGet()
var lastError: Throwable? = null
repeat(maxAttempts) { attempt ->
try {
return block(attempt)
} catch (e: Throwable) {
lastError = e
// ③ 只重试可重试的异常
if (!isRetryable(e)) throw e
// ① 次数上限
if (attempt >= maxAttempts - 1) throw e
// ④ 全局预算
if (!hasBudget()) {
throw RetryBudgetExceededException(
"重试预算已用尽($operation),放弃重试", e)
}
retriedRequests.incrementAndGet()
// ② 指数退避 + 抖动
val base = min(initialDelayMs shl attempt, maxDelayMs)
val jitter = Random.nextLong(0, base / 2)
delay(base + jitter)
}
}
throw lastError ?: IllegalStateException("unreachable")
}
fun stats(): Pair<Long, Long> = totalRequests.get() to retriedRequests.get()
}
class RetryBudgetExceededException(message: String, cause: Throwable?) : RuntimeException(message, cause)
关键说明:
| 纪律 | 实现 |
|---|---|
| ① 次数上限 | maxAttempts(默认 3) |
| ② 指数退避 + 抖动 | min(initialDelayMs shl attempt, maxDelayMs) + 随机抖动 |
| ③ 只重试可重试异常 | isRetryable 谓词(NonRetryableException 不重试) |
| ④ 全局重试预算 | retriedRequests / totalRequests < retryBudgetRatio |
重试预算的价值:防止"下游已经过载,所有客户端都在重试"的正反馈——这是重试风暴的根治手段。
四、动手改造
| 改动 | 观察什么 |
|---|---|
在 ResilientClient 里把 maxAttempts 改成 1 |
失败率上升但下游压力下降——理解重试的代价 |
把 initialBackoffMs 改成 0(无退避) |
下游被打得更狠(用压测观察下游的错误率) |
去掉 RetryPolicy 的预算检查 |
在下游故障时,重试流量会倍增——这就是重试风暴 |
把 timeoutMs 从 100 改成 5000 |
恢复时间显著变长(第 3 章 3.6 节的实验) |
用 auditTimeouts 审计你项目的真实配置 |
大概率会找到"默认 30 秒"的坑 |
| 故意让下游持续失败,观察熔断 | 熔断打开后,调用应该快速失败(而不是每次等超时) |
五、这段代码的局限
ResilientClient的限流是简化实现(固定窗口 + 每秒补充):生产建议用 GuavaRateLimiter(令牌桶)或 Resilience4j。- 熔断器的实现是手写的:生产建议用 Resilience4j(有更完善的半开、滑动窗口等语义)。
- 重试预算的检查有竞态:
hasBudget()用的是近似统计,高并发下可能略微超出预算——这是可接受的(预算本身是模糊约束)。 withTimeout依赖协程的取消机制:如果call里有不可取消的阻塞调用(比如 JDBC),超时不会真正中断它——这也是为什么必须把阻塞调用切到 IO 并设置驱动级超时(第 2 章 2.7)。- 五种手段组合起来复杂度高:不是所有调用都需要全套——按依赖的重要性选择(核心依赖用全套,非核心只用超时 + 熔断)。