You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kotlin协程中如何加入当前活跃的Continuation?

解决方案:复用正在执行的任务避免重复调用第三方方法

要实现并发调用foo()时复用同一个执行实例,核心思路是用共享的Deferred对象缓存当前正在运行的任务,让后续调用直接等待这个任务的结果,而非重新发起新的第三方调用。

具体实现代码

import kotlinx.coroutines.*
import java.util.concurrent.atomic.AtomicReference

// 原子引用保存当前任务,保证多线程/协程下的线程安全
private val currentFooTask = AtomicReference<Deferred<String>?>(null)

suspend fun foo(): String {
    // 先检查是否有正在运行的任务,有就直接等待其完成
    currentFooTask.get()?.let { return it.await() }

    // 创建延迟启动的任务,确保只有需要时才执行第三方操作
    val newTask = CoroutineScope(Dispatchers.IO).async(start = CoroutineStart.LAZY) {
        suspendCancellableCoroutine<String> { cont ->
            val listener = object : Listener {
                override fun theThingHappened() {
                    cont.resume("The result")
                }
            }

            // 处理协程取消:清理监听器并尝试取消第三方操作
            cont.invokeOnCancellation {
                thirdPartyThing.removeListener(listener)
                thirdPartyThing.cancelCurrentOperation() // 假设第三方提供取消方法
            }

            // 启动第三方耗时流程
            thirdPartyThing.doSomethingWithListener(listener)
        }
    }

    // 原子性设置任务,防止并发场景下重复创建
    return if (currentFooTask.compareAndSet(null, newTask)) {
        try {
            newTask.await()
        } finally {
            // 任务完成(成功/失败)后清空引用,允许下次调用重新发起任务
            currentFooTask.set(null)
        }
    } else {
        // 有其他线程先创建了任务,直接等待该任务结果
        currentFooTask.get()!!.await()
    }
}

关键逻辑说明

  • 线程安全的任务共享:AtomicReference确保多协程/线程场景下,不会同时创建多个第三方任务,避免触发重复调用的报错。
  • 延迟启动任务:CoroutineStart.LAZY模式保证只有第一个调用await()时,才会真正执行第三方的耗时操作。
  • 取消处理:通过invokeOnCancellation在协程取消时清理监听器、取消第三方操作,避免资源泄漏。
  • 自动清理任务:任务完成后清空currentFooTask,确保下次调用可以重新发起新任务(如果需要缓存成功结果,可修改逻辑把结果存入引用,失败时再允许重试)。

扩展场景:缓存成功结果

如果希望成功的结果能被后续调用直接复用,无需重复执行第三方操作,可以调整逻辑缓存结果:

private val cachedFooResult = AtomicReference<String?>(null)
private val currentFooTask = AtomicReference<Deferred<String>?>(null)

suspend fun foo(): String {
    // 优先返回缓存结果
    cachedFooResult.get()?.let { return it }
    // 再检查是否有正在运行的任务
    currentFooTask.get()?.let { return it.await() }

    val newTask = CoroutineScope(Dispatchers.IO).async(start = CoroutineStart.LAZY) {
        val result = suspendCancellableCoroutine<String> { /* 原有逻辑 */ }
        // 缓存成功结果
        cachedFooResult.set(result)
        result
    }

    // 后续逻辑同基础实现...
}

内容的提问来源于stack exchange,提问作者Codepunk

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 14:00:29