位置:首页 > Kotlin > Kotlin协程调度细节与取消机制闭环深度解析

Kotlin协程调度细节与取消机制闭环深度解析

时间:2026-08-15  |  作者:极客少年  |  阅读:0

Kotlin 协程的官方框架 kotlinx.coroutines 主要由以下几个部分构成:

深入理解 Kotlin 协程 (八):拾遗补阙,探秘官方框架的调度细节与取消闭环
            <!----></a> <!---->

  • core: 框架的核心逻辑,包括了复合协程以及 Channel、Flow 等特性。

  • ui: 包括 android、ja vafx、swing 库,用于提供各个平台的 UI 调度器以及特有的逻辑。

  • reactive 相关: 提供了对各种响应式编程框架的协程支持,如 rx2 提供了对 RxJa va 2.x 的协程支持。

  • integration 相关: 提供了其他框架中异步回调的集成,如 jdk8 集成中新增了 CompletableFuture 协程 API。

  • test: 提供了测试模块,用于控制协程的虚拟时间、测试挂起函数等,对于编写单元测试来说必不可少。

下面继续看看前文未涉及到的一些官方细节。

协程的启动模式(CoroutineStart)

在官方的协程构建器中,还可以传入一个启动模式:start: CoroutineStart

public fun CoroutineScope.launch(
    context: CoroutineContext = EmptyCoroutineContext,
    start: CoroutineStart = CoroutineStart.DEFAULT, // 启动模式
    block: suspend CoroutineScope.() -> Unit
): Job {
    val newContext = newCoroutineContext(context)
    val coroutine = if (start.isLazy)
        LazyStandaloneCoroutine(newContext, block) else
        StandaloneCoroutine(newContext, active = true)
    coroutine.start(start, coroutine, block)
    return coroutine
}

四种启动模式

  • DEFAULT:创建协程后,立即开始调度。如果在调度前协程被取消,将进入取消响应状态。

  • ATOMIC:也是创建后立即调度。不过在执行到第一个挂起点之前不响应取消。

  • LAZY:只有协程被需要时,才会开始调度。这里的需要包括主动调用 startjoin 或者 await 函数。如果调度前被取消,协程将进入异常结束状态。

  • UNDISPATCHED:协程创建后会立即在当前函数调用栈中执行,直到遇到第一个真正挂起点。

注意:立即调度不等于立即执行。

立即调度表示调度器会马上执行调度动作,但协程具体何时真正执行,仍由调度器决定。就像老妈安排你去扫地,任务已经下发了,但什么时候扫,由你决定。

UNDISPATCHEDATOMIC 都一定会执行,但二者有区别。

  • UNDISPATCHED:到第一个挂起点之前,少了一次线程调度。它会直接在创建协程所在的调用者线程上立即执行。

  • ATOMIC:依然会被调度器分发到指定线程上执行,只是在遇到第一个挂起点之前对取消状态“免疫”。

总之,这些启动模式更多是为了应对特殊场景。

在业务开发中,通常使用 DEFAULTLAZY 就够了。

虽然我们之前实现的协程没有启动模式,但效果等同于 ATOMIC 模式。

官方提供的调度器(Dispatchers)

官方框架中提前准备好了四个调度器,我们可以通过 Dispatchers 单例来访问。

  • Default:默认的调度器,适合后台计算任务。

  • IO:IO 调度器,适合执行 IO 操作。

  • Main:UI 调度器。平台不同,UI 线程的调度器也不同。在 Android 上,会将协程调度到 UI 事件循环(主线程)中执行。

注意:除了引入核心依赖外,还要引入 Android 平台的依赖,否则在使用 Dispatchers.Main 调度器时会抛出异常:

implementation("org.jetbrains.kotlinx:kotlinx-coroutines-core:1.10.2") // 核心库
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-android:1.10.2") // Android 平台依赖
  • Unconfined:不受限制的调度器。也就是说,协程执行在哪个线程上无所谓。本质上来说,这个相当于调度器为空,在哪恢复,就在恢复的线程上执行。但嵌套 Unconfined 协程时,会进行特殊处理:放到协程框架内部的事件循环上,防止栈溢出。

Unconfined 的示例

下面用一个例子说明 Unconfined 调度器:

