文档目录

6.4 配套代码:线程与协程快照分析

对应小节:6.4 线程与协程快照 两件事:① 多份快照的对比分析;② 协程快照与协程泄漏检测。

一、多份快照的对比分析器

# tools/analyze-threads.py <threads-*.txt>...
"""
从多份线程快照里找出"持续存在"的问题。

核心思路:
  单张快照是瞬时状态,多张对比才能区分"路过"与"堵住"。
"""
import re
import sys
from collections import Counter, defaultdict


def parse_dump(path):
    """解析 jstack/jcmd Thread.print 输出"""
    threads = []
    current = None
    with open(path, errors="ignore") as f:
        for line in f:
            # 线程头:"name" #id prio=... tid=... nid=...
            m = re.match(r'^"([^"]+)"\s+#(\d+)\s+.*nid=(\S+)', line)
            if m:
                if current:
                    threads.append(current)
                current = {"name": m.group(1), "id": m.group(2),
                           "nid": m.group(3), "state": None, "stack": []}
                continue
            if current is not None:
                ms = re.search(r"java\.lang\.Thread\.State:\s*([A-Z_]+)", line)
                if ms:
                    current["state"] = ms.group(1)
                    continue
                mt = re.match(r"^\s+at\s+(\S+)", line)
                if mt:
                    current["stack"].append(mt.group(1))
    if current:
        threads.append(current)
    return threads


def main(paths):
    dumps = [parse_dump(p) for p in paths]
    if not dumps:
        print("❌ 没有读到快照")
        return

    print("═" * 84)
    print(f"线程快照对比分析({len(dumps)} 份)")
    print("═" * 84)
    print()

    # ① 每份的状态分布
    print("【① 状态分布(每份快照)】")
    for path, threads in zip(paths, dumps):
        states = Counter(t["state"] for t in threads if t["state"])
        total = sum(states.values())
        print(f"  {path}  线程数={len(threads)}")
        for st, c in states.most_common():
            pct = c / total * 100 if total else 0
            bar = "█" * int(pct / 3)
            print(f"     {st:<20}{c:>4}  {pct:5.1f}%  {bar}")
    print()

    # ② 持续存在的栈帧(在所有快照里都出现的)
    print("【② 持续存在的栈帧 Top 20】")
    print("    (在所有快照里都出现 → 不是「路过」,是「堵住」)")
    stack_sets = []
    for threads in dumps:
        frames = set()
        for t in threads:
            frames.update(t["stack"])
        stack_sets.append(frames)

    persistent = set.intersection(*stack_sets) if stack_sets else set()
    # 统计每个持续帧在多少线程里出现(取最后一份快照的数量)
    last_counts = Counter()
    for t in dumps[-1]:
        for f in t["stack"]:
            if f in persistent:
                last_counts[f] += 1

    for frame, _ in last_counts.most_common(20):
        # 在所有快照里的出现次数
        occurrences = sum(1 for ss in stack_sets if frame in ss)
        print(f"  [{occurrences}/{len(dumps)} 份]  {frame[:70]}")
    print()

    # ③ 关键模式检测
    print("【③ 关键模式检测】")
    patterns = {
        "连接池排队": ["HikariPool.getConnection", "getConnection"],
        "等网络 IO": ["socketRead0", "socketRead", "SocketInputStream"],
        "等锁(park)": ["Unsafe.park", "LockSupport.park"],
        "等队列/任务": ["LinkedBlockingQueue.take", "ArrayBlockingQueue"],
        "队列满(生产者阻塞)": ["ArrayBlockingQueue.put", "LinkedBlockingQueue.put"],
        "同步日志刷盘": ["FileOutputStream.write", "FileDispatcherImpl.write"],
        "主动 sleep": ["Thread.sleep"],
        "类加载锁": ["ClassLoader.loadClass", "defineClass"],
        "协程调度": ["DispatchedTask", "CoroutineScheduler"],
    }
    for label, keywords in patterns.items():
        count_by_dump = []
        for threads in dumps:
            c = sum(1 for t in threads
                    if any(any(k in f for k in keywords) for f in t["stack"]))
            count_by_dump.append(c)
        if any(count_by_dump):
            trend = "→".join(str(c) for c in count_by_dump)
            flag = " ⚠️ 持续" if all(c > 0 for c in count_by_dump) else ""
            print(f"  {label:<22}线程数: {trend}{flag}")
    print()

    # ④ BLOCKED 的详细信息
    blocked_in_last = [t for t in dumps[-1] if t["state"] == "BLOCKED"]
    if blocked_in_last:
        print(f"【④ BLOCKED 线程详情({len(blocked_in_last)} 个)】")
        for t in blocked_in_last[:5]:
            print(f"  {t['name']} (nid={t['nid']})")
            for f in t["stack"][:5]:
                print(f"      at {f}")
            print()
    else:
        print("【④ 无 BLOCKED 线程】→ 排除锁竞争")
    print()

    # ⑤ 线程数趋势(是否泄漏)
    counts = [len(d) for d in dumps]
    print(f"【⑤ 线程数趋势】{' → '.join(map(str, counts))}")
    if len(counts) >= 3 and counts[-1] > counts[0] * 1.5:
        print("  ❌ 线程数持续增长 → 疑似线程泄漏")
    else:
        print("  ✅ 线程数稳定")
    print()

    print("═" * 84)
    print("判读要点:")
    print("  ① 大量 BLOCKED → 锁竞争")
    print("  ② 大量 RUNNABLE 且 CPU 高 → 真的在算")
    print("  ③ 「持续存在」的栈才是真问题(单份快照里的可能只是路过)")
    print("  ④ 栈特征比状态更有指示性(HikariPool → 池排队;socketRead → 等下游)")
    print("  ⑤ RUNNABLE 线程数 == CPU 核数 且 CPU 低 → 疑似调度器被占满")
    print("═" * 84)


