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。