文档目录

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 版本文档为准。