CoroutineScheduler CPU 许可证系统详细实现原理

  1. 核心数据结构

controlState 的位布局

⏺ controlState 是一个 64 位原子长整型,包含三个字段:

// 初始化:CPU许可证数量 = corePoolSize
private val controlState = atomic(corePoolSize.toLong() shl
CPU_PERMITS_SHIFT)

// 位掩码定义
private const val BLOCKING_SHIFT = 21
private const val CPU_PERMITS_SHIFT = BLOCKING_SHIFT * 2 // = 42
private const val CREATED_MASK: Long = (1L shl BLOCKING_SHIFT) -1
private const val CPU_PERMITS_MASK = CREATED_MASK shl
CPU_PERMITS_SHIFT

位布局图示:
63 42 21 0
|---------|--------|---------|
| CPU | BLK | CRT |

CPU 许可证计数
BLK 阻塞任务计数
CRT 已创建线程数

  • [0-20] 位:已创建线程数 (createdWorkers)
  • [21-41] 位:正在执行的阻塞任务数 (blockingTasks)
  • [42-62] 位:可用 CPU 许可证数 (availableCpuPermits)

⏺ 2. 许可证的获取和释放机制

tryAcquireCpuPermit() - 调度器级别

⏺ 调度器级别的许可证操作

// 尝试获取CPU许可证(调度器级别)
private inline fun tryAcquireCpuPermit(): Boolean =
controlState.loop { state ->
val available = availableCpuPermits(state)
if (available == 0) return false // 没有可用许可证
val update = state - (1L shl CPU_PERMITS_SHIFT) // 减1
if (controlState.compareAndSet(state, update)) return true
}

// 释放CPU许可证
private inline fun releaseCpuPermit() =
controlState.addAndGet(1L shl CPU_PERMITS_SHIFT)

Worker 级别的许可证操作

⏺ Worker 级别的许可证管理:

// Worker 尝试获取CPU许可证
private fun tryAcquireCpuPermit(): Boolean = when {
state == WorkerState.CPU_ACQUIRED -> true // 已经有许可证
this@CoroutineScheduler.tryAcquireCpuPermit() -> {
state = WorkerState.CPU_ACQUIRED // 获取成功,更新状态
true
}
else -> false // 获取失败
}

// Worker 释放CPU许可证并切换状态
fun tryReleaseCpu(newState: WorkerState): Boolean {
val previousState = state
val hadCpu = previousState == WorkerState.CPU_ACQUIRED
if (hadCpu) releaseCpuPermit() // 调用调度器的释放方法
if (previousState != newState) state = newState
return hadCpu // 返回是否确实释放了许可证
}

关键特性

  1. 原子操作:使用 CAS 操作确保多线程安全
  2. 状态同步:Worker 状态与许可证持有状态保持一致
  3. 无锁设计:通过循环重试避免锁竞争

⏺ 3. 阻塞任务执行时的许可证转移

beforeTask() - 任务执行前的许可证交出

⏺ 阻塞任务的许可证转移流程

执行前 - beforeTask():
private fun beforeTask(taskMode: Int) {
if (taskMode == TASK_NON_BLOCKING) return
// 执行阻塞任务时总是通知有新工作,释放CPU许可证
if (tryReleaseCpu(WorkerState.BLOCKING)) {
signalCpuWork() // 立即信号通知其他线程有CPU工作可做
}
}

执行后 - afterTask():
private fun afterTask(taskMode: Int) {
if (taskMode == TASK_NON_BLOCKING) return
decrementBlockingTasks() // 减少阻塞任务计数
val currentState = state
if (currentState !== WorkerState.TERMINATED) {
assert { currentState == WorkerState.BLOCKING }
state = WorkerState.DORMANT // 变为休眠状态
}
}

关键的许可证转移机制

