文档目录

7.3 配套代码:并行度调优与隔板

对应小节:7.3 并发与并行 两件事:① 用实验找到最优并行度;② 用隔板实现资源隔离。

一、并行度调优实验

// src/main/kotlin/experiments/ParallelismTuning.kt
package experiments

import kotlinx.coroutines.*
import kotlinx.coroutines.sync.Semaphore
import kotlinx.coroutines.sync.withPermit
import kotlin.system.measureNanoTime

private var sink = 0L

/** 模拟一次"要 5ms 的 IO 操作" */
suspend fun simulatedIo(): Int {
    delay(5)
    return 1
}

/** 模拟一次"要 1ms 的 CPU 操作" */
fun simulatedCpu(): Int {
    var s = 0
    repeat(10_000) { s += it }
    return s
}

/**
 * 测量不同并行度下的吞吐与延迟。
 *
 * @param parallelism 并发上限
 * @param totalOps 总操作数
 * @param ioBound true = IO 密集(delay),false = CPU 密集
 */
suspend fun measure(parallelism: Int, totalOps: Int, ioBound: Boolean): Triple<Double, Double, Double> {
    val limiter = Semaphore(parallelism)
    val latencies = java.util.concurrent.ConcurrentLinkedQueue<Long>()

    val elapsedNs = measureNanoTime {
        coroutineScope {
            val dispatcher = if (ioBound) Dispatchers.IO else Dispatchers.Default
            (0 until totalOps).map { i ->
                async(dispatcher) {
                    limiter.withPermit {
                        val t0 = System.nanoTime()
                        val r = if (ioBound) simulatedIo() else simulatedCpu()
                        latencies.add(System.nanoTime() - t0)
                        sink += r
                    }
                }
            }.awaitAll()
        }
    }

    val lats = latencies.toLongArray().apply { sort() }
    val qps = totalOps / (elapsedNs / 1e9)
    val p50 = lats[lats.size / 2] / 1e6
    val p99 = lats[(lats.size * 99) / 100] / 1e6
    return Triple(qps, p50, p99)
}

fun main() = runBlocking {
    val cores = Runtime.getRuntime().availableProcessors()
    println("CPU 核数 = $cores")
    println()

    for ((label, ioBound) in listOf("IO 密集(每次 5ms)" to true, "CPU 密集(每次 ~1ms)" to false)) {
        println("═══ $label ═══")
        println()
        val header = "%-8s%12s%12s%12s   %s".format("并行度", "QPS", "P50(ms)", "P99(ms)", "判读")
        println(header)
        println("-".repeat(64))

        var bestQps = 0.0
        for (p in listOf(1, 2, 4, 8, 16, 32, 64)) {
            // 跳过对 CPU 密集无意义的过大并行度
            if (!ioBound && p > cores * 4) continue

            val (qps, p50, p99) = measure(p, totalOps = if (ioBound) 2000 else 5000, ioBound = ioBound)

            val verdict = when {
                p == 1 -> "基线"
                qps > bestQps * 1.1 -> "✅ 仍在提升"
                qps > bestQps * 0.95 -> "⚠️  接近拐点"
                else -> "❌ 已过拐点(QPS 下降)"
            }
            if (qps > bestQps) bestQps = qps

            println("%-8d%12.0f%12.2f%12.2f   %s".format(p, qps, p50, p99, verdict))
        }
        println()
        if (ioBound) {
            println("IO 密集的判读:")
            println("  → QPS 随并行度上升,直到下游/调度器成为瓶颈")
            println("  → 注意 P99:并行度提高后 P99 通常上升(排队)")
            println("  → 选在「QPS 不再显著提升」之前,且 P99 可接受的位置")
        } else {
            println("CPU 密集的判读:")
            println("  → 最优值 ≈ CPU 核数(%d);超过后 QPS 下降".format(cores))
            println("  → 因为上下文切换成本超过并行收益")
        }
        println()
    }

    println("sink = $sink")
    println()
    println("⚠️  这是开发机上的简化实验,绝对值不可信;看【趋势与拐点位置】。")
}

预期输出形态:

═══ IO 密集(每次 5ms)═══

