2.7 配套代码:协程开销与调度器饥饿实验
对应小节:2.7 协程的真实开销模型 第二个实验是本章最重要的一个:它让你亲眼看到「阻塞调用污染调度器」造成的全局卡顿。
一、实验一:协程的成本
// src/main/kotlin/lesson02/CoroutineCost.kt
package lesson02
import kotlinx.coroutines.*
import kotlin.system.measureNanoTime
private var sink = 0L
suspend fun noopSuspend(): Int = 1 // 不真正挂起
suspend fun realSuspend(): Int = withContext(Dispatchers.IO) { 1 } // 真正切换线程
fun main() = runBlocking {
val n = 200_000
fun bench(label: String, body: () -> Unit): Double {
repeat(n / 4) { body() }
val ns = measureNanoTime { repeat(n) { body() } }
val per = ns.toDouble() / n
println("%-44s %10.1f ns/op".format(label, per))
return per
}
println("=== 协程相关操作的成本 ===")
val plain = bench("普通函数调用") { sink += 1 }
val suspendNoop = bench("suspend 函数(不真正挂起)") { sink += runBlocking { noopSuspend() } }
val dispatch = bench("真正挂起:withContext(IO)") { sink += runBlocking { realSuspend() } }
val thread = bench("线程创建 + join") {
val t = Thread { }
t.start(); t.join()
}
println()
println("=== 结论 ===")
println("suspend(不挂起)vs 普通调用 : %.1f 倍".format(suspendNoop / plain))
println("真正挂起 vs 普通调用 : %.1f 倍".format(dispatch / plain))
println("线程创建 vs 真正挂起 : %.0f 倍".format(thread / dispatch))
println()
println("注意:这里的 runBlocking 本身也有开销(它是被测的一部分)。")
println(" 严格的测量需要把 runBlocking 的空壳作为对照组单独测。")
println()
println("sink = $sink")
}
预期形态:
| 操作 | 量级 |
|---|---|
| 普通函数调用 | ~1 ns |
suspend 不真正挂起 |
几 ns(接近普通调用) |
真正挂起(withContext(IO)) |
微秒级(涉及线程调度) |
| 线程创建 + join | 几十微秒 |
注意
withContext(Dispatchers.IO)比「协程挂起」贵得多——因为它包含了线程切换。纯挂起(在同一线程上恢复,比如delay()的实现路径)要便宜得多。这个区分很重要:协程的便宜指的是「挂起/恢复」,不是「跨线程调度」。
二、实验二:调度器饥饿(本章最重要)
// src/main/kotlin/lesson02/DispatcherStarvation.kt
package lesson02
import kotlinx.coroutines.*
import kotlin.system.measureNanoTime
private var sink = 0L
/** 模拟一次"快速工作"(CPU 计算),并返回它的耗时(毫秒) */
suspend fun fastWork(): Double {
val ns = measureNanoTime {
var s = 0L
repeat(10_000) { s += it }
sink += s
}
return ns / 1_000_000.0
}
/** 在指定调度器上并发跑 n 个快速工作,返回各自耗时 */
suspend fun runFastWorks(n: Int): List<Double> = coroutineScope {
(0 until n).map { async(Dispatchers.Default) { fastWork() } }.awaitAll()
}
fun median(values: List<Double>): Double {
val sorted = values.sorted()
return sorted[sorted.size / 2]
}
fun main() = runBlocking {
val cores = Runtime.getRuntime().availableProcessors()
println("CPU 核数 = $cores")
println("Dispatchers.Default 并行度 ≈ max(2, 核数) = $cores")
println("Dispatchers.IO 并行度 ≈ max(64, 核数)")
println()
// ── 阶段 1:基线(只有快速工作) ──────────────────────────
runFastWorks(cores * 4) // 预热
val baseline = runFastWorks(cores * 4)
println("阶段 1|基线(只有快速工作) : 中位 %8.2f ms".format(median(baseline)))
// ── 阶段 2:阻塞调用放在 Dispatchers.Default ──────────────
val blockersOnDefault = List(cores) {
launch(Dispatchers.Default) {
Thread.sleep(1_000) // ❌ 阻塞调用!占住 Default 的线程
}
}
delay(100) // 给阻塞协程一点时间占住线程
val polluted = runFastWorks(cores * 4)
blockersOnDefault.forEach { it.cancel() }
println("阶段 2|阻塞在 Default 上 : 中位 %8.2f ms ← 全局卡顿!".format(median(polluted)))
delay(300) // 等调度器恢复
// ── 阶段 3:阻塞调用放到 Dispatchers.IO ───────────────────
val blockersOnIo = List(cores) {
launch(Dispatchers.IO) {
Thread.sleep(1_000) // ✅ 放到允许阻塞的调度器
}
}
delay(100)
val healthy = runFastWorks(cores * 4)
blockersOnIo.forEach { it.cancel() }
println("阶段 3|阻塞在 IO 上 : 中位 %8.2f ms ← 恢复正常".format(median(healthy)))
println()
println("=== 结论 ===")
println("阻塞污染的放大倍数 : %.0f 倍".format(median(polluted) / median(baseline)))
println("切到 IO 后的恢复 : %.2f 倍(相对基线)".format(median(healthy) / median(baseline)))
println()
println("同时观察 CPU:跑这个实验时 CPU 利用率很低,")
println("但『快速工作』的延迟涨了几百倍 —— 这就是这类故障最难排查的原因。")
}
三、预期输出
CPU 核数 = 8
Dispatchers.Default 并行度 ≈ max(2, 核数) = 8
Dispatchers.IO 并行度 ≈ max(64, 核数)
阶段 1|基线(只有快速工作) : 中位 0.08 ms
阶段 2|阻塞在 Default 上 : 中位 892.31 ms ← 全局卡顿!
阶段 3|阻塞在 IO 上 : 中位 0.09 ms ← 恢复正常
=== 结论 ===
阻塞污染的放大倍数 : 11154 倍
切到 IO 后的恢复 : 1.13 倍(相对基线)
同时观察 CPU:跑这个实验时 CPU 利用率很低,
但『快速工作』的延迟涨了几百倍 —— 这就是这类故障最难排查的原因。
这个实验揭示了三件事:
| 观察 | 含义 |
|---|---|
| 阶段 2 的延迟从 0.08 ms 涨到近 900 ms | 8 个阻塞协程占满了 8 个 Default 线程,所有快速工作只能排队 |
| 阶段 3 恢复正常 | 换到 IO 调度器后,Default 不再被占用 |
| 期间 CPU 利用率很低 | 这就是为什么这个故障难以排查——所有常规指标都正常 |
四、观察实验期间的 CPU
# 另开一个终端,观察 Java 进程的 CPU
pidstat -u -p $(jcmd | grep lesson02 | awk '{print $1}') 1
# 或者
top -pid $(jcmd | grep lesson02 | awk '{print $1}')
你会看到:CPU 只有个位数百分比,但「快速工作」的延迟暴涨。「CPU 不高但很慢」是等待类瓶颈的典型特征(第 0 章 0.1 节讲过:这时要用 wall-clock 火焰图,而不是 CPU 火焰图)。
五、在你的项目里怎么排查这个问题
代码审查:搜索这些模式
# 在代码里搜索可疑的阻塞调用
grep -rn "Dispatchers.Default" src/ | grep -v "limitedParallelism"
grep -rn "Thread.sleep" src/
grep -rn "runBlocking" src/
grep -rn "jdbcTemplate\|DriverManager\|\.executeQuery" src/ | grep -v "withContext"
运行时诊断
// 仅在排障时启用,开销很大
import kotlinx.coroutines.debug.DebugProbes
DebugProbes.install()
DebugProbes.dumpCoroutines() // 打印所有活跃协程及挂起点
DebugProbes.uninstall()
// build.gradle.kts
dependencies {
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-debug:1.9.0")
}
判读:如果 dump 出大量协程挂在同一个阻塞调用栈上,且数量 ≈ Default 的并行度,就坐实了这个问题。
正确的写法模板
// ✅ 模板:阻塞调用必须切到 IO,并限制对下游的并发
private val dbLimiter = Semaphore(permits = 64)
suspend fun queryOrder(id: String): Order = dbLimiter.withPermit {
withContext(Dispatchers.IO) {
jdbcTemplate.queryForObject("select ... where id = ?", id.toLong())
}
}
// ✅ 更好:用真正的异步驱动,根本不阻塞
suspend fun queryOrder(id: String): Order =
r2dbcClient.sql("select ... where id = :id").bind("id", id).awaitOne()
六、动手改造
| 改动 | 观察什么 |
|---|---|
把阻塞协程数量从 cores 改成 cores / 2 |
还有空闲线程,卡顿会小很多——说明必须占满才饿死 |
把 Thread.sleep(1000) 改成 delay(1000) |
完全没有卡顿——因为 delay 是真挂起,会释放线程 |
把阶段 2 的阻塞放进 runBlocking |
可能观察不到卡顿(因为它阻塞的是调用线程,不是 Default worker) |
把阶段 3 的 Dispatchers.IO 换成 Dispatchers.Default.limitedParallelism(64) |
也能避免污染,但要注意总并发上限 |
在 fastWork 里加上 Thread.sleep(1) |
快速工作自己也变成阻塞的了——两个问题叠加 |
七、这段代码的局限
Thread.sleep是简化的阻塞调用。真实场景是 JDBC、同步 HTTP 客户端、文件 IO——机制相同,但耗时分布不同。runBlocking本身会引入开销,实验一的绝对值不可信,只看相对倍数。- 阶段 2 的延迟受
delay(100)的影响:如果阻塞协程还没完全占住线程,放大倍数会被低估。 - 协程调度器的实现细节随版本变化(例如
Dispatchers.IO的默认并行度),请以你的 kotlinx-coroutines 版本文档为准。