⏺ signalCpuWork() 的作用:
fun signalCpuWork() {
if (tryUnpark()) return // 1. 尝试唤醒已停泊的线程
if (tryCreateWorker()) return // 2. 尝试创建新的工作线程
tryUnpark() // 3. 再次尝试唤醒(防止竞争)
}

完整的许可证转移流程

  1. Worker A 获得阻塞任务,调用
    beforeTask(TASK_PROBABLY_BLOCKING)
  2. 许可证释放:tryReleaseCpu(WorkerState.BLOCKING)
    释放许可证到池中
  3. 立即通知:signalCpuWork() 立即寻找其他线程来利用释放的许可证
  4. 状态转换:Worker A 从 CPU_ACQUIRED → BLOCKING
  5. 执行阻塞任务:Worker A 执行可能阻塞的操作(IO、锁等)
  6. 任务完成:调用 afterTask() 清理状态,变为 DORMANT

关键优势:

  • 即时转移:许可证在任务执行前就释放,不等任务完成
  • 主动通知:释放许可证后主动寻找可以利用它的线程
  • 无饥饿:确保始终有 corePoolSize 个线程可以执行 CPU 任务

⏺ 4. 许可证系统如何防止 CPU 任务饥饿

findTask() 中的许可证检查

⏺ 任务查找的双轨制设计

findTask() 的核心逻辑:
fun findTask(scanLocalQueue: Boolean): Task? {
if (tryAcquireCpuPermit()) return
findAnyTask(scanLocalQueue)
// 无法获取CPU许可证 -- 只能查找阻塞任务
val task = if (scanLocalQueue) {
localQueue.poll() ?:
globalBlockingQueue.removeFirstOrNull()
} else {
globalBlockingQueue.removeFirstOrNull()
}
return task ?: trySteal(blockingOnly = true) //
只窃取阻塞任务
}

两条执行路径:

  1. 有 CPU 许可证:可以执行任何类型的任务(CPU + 阻塞)
  2. 无 CPU 许可证:只能执行阻塞任务

工作窃取中的限制

⏺ trySteal() 中的许可证限制:
val stealResult = if (blockingOnly) {
localQueue.tryStealBlockingFrom(victim = worker.localQueue)
// 只窃取阻塞任务
} else {
localQueue.tryStealFrom(victim = worker.localQueue)
// 可窃取任何任务
}

CPU 任务饥饿的防护机制

核心不变式:
调度器始终保持至少 min(待处理CPU任务数, corePoolSize) 且最多
corePoolSize 个线程来执行 CPU 任务

实现机制:

  1. 许可证数量恒定:始终有 corePoolSize 个 CPU 许可证
  2. 阻塞任务隔离:无许可证的线程无法执行 CPU 任务
  3. 即时转移:阻塞任务执行前立即释放许可证
  4. 主动唤醒:释放许可证后主动寻找可利用的线程

饥饿防护的数学原理

假设场景:有 4 个核心线程,当前所有线程都在执行阻塞任务

传统线程池问题:

  • 4 个线程全部阻塞 → CPU 任务无法执行 → 饥饿

CPU 许可证解决方案:

  1. 每个线程执行阻塞任务前释放许可证 → 4 个许可证回到池中
  2. signalCpuWork() 唤醒线程或创建新线程来获取许可证
  3. 新线程获取许可证后可以执行 CPU 任务
  4. 系统保持 4 个 CPU 许可证的活跃状态

结果:无论有多少阻塞任务,始终有 corePoolSize 个"CPU
工作槽位"可用

⏺ 5. 许可证与线程创建的协调机制

tryCreateWorker() 中的智能判断

⏺ 动态线程创建逻辑

核心计算公式:
val created = createdWorkers(state) // 总线程数
val blocking = blockingTasks(state) //
正在执行阻塞任务的线程数
val cpuWorkers = (created - blocking).coerceAtLeast(0) //
可用于CPU工作的线程数

