位置:首页 > Kotlin > 深入理解Kotlin协程:Channel缓冲策略与无锁机制解析

深入理解Kotlin协程:Channel缓冲策略与无锁机制解析

时间:2026-08-15  |  作者:骑光打字机  |  阅读:0

Channel 的基本用法

Channel 是 Kotlin 协程之间进行安全通信的工具。本质上,它就是一个并发安全的队列。以下是其基本用法:

fun main() = runBlocking {
    val channel = Channel<Int>()
    val producer = launch {
        while (isActive) {
            delay(timeMillis = 1000L)
            // 发送数据
            channel.send(Random.nextInt())
        }
    }

    val consumer = launch {
        while (isActive) {
            // 获取数据
            val element = channel.receive()
            println("The received element is $element")
        }
    }

    producer.join()
    consumer.join()
}

当队列为空时,接收端获取数据会挂起等待,直到新元素到来。这个过程不会阻塞线程。

发送端也是一样。当队列塞满后,发送数据会挂起,直到有元素被取走。

容量与缓冲区策略

在阻塞队列(BlockingQueue)中,如果队列空间不足,还往其中添加元素,通常会出现两种情况:

  1. 阻塞等待,直到队列腾出空间。

  2. 抛异常,拒绝此次添加。

Channel 队列同样有缓冲区。先看看它的设置:

public fun  Channel(
    capacity: Int = RENDEZVOUS,
    onBufferOverflow: BufferOverflow = BufferOverflow.SUSPEND,
    onUndeliveredElement: ((E) -> Unit) = null
): Channel =
    when (capacity) {
        RENDEZVOUS -> {
            if (onBufferOverflow == BufferOverflow.SUSPEND)
                BufferedChannel(RENDEZVOUS, onUndeliveredElement) // an efficient implementation of rendezvous channel
            else
                ConflatedBufferedChannel(1, onBufferOverflow, onUndeliveredElement) // support buffer overflow with buffered channel
        }
        CONFLATED -> {
            require(onBufferOverflow == BufferOverflow.SUSPEND) {
                "CONFLATED capacity cannot be used with non-default onBufferOverflow"
            }
            ConflatedBufferedChannel(1, BufferOverflow.DROP_OLDEST, onUndeliveredElement)
        }
        UNLIMITED -> BufferedChannel(UNLIMITED, onUndeliveredElement) // ignores onBufferOverflow: it has buffer, but it never overflows
        BUFFERED -> { // uses default capacity with SUSPEND
            if (onBufferOverflow == BufferOverflow.SUSPEND) BufferedChannel(CHANNEL_DEFAULT_CAPACITY, onUndeliveredElement)
            else ConflatedBufferedChannel(1, onBufferOverflow, onUndeliveredElement)
        }
        else -> {
            if (onBufferOverflow === BufferOverflow.SUSPEND) BufferedChannel(capacity, onUndeliveredElement)
            else ConflatedBufferedChannel(capacity, onBufferOverflow, onUndeliveredElement)
        }
    }

在 Channel 的快捷构造函数,也就是“工厂函数”中,会根据指定的容量(capacity)和溢出策略(onBufferOverflow)来决定缓冲区如何创建。

第三个参数 onUndeliveredElement 是未分发回调。它用于指定元素已经发送,但未传递给消费者时的处理逻辑。这些元素通常是由于 Channel 的关闭和溢出策略导致丢弃的。我们可以在这里释放资源和打印日志。

几种常见容量类型

先看 RENDEZVOUS。它的值是 0,本义是“会合”。

它表示无缓冲区。发送端不发送数据,接收端就会挂起等待;接收端不获取数据,发送端也会挂起等待。

想到这,我脑海中突然涌现了这个画面:

是不是非常贴切啊?两个人需要同时伸手才能够到。

为了更好地观察这个现象,我们改造一下之前的代码:

fun main() = runBlocking {
    val channel = Channel<Int>()
    val producer = launch {
        while (isActive) {
            val randomInt = Random.nextInt()
            delay(timeMillis = 3000L)
            println("before send $randomInt")
            channel.send(randomInt)
            println("after send $randomInt")
        }
    }

    val consumer = launch {
        while (isActive) {
            delay(timeMillis = 5000L)
            val element = channel.receive()
            println("received: $element")
        }
    }

    producer.join()
    consumer.join()
}

