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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 07:27:08