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

Kotlin协程如何实现Channel接收任务并行计算且可向输出Channel发消息

修改方案

你现在的代码串行执行的核心原因是consumeEach循环内同步调用了suspend方法calc(),必须等单次calc执行完成才会消费下一条Channel消息,只需要做两处核心修改即可实现并行计算:

  1. 给Engine类绑定独立协程作用域,复用你已经定义的缓存线程池作为计算调度器,方便管理并行协程的生命周期
  2. 每次收到inputChannel的消息时,启动独立子协程执行calc(),不阻塞Channel消费流程

修改后的完整代码

@ExperimentalCoroutinesApi
fun main() {
    val inputChannel = Channel<Unit>()
    val outputChannel = Channel<String>()
    val engine = Engine(outputChannel)
    val calculator = Logger()

    runBlocking(Dispatchers.Default) {
        launch {
            engine.connect(inputChannel)
        }

        launch {
            calculator.connect(outputChannel)
        }

        inputChannel.send(Unit)
        inputChannel.send(Unit)
        // 等待所有计算任务完成后再关闭程序,可根据实际场景调整等待逻辑
        delay(11000)
        engine.close()
    }
}

class Engine(private val logger: SendChannel<String>) : CoroutineScope {
    // 绑定自定义线程池作为协程调度器,加SupervisorJob避免单个任务失败取消所有并行任务
    private val job = SupervisorJob()
    private val pool = Executors.newCachedThreadPool().asCoroutineDispatcher()
    override val coroutineContext = job + pool

    @ExperimentalCoroutinesApi
    suspend fun connect(input: ReceiveChannel<Unit>) {
        input.consumeEach {
            println("${Instant.now()} [${Thread.currentThread().name}] Engine - Received input")
            // 启动独立子协程执行计算,不阻塞当前消费循环
            launch {
                calc()
            }
        }
    }

    suspend fun calc() {
        logger.send("Starting processing")
        for (i in 1..100) {
            delay(100)
            print(".")
        }
        println()
        logger.send("Finished processing")
    }

    // 对外暴露关闭方法,销毁时统一取消所有运行中的计算任务、释放线程池资源
    fun close() {
        job.cancel()
        pool.close()
    }
}

class Logger {

    @ExperimentalCoroutinesApi
    suspend fun connect(channel: ReceiveChannel<String>) {
        channel.consumeEach {
            println("${Instant.now()} [${Thread.currentThread().name}] Logger - $it")
        }
    }
}

效果说明

修改后两个输入消息触发的calc任务会并行执行,总耗时从原来的20秒左右缩短到10秒左右,日志会出现两条Starting processing先后输出、打印的.交替出现的并行特征,且所有日志消息都可以正常发送到outputChannel由Logger消费。


内容的提问来源于stack exchange,提问作者F.P

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 01:06:07