并行度           QPS     P50(ms)     P99(ms)   判读
----------------------------------------------------------------
1                199        5.01        5.12   基线
2                396        5.03        5.31   ✅ 仍在提升
4                782        5.05        5.68   ✅ 仍在提升
8               1520        5.08        6.21   ✅ 仍在提升
16              2410        5.11       11.42   ⚠️  接近拐点
32              2630        5.20       38.15   ⚠️  接近拐点
64              2680        5.34       95.20   ❌ 已过拐点(QPS 下降)

═══ CPU 密集(每次 ~1ms)═══

并行度           QPS     P50(ms)     P99(ms)   判读
----------------------------------------------------------------
1              1020        0.98        1.05   基线
2              2010        0.99        1.12   ✅ 仍在提升
4              3980        1.00        1.21   ✅ 仍在提升
8              7800        1.02        1.45   ✅ 仍在提升
16             7900        1.95        4.20   ⚠️  接近拐点
32             6100        5.20       12.80   ❌ 已过拐点(QPS 下降)

两个关键观察:

观察 结论
IO 密集:QPS 在 16–32 之间趋于饱和,但 P99 从 6 ms 涨到 95 ms 吞吐与延迟的权衡——“最优并行度"取决于你优化哪个指标
CPU 密集:超过核数(8)后 QPS 下降 超过核数的并行是负收益

二、隔板:为每个下游分配独立资源

// src/main/kotlin/resilience/Bulkhead.kt
package resilience

import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.TimeoutCancellationException
import kotlinx.coroutines.sync.Semaphore
import kotlinx.coroutines.sync.withPermit
import kotlinx.coroutines.withContext
import kotlinx.coroutines.withTimeout

/**
 * 隔板(Bulkhead):为每个下游分配独立的资源与调度器。
 *
 * 目的:一个下游变慢时,不会耗尽其他下游/其他接口的资源。
 */
class DownstreamBulkhead(
    val name: String,
    maxConcurrent: Int,
    private val timeoutMs: Long,
) {
    private val semaphore = Semaphore(maxConcurrent)

    // 每个下游独立的调度器(避免互相挤占 IO 池)
    private val dispatcher = Dispatchers.IO.limitedParallelism(maxConcurrent)

    /** 已拒绝的请求数(用于监控) */
    @Volatile var rejectedCount: Long = 0
        private set

    /** 当前等待中的请求数 */
    val waitingCount: Int get() = semaphore.availablePermits.let { maxConcurrent - it }

    suspend fun <T> execute(block: suspend () -> T): BulkheadResult<T> {
        // ① 快速失败(不等,直接拒绝)—— 可选策略
        if (semaphore.availablePermits == 0 && waitingCount > maxConcurrent * 2) {
            rejectedCount++
            return BulkheadResult.Rejected("$name 繁忙")
        }

        return try {
            // ② 限流 + 超时
            val result = withTimeout(timeoutMs) {
                semaphore.withPermit {
                    withContext(dispatcher) { block() }
                }
            }
            BulkheadResult.Success(result)
        } catch (e: TimeoutCancellationException) {
            rejectedCount++
            BulkheadResult.Timeout("$name 超时(${timeoutMs}ms)")
        } catch (e: Exception) {
            BulkheadResult.Failure(e)
        }
    }
}

sealed class BulkheadResult<out T> {
    data class Success<T>(val value: T) : BulkheadResult<T>()
    data class Rejected(val reason: String) : BulkheadResult<Nothing>()
    data class Timeout(val reason: String) : BulkheadResult<Nothing>()
    data class Failure(val error: Throwable) : BulkheadResult<Nothing>()
}

/**
 * 使用示例:每个下游一个独立的隔板
 */
class DownstreamClients {
    private val userService = DownstreamBulkhead("user-service", maxConcurrent = 32, timeoutMs = 100)
    private val orderService = DownstreamBulkhead("order-service", maxConcurrent = 32, timeoutMs = 150)
    private val paymentService = DownstreamBulkhead("payment-service", maxConcurrent = 16, timeoutMs = 300)

    suspend fun getUser(id: Long): String? =
        when (val r = userService.execute { httpGet("user", id) }) {
            is BulkheadResult.Success -> r.value
            else -> null                        // 降级:返回 null 或走兜底逻辑
        }

