位置:首页 > Kotlin > Kotlin协程Channel热数据通道原理与使用指南

Kotlin协程Channel热数据通道原理与使用指南

时间:2026-08-21  |  作者:风起客  |  阅读:0

种一颗树的最好时机是十年前,其次是现在。

学习也一样。

跟着霍老师的《深入理解 Kotlin 携程》学习一下协程。

kotlin协程-热数据通道Channel

直奔主题,认识 Channel

Channel 实际上就是一个并发安全的队列。

它可以用来连接协程,实现不同协程之间的通信。

suspend fun main() {
    val channel = Channel<Int>()
    val producer = GlobalScope.launch {
        var i = 0
        while (true) {
            delay(1000)
            channel.send(i++)
        }
    }
    val consumer = GlobalScope.launch {
        while (true) {
            val element = channel.receive()
            println(element)
        }
    }
    producer.join()
    consumer.join()
}

上述代码中构造了两个协程:producer 和 consumer。

我们没有为它们明确指定调度器,所以它们使用的都是默认调度器。

其中,producer 每隔 1 秒向 Channel 发送一个整数。

consumer 则持续读取 channel 中的数据并打印。

显然,发送端比接收端更慢。

在没有可读取值时,receive 会挂起,直到有新元素到达。

这样看来,receive 一定是一个挂起函数。

那么 send 呢?

Channel 的容量

我们查看 send 方法的声明,会发现它也是挂起函数。

那么发送端为什么也要挂起?

前面提到,Channel 实际上就是一个队列。

队列中存在缓冲区。

一旦缓冲区满了,并且一直没有人调用 receive 取走元素,send 就要挂起。

它会等待接收者取走元素后,再继续写入 Channel。

public fun  Channel(capacity: Int = RENDEZVOUS): Channel =
    when (capacity) {
        RENDEZVOUS -> RendezvousChannel()
        UNLIMITED -> LinkedListChannel()
        CONFLATED -> ConflatedChannel()
        BUFFERED -> ArrayChannel(CHANNEL_DEFAULT_CAPACITY)
        else -> ArrayChannel(capacity)
    }

我们在构造Channel时,调用了一个名为Channel的函数。

不过,它并非Channel的构造函数。

在Kotlin里,常常会定义一个顶级函数来模拟同名类型的构造器,其实质就是工厂函数。

这里有个Int类型的参数capacity,其默认值为RENDEZVOUS。

需注意,如果不调用receive,send就会一直处于挂起等待状态。

几种容量类型的含义

  • UNLIMITED:表示没有任何限制,对所有来的元素都接受。

  • CONFLATED:名字看起来像“合并”,但它的实际作用是只保留最后一个元素。

  • 也就是说,它的缓冲区只有一个元素大小。每次有新元素到来,都会把旧元素覆盖掉。

  • BUFFERED:效果和ArrayBlockingQueue类似,它接收一个值作为缓冲区的容量大小。

迭代 Channel

在发送和读取时,我们写了一个while(true)死循环。

因为这里需要不断进行读写操作。

实际上,我们也可以直接获取一个 Channel 的 Iterator。

    val consumer = GlobalScope.launch {
        val iter = channel.iterator()
        while (iter.hasNext()) {
            val element = iter.next()
            println(element)
            delay(1000)
        }
    }

其中 iter.hasNext() 是挂起函数。

因为在判断是否有下一个元素时,就需要去 Channel 中读取元素。

当然,也可以直接使用 for ... in ..

    val consumer = GlobalScope.launch {
        for(element in channel) {
            println(element)
            delay(1000)
        }
    }

produce 和 actor

再来看两个便捷的构造生产者和消费者的 api。

  • 可以通过produce启动生产者协程,并返回一个ReceiveChannel

  • 其他协程就可以通过这个 channel 获取数据。

  • 可以通过actor启动消费者协程,并返回一个SendChannel

  • 其他协程就可以通过这个 channel 发送数据。

val receiveChannel : ReceiveChannel<Int> = GlobalScope.produce{
    repeat(100){

        delay(100)
        send(it)
    }
}

val sendChannel : SendChannel<Int> = GlobalScope.actor {
    while(true){
        val element = receive()
        println(element)
    }
}

Channel 的关闭

以上面的produce方法为例,我们可以看到最终返回的是一个ProducerCoroutine

它的定义如下:

private class ProducerCoroutine<E>(
    parentContext: CoroutineContext, channel: Channel
) : ChannelCoroutine(parentContext, channel, true, active = true), ProducerScope {
    override val isActive: Boolean
        get() = super.isActive

    override fun onCompleted(value: Unit) {
        _channel.close()
    }

    override fun onCancelled(cause: Throwable, handled: Boolean) {
        val processed = _channel.close(cause)
        if (!processed && !handled) handleCoroutineException(context, cause)
    }
}

我们会发现,在它的完成取消方法中,都会调用_channel.close的方法。

也正是这样,Channel 才被称为热数据流。

这里有一点需要注意。

对于一个Channe,如果我们调用了它的close方法,它会立即停止接收新元素。

也就是说,这时候它的isClosedForSend会立即返回true

但由于Channel缓冲区的存在,这时候可能还有一些元素没有被处理完。

因此,要等所有元素都被读取之后,isClosedForReceive才会返回true

BroadcastChannel

在实际环境中,经常会出现一个发送对应多个接收的情况。

这时就需要 BroadcastChannel 了。

    val broadcastChannel =  BroadcastChannel<Int>(Channel.BUFFERED)

    val producer = GlobalScope.launch {
        List(3){
            delay(1000)
            broadcastChannel.send(it)
        }
    }
    List(3){
        index ->
        GlobalScope.launch {
            val receiveChannel = broadcastChannel.openSubscription()
            for(i in receiveChannel){
                println("[#$index] received $i")
            }
        }
    }.joinAll()

这里有个细节需要注意。

如果把发送端的dealy(100)去掉,可能会出现部分元素收不到,或者完全收不到的情况。

这是因为BroadcastChannel在发送的时候如果没有订阅者,这条消息就会被丢弃。

我们也可以通过普通的 Channel 进行转换:

val channel = Channel<Int>()
channel.broadcast(3)

这里也需要注意,BroadcastChannel被标记为过时了。

可以使用SharedFlowStateFlow代替。

channel.broadcast()方法也被标记为过时,同样使用SharedFlow来代替。

以上

免责声明:文中图文均来自网络,如有侵权请联系删除,心愿游戏发布此文仅为传递信息,不代表心愿游戏认同其观点或证实其描述。

相关文章

更多

精选合集

更多

大家都在玩

热门话题

大家都在看

更多