可以看到,producer 发送后会挂起。等到 consumer 接收后,它才会继续往下执行,并打印 "after send ..."。

UNLIMITED 的值是 Int.MAX_VALUE,表示无限制。缓冲区永远不会溢出,但它的元素数量仍然受到可用内存的限制。

CONFLATED 的字面意思是“合并、结合”。发送的元素会不断被替换。

例如同时发送了 3 个元素,最后只会保留一个元素。它本质上就是一个容量为 1 的缓冲区,新元素会替换掉旧元素。

BUFFERED 表示默认缓冲,将使用默认的容量 64。

缓冲区溢出策略

再来看缓冲区溢出策略。

  • SUSPEND:默认值。缓冲区满了,生产者会挂起。

  • DROP_LATEST:缓冲区满时,会丢弃最新发送的元素,保留旧数据。

  • DROP_OLDEST:缓冲区满时,会丢弃缓冲区里最旧的元素,保留最新数据。

注意:只有 ConflatedBufferedChannel 支持 DROP_OLDEST / DROP_LATEST,BufferedChannel 只支持 SUSPEND,所有非 SUSPEND 策略,最终都会强制创建 ConflatedBufferedChannel。

深入底层实现:BufferedChannel

接着来说最重要的两个底层实现类:BufferedChannelConflatedBufferedChannel

Channel 的工厂函数最终只会创建这两种通道的实例。

在 Kotlin 协程 1.7 之后,Channel 的内部结构经过了一次重构,底层从原本繁琐的各个实现类,改为了只由 BufferedChannel 和 ConflatedBufferedChannel 接管。

ConflatedBufferedChannel 本质上就是在 BufferedChannel 之上,额外补了一层特定的丢弃策略处理。

换句话说,把 BufferedChannel 吃透,ConflatedBufferedChannel 的逻辑也就不难理解了。

Channel 之所以能同时做到高性能和线程安全,关键正在于它内部这套 BufferedChannel 机制。

物理结构

它的物理结构,是一个由固定大小的数组分段(Segment)组成的链表,也叫块状链表(Unrolled Linked List)。

这种结构同时具备数组和链表的优点。数组缓存友好、随机访问快;链表插入删除灵活。

简单理解,就是固定大小的数组以链表的形式连在一起。就像这样:

虽然是链表,但在逻辑上,被当作了一个无限大的数组来使用。

Segment Node 1      Segment Node 2      Segment Node 3
[ A, B, C, D ] ---> [ E, F, G,   ] ---> [ H, I,  ,  ]

三个核心原子计数器

同时,它采用了非常高效的 FAA(Fetch-And-Add)无锁算法,并维护了三个非常核心的原子计数器:

  • sendersCounter: 记录了 Channel.send() 的调用次数。
  • receiversCounter: 记录了 Channel.receive() 的调用次数。
  • bufferEndCounter: 标记当前缓冲区允许发送者不挂起就能写入的最高索引。

核心流程

有了这些前置知识后,我们再来看它的核心流程:

  • 不管是获取还是发送数据,都会先利用原子的 FAA 操作将自己的计数器 +1。这样做是为了在逻辑无限数组中占用一个单元格。这个格子只属于自己,不会存在两个发送者拿到同一个格子。

  • 如果是 send():它会看分配到的格子里有没有正在挂起的 receive。如果有,就把数据交给并唤醒对方,也就是两者握手。

    如果没有接收者,并且当前索引小于 bufferEndCounter,说明缓冲区还有空位,就直接把数据存进去并返回,不会挂起。

    如果超出了缓冲区边界,就会把自己,也就是协程实例,存入格子并挂起等待。

  • 如果是 receive():当对应格子里已经有数据时,也就是 send 之前存入的内容,会直接取走并清空格子。

    如果格子里是一个挂起的 send,那么它在拿走数据的同时,也会唤醒发送者。

    如果格子是空的,它就会挂起等待发送者存入数据。

这种设计在无竞争时指令极少;在高并发竞争时,也没有线程会被锁阻塞。

所以,Channel 的缓冲区采用了这种架构。

用底层流程理解几种容量类型