if __name__ == "__main__":
    if len(sys.argv) < 2:
        print("usage: analyze-threads.py <threads-1.txt> [threads-2.txt ...]")
        sys.exit(1)
    main(sys.argv[1:])

二、自动采集多份快照并分析

#!/usr/bin/env bash
# tools/capture-threads.sh <EXP_ID> [COUNT] [INTERVAL]
set -uo pipefail

EXP_ID="${1:?usage: capture-threads.sh <exp_id> [count] [interval]}"
COUNT="${2:-5}"
INTERVAL="${3:-3}"
DIR="docs/experiments/${EXP_ID}/results/threads"
mkdir -p "$DIR"

PID=$(jcmd 2>/dev/null | grep app.jar | awk '{print $1}')
[ -z "$PID" ] && { echo "❌ 找不到应用进程"; exit 1; }

echo "采集 $COUNT 份线程快照(间隔 ${INTERVAL}s,PID=$PID)"
for i in $(seq 1 "$COUNT"); do
  jcmd "$PID" Thread.print > "$DIR/threads-$i.txt" 2>/dev/null
  echo "  threads-$i.txt"
  [ "$i" -lt "$COUNT" ] && sleep "$INTERVAL"
done

echo
python3 tools/analyze-threads.py "$DIR"/threads-*.txt | tee "$DIR/analysis.txt"

三、协程快照与泄漏检测

3.1 协程 dump(排障用)

// src/main/kotlin/debug/CoroutineDebug.kt
package debug

import kotlinx.coroutines.debug.DebugProbes

/**
 * 协程调试工具。
 *
 * ⚠️ 开销很大(10%~30%),只能在排障时临时启用。
 *    生产环境请用下面的 CoroutineMonitor(开销低)。
 */
object CoroutineDebug {

    fun install() {
        DebugProbes.install()
    }

    /** dump 所有活跃协程及挂起点 */
    fun dump() {
        DebugProbes.dumpCoroutines()
    }

    /**
     * 只打印摘要(协程数量 + 挂起点分布),比 dumpCoroutines 开销小。
     */
    fun summary() {
        val coroutines = DebugProbes.dumpCoroutinesInfo()
        println("活跃协程数: ${coroutines.size}")

        // 按挂起点分组统计
        val byLocation = coroutines
            .mapNotNull { it.lastObservedStackTrace().firstOrNull() }
            .groupingBy { it.toString().substringBefore("(") }
            .eachCount()
            .entries
            .sortedByDescending { it.value }
            .take(10)

        println("挂起点分布(Top 10):")
        byLocation.forEach { (loc, count) ->
            val pct = count * 100.0 / coroutines.size
            println("  %6.1f%%  %4d 个  %s".format(pct, count, loc))
        }

        // 关键模式检测
        val stackTraces = coroutines.flatMap { it.lastObservedStackTrace() }.map { it.toString() }
        listOf(
            "连接池排队" to listOf("HikariPool", "getConnection"),
            "JDBC 阻塞" to listOf("jdbc", "QueryExecutor", "PgConnection"),
            "网络阻塞" to listOf("socketRead", "SocketInputStream"),
            "锁等待" to listOf("Mutex", "lock"),
            "队列等待" to listOf("Channel", "receive", "send"),
        ).forEach { (label, keywords) ->
            val c = stackTraces.count { st -> keywords.any { it in st } }
            if (c > 0) {
                println("  ⚠️  $label: $c 个协程(%.1f%%)".format(c * 100.0 / coroutines.size))
            }
        }
    }

