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。