如何用Kotlin协程实现同Channel的消费与生产任务处理?
问题分析
你的代码抛出ClosedSendChannelException的核心原因是:runBlocking作用域在发送完初始3个任务后就执行完毕,导致整个协程作用域关闭,连带绑定的Channel也被关闭。此时Task(5000)执行完成后尝试向已关闭的Channel发送新任务,就会触发异常。同时,worker的for (task in channel)循环在Channel关闭且现有任务消费完后就会退出,无法处理后续生成的新任务。
修正方案
我们需要跟踪任务的活跃数量,确保所有任务(包括衍生的新任务)都处理完成后再关闭Channel,并且让runBlocking等待所有worker结束。以下是完整的修正代码:
import kotlinx.coroutines.* import kotlinx.coroutines.channels.Channel data class Task(val id: Int) { suspend fun run(): List<Int> { println("starting $id") delay(id.toLong()) println("done $id") if (id >= 5000) return listOf(6000, 7000, 8000) return listOf() } } fun main(args: Array<String>) = runBlocking { val channel = Channel<Task>() val jobs = mutableListOf<Job>() // 用互斥锁保护任务计数,初始为3个种子任务 val activeTasksMutex = kotlinx.coroutines.sync.Mutex() var taskCount = 3 // 启动2个worker协程 repeat(2) { val job = launch { for (task in channel) { val results = task.run() // 新增衍生任务时更新计数 activeTasksMutex.lock() taskCount += results.size activeTasksMutex.unlock() // 将衍生任务发送到队列 for (result in results) { channel.send(Task(result)) } // 当前任务完成,减少计数并检查是否所有任务都处理完毕 activeTasksMutex.lock() taskCount-- val allTasksDone = taskCount == 0 activeTasksMutex.unlock() if (allTasksDone) { channel.close() // 所有任务处理完成,关闭Channel } } } jobs.add(job) } // 发送初始种子任务 for (taskId in listOf(2000, 4000, 5000)) { channel.send(Task(taskId)) } // 等待所有worker协程执行完毕 jobs.forEach { it.join() } }
关键调整点
- 任务计数跟踪:用互斥锁保护
taskCount,准确统计所有活跃任务(初始任务+衍生任务),避免并发修改导致计数错误。 - 延迟Channel关闭:仅当
taskCount归0时才关闭Channel,确保所有衍生任务都能被发送到队列中。 - 等待worker结束:通过
jobs.forEach { it.join() }让runBlocking等待所有worker协程完成,保证所有任务都被处理完毕。
运行修正后的代码,会正常处理所有衍生任务,输出如下:
starting 2000 starting 4000 done 2000 starting 5000 done 4000 done 5000 starting 6000 starting 7000 done 6000 starting 8000 done 7000 done 8000
内容的提问来源于stack exchange,提问作者Random Guy
相关产品推荐
相关产品推荐