线程创建判断:
if (cpuWorkers < corePoolSize) {
val newCpuWorkers = createNewWorker()
//
特殊优化:如果创建了第一个CPU工作线程且核心池大小>1,创建第二个
if (newCpuWorkers == 1 && corePoolSize > 1)
createNewWorker()
if (newCpuWorkers > 0) return true
}

许可证与线程创建的协调机制详解

场景分析:

场景 1:CPU 任务到达
初始状态:corePoolSize=4, created=2, blocking=0, cpuWorkers=2,
availablePermits=4
CPU任务到达 → signalCpuWork()
→ tryUnpark() 失败(没有停泊线程)
→ tryCreateWorker(): cpuWorkers(2) < corePoolSize(4) →
创建新线程
→ 新线程启动后可以获取CPU许可证执行任务

场景 2:阻塞任务执行
当前状态:created=4, blocking=0, cpuWorkers=4
阻塞任务执行:blocking=1, cpuWorkers=3
beforeTask() → releasesCpuPermit() → signalCpuWork()
→ tryCreateWorker(): cpuWorkers(3) < corePoolSize(4) →
尝试创建新线程
→ 如果 created < maxPoolSize,创建成功,新线程获取许可证

关键的协调原则

⏺ signalBlockingWork() 的智能处理:
private fun signalBlockingWork(skipUnpark: Boolean) {
val stateSnapshot = incrementBlockingTasks() //
先增加阻塞任务计数
if (skipUnpark) return
if (tryUnpark()) return //
优先唤醒现有线程
if (tryCreateWorker(stateSnapshot)) return //
使用快照避免过度创建
tryUnpark() //
防止竞争条件的二次尝试
}

线程创建的限制条件

  1. CPU 工作线程不足:cpuWorkers < corePoolSize
  2. 总线程数限制:created < maxPoolSize
  3. 调度器未终止:!isTerminated

特殊优化:第二个线程的创建

// 如果创建了第一个CPU工作线程且corePoolSize > 1,创建第二个
if (newCpuWorkers == 1 && corePoolSize > 1) createNewWorker()

设计原因:

  • 单个线程无法进行工作窃取
  • 第二个线程使工作窃取算法生效
  • 提升整体并发性能

⏺ CPU 许可证系统总结

系统设计的精妙之处

CPU 许可证系统是 CoroutineScheduler
最核心的创新,它通过以下机制解决了传统线程池的根本问题:

  1. 核心问题解决

传统问题:阻塞任务占用所有线程 → CPU 任务饥饿 → 系统响应性下降

许可证解决方案:

  • 将线程资源与 CPU 执行权限分离
  • 阻塞任务执行前主动释放 CPU 权限
  • 保证始终有足够的 CPU 执行槽位
  1. 实现机制概览

许可证池: [CPU_PERMIT_1] [CPU_PERMIT_2] [CPU_PERMIT_3]
[CPU_PERMIT_4]
线程池: [Worker_1] [Worker_2] [Worker_3] [Worker_4] [Worker_5]
[Worker_6]

场景: Worker_1, Worker_2 执行阻塞任务
结果: 许可证1,2 被释放回池中
Worker_3,4,5,6 可以争抢许可证执行 CPU 任务

  1. 关键设计原则

  2. 许可证数量恒定:始终等于 corePoolSize

  3. 即时转移:阻塞任务执行前立即释放许可证

  4. 主动信号:释放后立即寻找可用线程

  5. 分离执行:有许可证=执行任何任务,无许可证=只能执行阻塞任务

  6. 动态扩展:线程数可以超过 corePoolSize,但 CPU 许可证数量固定

  7. 性能优势

  • 无饥饿:CPU 任务永远不会被完全阻塞
  • 高吞吐:阻塞任务不影响 CPU 任务处理
  • 低延迟:许可证转移是 O(1) 原子操作
  • 自适应:根据负载动态调整线程数量

这个系统展现了现代并发编程中资源管理的高级技巧,通过抽象化的"许
可证"概念,实现了复杂的线程池行为控制

©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

友情链接更多精彩内容