runBlocking {
    // runBlocking 在主线程中执行
    launch(Dispatchers.Unconfined) {
        println(Thread.currentThread().name) // 启动时,协程处于主线程
        
        delay(100) // 挂起
        
        // 等外部恢复执行时,协程会在唤醒当前协程的线程上执行(不是主线程)
        println(Thread.currentThread().name) 
    }
}

调度器的使用场景

  • 如果涉及 UI 操作,就必须使用 Main

  • 如果只是普通的后台任务,应该使用 Default

  • 如果涉及 IO 操作,比如网络或文件,就使用 IO

其实 DefaultIO 调度器的背后都是线程池,只是配置不同。

更多关于 Default 和 IO 的区别和使用场景,可以看这篇博客:Kotlin 协程中的 IO 与 Default 调度器详解 | Baeldung中文网

自定义调度器

如果需要自定义调度器,可以通过继承 CoroutineDispatcher 抽象类来完成。

不过更常见的做法,是将一个已经存在的线程池转换成调度器:

Executors
    .newSingleThreadExecutor()
    .asCoroutineDispatcher()
    .use { dispatcher ->
        val result = GlobalScope.async(dispatcher) {
            delay(timeMillis = 100)
            123456
        }.await()
    } 

注意:这个转换而来的调度器必须主动关闭,以免造成线程泄漏。

另外,使用 withContext 也可以很方便地切换调度器。

它的作用等价于 async{ ... }.await()。如果调用 async 后立即调用了 await,可以使用 withContext 替代,因为其内存开销更低。

协程的取消响应点

我们知道,自定义的挂起函数可以通过 suspendCancellableCoroutine 来响应取消,例如:

// 异步定时任务
suspend fun delayWithCallback(delayMs: Long): String = suspendCancellableCoroutine { cont ->
    // 创建定时器
    val timer = Timer()
    timer.schedule(object : TimerTask() {
        override fun run() {
            // 回调执行,先判断协程是否已取消
            if (cont.isActive) {
                cont.resume("定时 $delayMs ms 执行完成") { cause, value, context -> }
            }
            timer.cancel()
        }
    }, delayMs)

    // 监听协程取消
    cont.invokeOnCancellation {
        // 销毁定时器
        println("协程取消,关闭定时器")
        timer.cancel()
    }
}

另外,自带的挂起函数,如 delaywithContext 等,也会响应取消。

所以协程通常默认只会在这些挂起点检查取消并退出。

没有挂起点时怎么办

如果没有任何挂起点,例如一个文件复制的函数:

这是标准库提供的扩展函数。

public fun InputStream.copyTo(out: OutputStream, bufferSize: Int = DEFAULT_BUFFER_SIZE): Long {
    var bytesCopied: Long = 0
    val buffer = ByteArray(bufferSize)
    var bytes = read(buffer)
    while (bytes >= 0) {
        out.write(buffer, 0, bytes)
        bytesCopied += bytes
        bytes = read(buffer)
    }
    return bytesCopied
}

这时,我们可以使用 suspendCancellableCoroutine 包装这段逻辑。

也可以手动轮询判断 isActive 标志:

// 可取消的同步流复制函数
@OptIn(InternalCoroutinesApi::class)
suspend fun InputStream.copyToSuspend(
    out: OutputStream,
    bufferSize: Int = DEFAULT_BUFFER_SIZE
): Long {
    var bytesCopied: Long = 0
    val buffer = ByteArray(bufferSize)
    var bytes = read(buffer)
    val job = currentCoroutineContext()[Job]
    while (bytes >= 0) {
        // 检查协程是否取消
        job.let {
            if (!it.isActive) {
                throw job.getCancellationException() // 内部 API
            }
        }
        out.write(buffer, 0, bytes)
        bytesCopied += bytes
        bytes = read(buffer)
    }
    return bytesCopied
}

如果 Job 为空,说明当前所在的是一个简单协程,没有实现取消逻辑。

上面这段判断逻辑,也可以使用 ensureActive 扩展函数代替:

public fun Job.ensureActive(): Unit {
    if (!isActive) throw getCancellationException()
}

如果不希望影响内部逻辑,可以使用 yield 挂起函数。

它内部就有取消响应点,而且就在函数开头:

val context = uCont.context
context.ensureActive()
...

此外,yield 还会尝试让出线程的执行权,让其他协程获得执行机会。

