文档目录

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 的限流是简化实现(固定窗口 + 每秒补充):生产建议用 Guava RateLimiter(令牌桶)或 Resilience4j。
  • 熔断器的实现是手写的:生产建议用 Resilience4j(有更完善的半开、滑动窗口等语义)。
  • 重试预算的检查有竞态:hasBudget() 用的是近似统计,高并发下可能略微超出预算——这是可接受的(预算本身是模糊约束)。
  • withTimeout 依赖协程的取消机制:如果 call 里有不可取消的阻塞调用(比如 JDBC),超时不会真正中断它——这也是为什么必须把阻塞调用切到 IO 并设置驱动级超时(第 2 章 2.7)。
  • 五种手段组合起来复杂度高:不是所有调用都需要全套——按依赖的重要性选择(核心依赖用全套,非核心只用超时 + 熔断)。