Kotlin协程如何实现Channel接收任务并行计算且可向输出Channel发消息
修改方案
你现在的代码串行执行的核心原因是consumeEach循环内同步调用了suspend方法calc(),必须等单次calc执行完成才会消费下一条Channel消息,只需要做两处核心修改即可实现并行计算:
- 给
Engine类绑定独立协程作用域,复用你已经定义的缓存线程池作为计算调度器,方便管理并行协程的生命周期 - 每次收到
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
相关产品推荐
相关产品推荐