    fun uninstall() {
        DebugProbes.uninstall()
    }
}
// build.gradle.kts
dependencies {
    implementation("org.jetbrains.kotlinx:kotlinx-coroutines-debug:1.9.0")
}

3.2 低开销的协程监控(生产可用)

// src/main/kotlin/metrics/CoroutineMonitor.kt
package metrics

import io.micrometer.core.instrument.Gauge
import io.micrometer.core.instrument.MeterRegistry
import kotlinx.coroutines.Job
import kotlinx.coroutines.launch
import kotlinx.coroutines.CoroutineScope
import java.util.concurrent.atomic.AtomicInteger

/**
 * 低开销的协程监控:只统计数量,不采集栈。
 * 可以长期开启(生产环境可用)。
 */
class CoroutineMonitor(registry: MeterRegistry) {

    private val active = AtomicInteger(0)
    private val maxObserved = AtomicInteger(0)

    init {
        Gauge.builder("app.coroutines.active") { active.get().toDouble() }
            .description("当前活跃协程数")
            .register(registry)

        Gauge.builder("app.coroutines.max_observed") { maxObserved.get().toDouble() }
            .description("观测到的最大活跃协程数(用于判断是否堆积)")
            .register(registry)
    }

    /** 包装一个协程,统计它的活跃状态 */
    fun CoroutineScope.tracked(block: suspend () -> Unit): Job = launch {
        val n = active.incrementAndGet()
        maxObserved.updateAndGet { maxOf(it, n) }
        try {
            block()
        } finally {
            active.decrementAndGet()
        }
    }

    val activeCount: Int get() = active.get()
}

判读规则:

观察 结论
app.coroutines.active 持续增长 协程泄漏(GlobalScope 未取消、Effect 未释放)
active 稳定但很高 正常(高并发)
active 与 QPS 成正比 正常
active 增长但 QPS 未变 有协程没结束(泄漏)

3.3 在 CI 里检测协程泄漏

// src/test/kotlin/CoroutineLeakTest.kt
package test

import kotlinx.coroutines.*
import kotlinx.coroutines.test.*
import org.junit.jupiter.api.Test
import kotlin.test.assertEquals

class CoroutineLeakTest {

    /**
     * 用 runTest 检测协程泄漏。
     * 如果测试结束时还有未完成的协程,runTest 会失败。
     */
    @Test
    fun `处理请求后不应该留下活跃协程`() = runTest {
        val monitor = metrics.CoroutineMonitor(io.micrometer.core.instrument.simple.SimpleMeterRegistry())

        repeat(100) {
            monitor.tracked { delay(10) }
        }
        advanceUntilIdle()

        assertEquals(0, monitor.activeCount, "有 ${monitor.activeCount} 个协程未结束")
    }
}

这是防止协程泄漏最有效的手段——把检测放进 CI,而不是等线上内存涨了才发现。

四、动手改造

改动 观察什么
用 analyze-threads.py 分析一个正常服务的快照 大部分线程应该在 WAITING (parking)(空闲的池线程)——这是健康的
故意在某个接口加 Thread.sleep 并施压 快照里会出现 Thread.sleep 的持续栈
故意制造锁竞争 出现大量 BLOCKED,且 analyze-threads.py 会列出详情
把 CoroutineMonitor 接进你的项目 得到一个低开销的协程数指标
把协程泄漏测试加进 CI 从流程上防止泄漏
对比「协程 dump」与「线程 dump」 体会为什么协程应用必须补协程快照

五、这段代码的局限

  • analyze-threads.py 的解析依赖 jstack 输出格式:不同 JDK 版本格式略有差异,可能需要调整正则。
  • 「持续存在」的判定需要足够多份快照:3 份是下限,5 份更可靠;间隔太短可能看不出变化。
  • 栈帧名字是 Java 签名(带包路径),可读性一般——可以用 c++filt 类似的方式做简化,或直接看 HTML 化的快照工具。
  • DebugProbes 开销 10%–30%:绝不能在基线组和实验组之间不一致地使用(第 5 章 5.1 节)。
  • CoroutineMonitor 只能统计数量:不能告诉你协程卡在哪——需要时还要用 DebugProbes 做一次深度 dump。