不过要注意性能差异。

  • yield 的主要职责是让出线程,这会触发底层的上下文切换和重新排队,性能开销较大。

  • 如果只是为了在密集计算或循环中响应取消,使用 ensureActive() 才是最优解。

异步任务的超时控制

对于异步任务的超时取消,我们可以同步发射一个协程,到时间后取消异步任务。

不过这样有些冗余,因此可以直接使用官方提供的超时函数 withTimeout

runBlocking {
    val time = withTimeout(timeMillis = 5000L) {
        val delayTime = (4000L..6000L).random()
        delay(timeMillis = delayTime)
        val formatter =
            DateTimeFormatter.ofPattern("HH:mm:ss")
        LocalDateTime.now().format(formatter)
    }
    println(time)
}

withTimeout 在超时后,会取消 block 代码块的执行,并且抛出一个取消异常。

如果不希望在超时时抛出取消异常,可以使用 withTimeoutOrNull。它在超时后会返回 null

禁止取消(NonCancellable 上下文)

最后来看一下禁止取消。

举个例子,我们希望使用 delay 模拟耗时任务,同时又不希望它响应外部取消。例如,我们想观察其他挂起函数的取消效果:

runBlocking {
    val job = launch {
        listOf(1, 2, 3, 4).forEach {
            yield()
            delay(timeMillis = it * 100L)
        }
    }
    delay(timeMillis = 200L)
    job.cancelAndJoin()
}

这时,可以使用 NonCancellable 上下文。

它能够禁止作用范围内的取消响应:

GlobalScope.launch {
    val job = launch {
        listOf(1, 2, 3, 4).forEach {
            yield()
            withContext(NonCancellable) {
                delay(timeMillis = it * 100L)
            }
        }
    }
    delay(timeMillis = 200L)
    job.cancelAndJoin()
}

它的原理是什么

先说原理。其实它就是一个实现了 Job 接口的 Job 上下文。

不过 Job 的一些能力都是废弃、阉割的,比如:

@Deprecated(level = DeprecationLevel.WARNING, message = message)
override val parent: Job
    get() = null

@Deprecated(level = DeprecationLevel.WARNING, message = message)
override fun cancel(cause: CancellationException) {}

@Deprecated(level = DeprecationLevel.HIDDEN, message = "Since 1.2.0, binary compatibility with versions <= 1.1.x")
override fun cancel(cause: Throwable): Boolean = false // never handles exceptions

@Deprecated(level = DeprecationLevel.WARNING, message = message)
override val children: Sequence
    get() = emptySequence()

在协程内部调用 ensureActive() 检查取消状态时,本质上查的是 Job.isActive

问题就在这里:进入 withContext 之后,拿到的是一份新的协程上下文,而这个上下文里的 Job 实际上是 NonCancellable

这个对象的 isActive 被直接写死为始终返回 true。因此,取消检查永远通不过“已取消”这道判断,异常自然也不会被抛出。

换言之,withContext 内部不知道当前协程已经取消了,所以内部能够不响应取消。

public suspend fun  withContext(
    context: CoroutineContext,
    block: suspend CoroutineScope.() -> T
): T {
    contract {
        callsInPlace(block, InvocationKind.EXACTLY_ONCE)
    }
    return suspendCoroutineUninterceptedOrReturn sc@ { uCont ->
        val oldContext = uCont.context
        val newContext = oldContext.newCoroutineContext(context)
        newContext.ensureActive() // 内部获取(get(Job))到的 Job 对象是 NonCancellable
        // ...
    }
}

正确使用方式

NonCancellable 要和 withContext 搭配使用。

不要把它直接塞进 launch 这类协程构建器的上下文里。

原因很直接:这样做会把协程的结构化并发机制拦腰打断。

  • 子协程不会再随着父协程的取消而结束。

  • 父协程也不会等它执行完。

  • 哪怕子协程中途崩溃,父协程同样不会因此被取消。

说到底,NonCancellable 会把这层父子协程关系彻底切开。

NonCancellable 的顶部注释非常详细记录了正确的使用方法,可以在那学习。

免责声明:文中图文均来自网络,如有侵权请联系删除,心愿游戏发布此文仅为传递信息,不代表心愿游戏认同其观点或证实其描述。

相关文章

更多

精选合集

更多

大家都在玩

热门话题

大家都在看

更多