我们再用这个流程来解释一下 Channel 的所有容量类型:

  • Channel.RENDEZVOUS (容量 0):缓冲区边界为 0,send 只要拿到坑位,发现永远越界。此时如果没有现成的 receive,就一定会挂起。

  • Channel.UNLIMITED:缓冲区边界为无穷大。此时 send 永远不会越界,所以总是直接把数据塞进数组然后返回,绝对不会挂起。

  • Channel.BUFFERED(指定容量 N):一开始,缓冲区边界就设在 N。前 N 次 send 都还没碰到这条边界,所以会直接写进数组里。

    等到第 N+1 次操作时,就已经越界了。这时候发送方才会进入挂起状态。

    只有当 receive 把数据消费掉,bufferEndCounter 才会继续往前推进,从而腾出新的缓冲空间。

Channel 的迭代与消费

迭代 Channel 时,不需要像上面那样使用 while(isActive) 循环。

我们可以直接获取一个 Channel 的迭代器:

runBlocking {
    val consumer = launch {
        val iterator = channel.iterator()
        while (iterator.hasNext()) { // 挂起点
            val element = iterator.next()
            println(element)
        }
    }
}

hasNext() 是一个挂起函数。它会去 Channel 中读取元素,以判断是否有下一个元素。

在判断过程中,会让对应的协程恢复执行,直到挂起或完成,hasNext 才会结束挂起。

因此,你能看到一个不太符合直觉的现象:输出结果中,B 比 Got 1 还早输出,同样 Done 比 Got 早输出。

runBlocking {
    // 为了赶在第一次调用 hasNext 前,让协程在 send(1) 挂起点,我们使用了 Unconfined 调度器,它会让协程在当前线程同步执行直到挂起 
    val channel = produce(Dispatchers.Unconfined) {
        println("A")
        send(1)
        println("B")
        send(2)
        println("Done")
    }
    for (item in channel) {
        println("Got $item")
    }
}

这个写法还可以简化为 for ... in

runBlocking {
    val consumer = launch {
        for (element in this) {
            println(element)
        }
    }
}

协程构造器:produce 与 actor

如果想快速构造生产者和消费者协程,可以使用 produceactor 方法:

runBlocking {
    val receiveChannel: ReceiveChannel<Int> = produce {
        // 构造生产者协程
        repeat(10) {
            delay(timeMillis = 1000L)
            send(it)
        }
    }
    
    launch {
        // 通过返回的 Channel 获取数据
        for (element in receiveChannel) { // 这是一个挂起操作
            println(element)
        }
    }

    // ============================
    val sendChannel: SendChannel<Int> = actor {
        // 构造消费者协程
        for (element in this) {
            println(element)
        }
    }
    // 通过返回的 Channel 发送数据
    repeat(5) {
        delay(timeMillis = 500L)
        sendChannel.send(-it)
    }

}

ReceiveChannelSendChannel 是 Channel 的父接口,因此 Channel 才既能发也能收。

produceactor 也是协程构造器,只不过它们专用于 Channel。协程结束后,对应的 Channel 也会关闭。

以生产者协程为例:

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() // 在完成时关闭 _channel
    }

    override fun onCancelled(cause: Throwable, handled: Boolean) {
        // 在取消时,同样关闭 _channel
        val processed = _channel.close(cause)
        if (!processed && !handled) handleCoroutineException(context, cause)
    }
}

注意:actor 目前已经被官方标记为了 @ObsoleteCoroutinesApi (即过时 API)。

Channel 的关闭

前面我们多次提到了 Channel 的关闭。只需调用它的 close 方法即可。

关闭后,SendChannel 会立即停止发送新元素,对应的 isClosedForSend 标志会返回 true。

此时缓冲区里可能还有未消费的数据。而 isClosedForReceive 标志只有在缓冲区为空时,才会返回 true。

因此,关闭通道后将无法再发送数据。

  • 如果此时继续调用 send(),会立即抛出 ClosedSendChannelException 异常。

  • 当缓冲区里的数据被全部取完后,如果还继续调用 receive(),会抛出 ClosedReceiveChannelException 异常。

关于 Channel 的取消cancel)可以看我的这篇博客:协程间的通信管道 —— Kotlin Channel 详解

关闭的意义

那关闭有什么意义呢?

  • 一是前面提到的释放资源。

  • 二是在没有数据要发送后,让接收端停止挂起等待。

Channel 的关闭表示没有更多数据要发送了,所以 close 通常由生产者来调用。

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

相关文章

更多

精选合集

更多

大家都在玩

热门话题

大家都在看

更多