    suspend fun getOrder(id: Long): String? =
        when (val r = orderService.execute { httpGet("order", id) }) {
            is BulkheadResult.Success -> r.value
            else -> null
        }

    private suspend fun httpGet(service: String, id: Long): String {
        // 真实的 HTTP 调用
        return "$service-$id"
    }
}

三、验证隔板的效果

// src/test/kotlin/resilience/BulkheadTest.kt
package resilience

import kotlinx.coroutines.*
import kotlin.test.Test
import kotlin.test.assertTrue

class BulkheadTest {

    /**
     * 验证:一个下游变慢时,另一个下游不受影响。
     * 这是隔板的核心价值。
     */
    @Test
    fun `慢下游不应影响快下游`() = runBlocking {
        val slow = DownstreamBulkhead("slow", maxConcurrent = 4, timeoutMs = 5_000)
        val fast = DownstreamBulkhead("fast", maxConcurrent = 4, timeoutMs = 1_000)

        // 让 slow 占满(4 个协程各跑 2 秒)
        val slowJobs = List(4) {
            launch { slow.execute { delay(2_000); "slow" } }
        }
        delay(100)      // 等它们占满配额

        // 此时 fast 应该【完全不受影响】
        val fastLatencies = (1..20).map {
            async {
                val t0 = System.nanoTime()
                fast.execute { delay(10); "fast" }
                (System.nanoTime() - t0) / 1e6
            }
        }.awaitAll()

        slowJobs.forEach { it.cancel() }

        val maxFast = fastLatencies.max()
        println("fast 的最大延迟: %.1f ms(应该是 ~10-20ms,不受 slow 影响)".format(maxFast))
        assertTrue(maxFast < 100, "fast 被 slow 影响了:最大延迟 $maxFast ms")
    }

    /**
     * 验证:超过配额时快速拒绝(而不是无限等待)
     */
    @Test
    fun `超过配额应该快速失败`() = runBlocking {
        val bh = DownstreamBulkhead("limited", maxConcurrent = 2, timeoutMs = 100)

        // 占满配额
        val blockers = List(2) { launch { bh.execute { delay(500) } } }
        delay(50)

        // 后续请求应该很快超时(100ms)而不是等 500ms
        val t0 = System.nanoTime()
        val result = bh.execute { delay(500) }
        val elapsed = (System.nanoTime() - t0) / 1e6

        blockers.forEach { it.cancel() }

        println("超过配额时的耗时: %.1f ms".format(elapsed))
        assertTrue(result is BulkheadResult.Timeout, "应该超时,实际是 $result")
        assertTrue(elapsed < 200, "应该在 timeoutMs(100ms) 附近失败,实际 $elapsed ms")
    }
}

四、动手改造

改动 观察什么
在 ParallelismTuning 里把 IO 的 delay(5) 改成 delay(50) 最优并行度会往后移(Little’s Law)
把 CPU 密集的并行度扩到 cores * 4 能看到 QPS 明显下降——这就是「超过核数是负收益」
在 BulkheadTest 里把 slow 的 maxConcurrent 改成 100 fast 仍然不受影响(隔板生效);但如果去掉隔板(共享 IO 池),fast 会被拖慢
把 DownstreamBulkhead 的快速失败策略去掉 请求会排队到超时——体验变差(但不会丢请求)
给你的真实项目加上隔板 用「让一个下游变慢」的方式验证其他接口不受影响

五、这段代码的局限

  • ParallelismTuning 的数字来自开发机:看趋势与拐点位置,不看绝对值。
  • simulatedIo/simulatedCpu 是理想化的:真实操作有更复杂的分布(重尾),最优并行度会更难判断。
  • DownstreamBulkhead 的快速失败策略是自适应的(靠 availablePermits 判断),生产中可能需要更明确的拒绝阈值与监控。
  • 隔板会降低资源利用率:分配独立配额意味着某个下游空闲时,它的配额不能给别人用——这是为了隔离性付出的代价(典型的权衡)。
  • limitedParallelism 需要 Kotlin 1.6+;更早的版本要用自定义的 Executor。