Kotlin 协程 - 热流 Channel、SharedFlow、StateFlow

发布时间:2026/9/12 10:53:13
Kotlin 协程 - 热流 Channel、SharedFlow、StateFlow 一、概念1.1 通信模式根据接收端节点数量以及接收方式可分为三类。单播 Unicast一对一对单个订阅者发送。组播 Multicast分配性队列分摊消费一对多依次对单个订阅者发送。多个订阅者互斥收到的值不是同一个。同步性元素只能被一个接收端消费。广播 BroadCast共享性广播共享订阅一对多同时对全体订阅者发送。多个订阅者共享接收到的值是同一个。并发性一次发送多处消费。1.2 扇入扇出扇入 Fan-in指直接调用该模块的上级模块的个数多个发送端即多个协程可能会向同一个Channel发送值。扇入大表示模块的复用程序高。扇出 Fan-out指该模块直接调用的下级模块的个数多个接收端即多个协程可能会从同一个Channel中接收值。扇出大表示模块的复杂度高。二、协程间通信 Channel出现在Flow之前现在退居幕后职责单一仅作为协程间通信的并发安全的缓冲队列而存在。两个消费者从同一个队列里取数据一个取走一条另一个再取下一条它们是在分摊消费不是在共享订阅。SendChannel在创建时定义了消费方式外界只能往里发送值ReceiveChannel在创建时定义了生产方式外界只能从中获取值Channel继承了它俩既能发送也能接收根据实际需求暴露不同类型收窄功能。非阻塞类似于 Java 中的 BlockingQueue 队列不同的是 put() 和 take() 读取写入数据是阻塞的而 Channel 中的 send() 和 receive() 是挂起的。同步性每个值只能被众多订阅者中的一个消费。并发安全没有检测到 receive() 的话 send() 就会挂起不会发送值默认模式 RENDEZVOUS没有缓冲区的Channel是同步的。公平性在多个协程中发送或接收多线程竞争遵循先进先出FIFO即队列。SendChannel生产者通道send()public suspend fun send(element: E)挂起函数式发送如果缓冲区已满挂起协程直到有旧元素被消费腾出空间。适用于需要确保事件一定被发送不丢失数据。trySend()public fun trySend(element: E): ChannelResultUnit普通函数式发送立即返回结果成功返回Success失败返回Closed或队列已满适用于非关键事件允许丢失避免协程挂起。close()public fun close(cause: Throwable? null): Boolean调用后表示关闭发送功能此时 isClosedForSend() 会返回true继续发送元素会报错ClosedSendChannelException。缓冲区里的元素可以继续被消费等所有元素被消费后 isClosedForReceive() 会返回truefor循环会自动结束继续消费会报错ClosedReceiveChannelException。具有原子性。isClosedForSend()public val isClosedForSend: Boolean判断通道是否关闭了发送功能若为true调用 send() 继续发送元素会报错ClosedSendChannelException。ReceiveChannel消费者通道receive()public suspend fun receive(): E从通道中消费一个元素并移除通道中没有元素时将被挂起直到有新元素被发送进来。tryReceive()public fun tryReceive(): ChannelResultE当通道中有值时就从中消费并返回成功结果。否则返回失败或关闭的结果。receiveCatching()public suspend fun receiveCatching(): ChannelResultE如果此通道不为空则从中检索并删除元素返回成功结果如果通道为空则返回失败结果如果通道关闭则返回关闭的原因。比 receive() 更温和。它仍然是挂起函数但遇到关闭时不会直接把异常抛出来而是把结果包进 ChannelResult让调用方自己判断。isEmpty()public val isEmpty: Boolean若通道中没有元素且接收未被关闭则返回true。isClosedForReceive()public val isClosedForReceive: Boolean判断通道是否关闭了消费功能若为ttrue调用 receive() 继续消费元素会报错ClosedReceiveChannelException。cancel()public fun cancel(cause: CancellationException? null)以可选原因直接关闭通道缓存中未消费的元素会被丢弃可在通过构造方式创建通道时指定 onUndeliveredElement() 进行处理。iterator()public operator fun iterator(): ChannelIteratorE返回通道的迭代器。send() 和 receive()都是挂起函数。一个协程往队列里放数据另一个协程从队列里取数据。队列为空时接收方挂起队列满时发送方挂起。trySend() 和 tryReceive()普通函数。从非挂起函数中发送或接收元素操作是即时的并返回 ChannelResult 对象包含了有关操作成功或失败的信息。能发就发发不了就立刻返回能收就收收不到也立刻返回。绝不等待不像挂起版会挂起等待来确保完成操作。close() 和 cancel()当创建的是 Channel 类型的时候两个都可以用来关闭。close() 是从生产端关闭表示不再需要新数据会消费完缓冲区中已有的元素。cancel() 是从消费端取消表示不再需要数据会直接关闭并丢未消费的元素若元素持有资源考虑通过 onUndeliveredElement() 处理这些被丢弃的元素进行收尾。通道使用完一定要关闭否则代码块不会停止造成协程阻塞和内存泄漏。send() 和 trySend()由于 emit() 会等待收集器处理完数据在 collect 处理数据的这 5 秒内emit() 所在的协程会被阻塞无法继续执行后续的 emit() 操作这可能会影响整个程序的性能和响应速度 。因此在使用 emit() 时需要充分考虑收集器的性能确保不会因为收集器的问题导致 emit() 长时间阻塞进而影响其他操作的执行 。tryEmit() 会尝试立即发送数据如果发送失败例如缓冲区已满不会进行重试或等待直接返回 false。这就像在网络信号不稳定时发送消息发送者尝试发送后不管消息是否真正送达就继续做其他事情 。在一些对数据实时性要求不高允许部分数据丢失的场景如日志记录、实时监控数据的快速上报等tryEmit 的这种策略能提高系统的性能和响应速度 。fun main() runBlocking { private val _channel ChannelString() val receiveChannel: ReceiveChannelString _channel val sendChannel: SendChannelString _channel _channel.send(Channel发送) sendChannel.send(SendChannel发送) delay(1000) _channel.consumeEach { println(Channel接收$it) } receiveChannel.consumeEach { println(ReceiveChannel接收$it) } }2.1 创建produce() 和 actor() 被定义成协程构建器因此只能在协程环境中调用会在异常、完成、取消时自动关闭同 launch、async 一样作为 CoroutineScope 的扩展函数。actor() 在 kotlinx.coroutines 1.0.02018年10月被标记为 ObsoleteCoroutinesApi处于不推荐使用也暂时不会消失的状态。produce() 依然可用但 Flow 通常是更现代、更主流的选择。例如 ChannelFlow 或 CallbackFlow 来构建异步数据流。2.1.1 actor() 创建 SendChannel创建时在协程构建器 actor() 的 Lambda 中定义了数据的消费方式返回一个生产者通道 SendChannel其它协程通过该对象往里发送数据。把共享状态藏进 actor() 里外部协程只能往里发消息actor() 内部串行处理天然没有竞态。public fun E CoroutineScope.actor(context: CoroutineContext EmptyCoroutineContext,capacity: Int 0,start: CoroutineStart CoroutineStart.DEFAULT,onCompletion: CompletionHandler? null,block: suspend ActorScopeE.() - Unit): SendChannelEfun main() runBlocking { //创建生产者通道 val send: SendChannelInt actor { for (i in channel) println(i) //定义了消费方式 } //其它协程拿来发送数据 launch { (1..3).forEach { sendChannel.send(it) } } }2.1.2 produce() 创建 ReceiveChannel创建时在协程构建器 produce() 的 Lambda 中定义了数据的生产方式返回一个消费者通道 ReceiveChannel其它协程通过该对象从中取出数据。public fun E CoroutineScope.produce(context: CoroutineContext EmptyCoroutineContext,capacity: Int 0,BuilderInference block: suspend ProducerScopeE.() - Unit): ReceiveChannelEfun main() runBlocking { //创建消费者通道 val receiveChannel: ReceiveChannelInt produce { (1..3).forEach { send(it) } //定义了生产方式 } //其它协程拿来接收数据 launch { for (i in receiveChannel) println(i) } }2.1.3 构造创建 Channelpublic fun E Channel(capacity: Int RENDEZVOUS, //缓冲区容量onBufferOverflow: BufferOverflow BufferOverflow.SUSPEND, //缓冲区溢出策略onUndeliveredElement: ((E) - Unit)? null //异常回调): ChannelEcapacity缓冲区容量RENDEZVOUS 同步值为0。同步通信并发安全没有缓冲当面交接。在默认缓冲溢出策略SUSPEND下 send 会挂起直到调用了 receive可看作是同步交替执行并发安全。若更改了默认缓冲溢出策略则会创建缓冲容量为1的通道。CONFLATED 最新值为-1。只关心最新值新值覆盖旧值send永不挂起。虽然效果类似缓冲溢出策略DROP_OLDEST但更改默认缓冲溢出策略SUSPEND会抛异常。BUFFERED 缓冲值为-2大小64或设为1的具体值。节省内存缓冲区满了后根据溢出策略决定send是否被挂起直到消费后腾出空间。UNLIMITED 无限值为Int.MAX_VALUE。绝对不能丢失数据send永不挂起能一直往里发送数据实际开发一般不用容易内存溢出。会忽略缓冲溢出策略设置。onBufferOverflow缓冲区溢出策略当缓存容量 0 或容量 Channel.BUFFERED 时才会触发。BufferOverflow.SUSPEND满了后send会挂起。BufferOverflow.DROP_OLDEST丢弃最旧的值发送新值。BufferOverflow.DROP_LATEST丢弃新值。onUndeliveredElement异常回调元素没有被消费时回调丢弃的元素会在这里收到。若元素持有资源就需要在这里处理被丢弃的元素进行收尾。val rendezvousChannel ChannelString() //约会类型 val bufferedChannel ChannelString(10) //指定缓存大小类型 val conflatedChannel ChannelString(Channel.CONFLATED) //混合类型 val unlimitedChannel ChannelString(Channel.UNLIMITED) //无限缓存大小类型fun main() runBlocking { val channel ChannelInt(capacity Channel.CONFLATED) { println(onUndeliveredElement: $it) } launch { (1..3).forEach { println(send: $it) channel.send(it) } channel.close() //不要忘记关闭 } launch { for (i in channel) { println(receive: $i) } } println(结束) } 打印 结束 send: 1 send: 2 onUndeliveredElement: 1 send: 3 onUndeliveredElement: 2 receive: 32.2 遍历元素接收端生命周期结束时考虑是否应该调用 cancel()。如果有未交付元素持有资源再考虑是否需要 onUndeliveredElement 做收尾。2.2.1 Iterate 迭代器可以理解为不断调用 receive()没有元素时会挂起等待直到有新元素或者通道关闭。fun main(): Unit runBlocking { val channel ChannelString(Channel.UNLIMITED) repeat(5) { channel.send($it) } val iterator channel.iterator() while (iterator.hasNext()) { println(【iterator】${iterator.next()}) delay(1000) } //5个元素打印完后程序没结束Channel不会关闭后面代码执行不到 println(这行执行不到) }2.2.2 for 循环可以理解为不断调用 receive()没有元素时会挂起等待直到有新元素或者通道关闭。fun main():Unit runBlocking { val channel ChannelString(Channel.UNLIMITED) repeat(5) { channel.send($it) } for (i in channel) { println(【for】$i,) delay(1000) } //5个元素打印完后程序没结束Channel不会关闭后面代码执行不到 println(这行执行不到) }2.2.3 扩展函数会在代码块执行完后关闭通道但无法保证在代码块执行完后不会有新元素进入新元素会被丢弃可在通过构造方式创建通道时指定 onUndeliveredElement() 进行处理。consume()public inline fun E, R ReceiveChannelE.consume(block: ReceiveChannelE.() - R): R执行完后关闭通道。consumeEach()public suspend inline fun E ReceiveChannelE.consumeEach(action: (E) - Unit): Unit将 action 应用于每一个元素执行完后会关闭通道。//只消费第一个元素并关闭通道 suspend fun E ReceiveChannelE.consumeFirst(): E consume { return receive() }2.2.4 转换成 Flow底层通过 ChannelFlow 实现Channel 转成 Flow 是为了使用那些方便的操作符同时 Flow 的很多功能扩展底层由 Channel 实现如flowOn、buffer。consumeAsFlow()public fun T ReceiveChannelT.consumeAsFlow(): FlowT ChannelAsFlow(this, consume true)只能被收集一次只能有一个消费者多次收集抛异常 IllegalStateException。receiveAsFlow()public fun T ReceiveChannelT.receiveAsFlow(): FlowT ChannelAsFlow(this, consume false)采用扇出模式可以有多个消费者一个元素只能被消费一次该元素的消费者不确定是哪个。三、事件流SharedFlowBroadcastChannel 在协程1.4版本已被废弃取而代之的是 SharedFlowStateFlow是它的特定配置版本。SharedFlow只能订阅消费FlowCollector只能发送生产MutableSharedFlow继承了它俩既能发送也能订阅根据实际需求暴露不同类型收窄功能。接收会一直监听通过取消协程来关闭它的收集。SharedFlowpublic interface SharedFlowout T : FlowT {//缓存的回放元素的快照public val replayCache: ListT//收集元素override suspend fun collect(collector: FlowCollectorT): Nothing}FlowCollectorpublic fun interface FlowCollectorin T {public suspend fun emit(value: T) //发送元素}MutableSharedFlowpublic interface MutableSharedFlowT : SharedFlowT, FlowCollectorT {// 发射元素注意这是个挂起函数override suspend fun emit(value: T)// 发射元素注意这是个普通函数如果缓存溢出策略是 SUSPEND溢出时就不会挂起了而是直接返回 falsepublic fun tryEmit(value: T): Boolean// 活跃订阅者数量将它设为0生产元素就会停止用来释放资源。public val subscriptionCount: StateFlowInt//清空当前回放里的历史记录public fun resetReplayCache()}fun main() runBlocking { //利用多态暴露不同父类限制功能给外部使用 private val _mutableSharedFlow MutableSharedFlowString() val sharedFlow: SharedFlowString _mutableSharedFlow //或调用 asSharedFlow() val flowCollector: FlowCollectorString _mutableSharedFlow launch { _mutableSharedFlow.collect { println(mutableSharedFlow接收$it) } } launch { sharedFlow.collect { println(mutableSharedFlow接收$it) } } delay(1000) mutableSharedFlow.emit(mutableSharedFlow发送) flowCollector.emit(flowCollector发送) }3.1 通过构造创建MutableSharedFlow()public fun T MutableSharedFlow(replay: Int 0, //回放观察者订阅后能得到几个订阅前已经发出过的旧元素。extraBufferCapacity: Int 0, //额外缓存容量。onBufferOverflow: BufferOverflow BufferOverflow.SUSPEND //缓存溢出策略。): MutableSharedFlowT3.1.1 回放新订阅者能不能拿到订阅前已经播过的旧值能拿几个。replay 0用作UI事件不给新订阅者补发旧值。如Toast、Snackbar、导航事件、点击事件、一次性弹窗事件。replay 1用作状态新订阅者可以收到最近一次值。如最近一次结果、最近一次配置、最近一次广播数据、需要“粘性”的事件。3.1.2 额外缓存容量在回放之外额外增加缓存量所以 replay extraBufferCapacity 才是总共的缓存数量。生产者最多可以额外缓冲几个值而不必马上挂起如果缓存溢出策略设置的是挂起的话。3.1.3 缓存溢出策略缓冲区满了以后事件太多出现背压是挂起、丢旧的、还是丢新的。只有存在消费者时才会触发缓存溢出策略否则直接丢弃。只有在回放或缓存容量0时才支持两种丢弃模式。例如两个launch虽然会并行执行若下方的消费launch在上方的生产launch后执行先发送的值都会被丢弃接收不到直到消费launch执行了才开始接收后面的值。这是和 Channel 最大的区别没有消费者 Channel 不会丢。onBufferOverflow缓存溢出策略BufferOverflow.SUSPEND挂起。适合不能丢事件的场景安全但可能让上游等待。如订单状态事件、支付流程事件、必须按顺序处理的业务事件。BufferOverflow.DROP_OLDEST丢弃最旧的。适合只关心最新事件的场景。如搜索内容变化、滚动位置变化、进度刷新、高频传感器数据。BufferOverflow.DROP_LATEST丢弃最新的。适合当前正在处理的东西更重要新来的可以暂时忽略。如正在处理上一个请求时新到来的请求可以暂时丢弃、某些日志采样保留已有记录而新纪录可以忽略、批量处理任务当前批次未完成时新任务可延后。suspend fun method(): Unit coroutineScope { //让3个launch同时进行 val shared MutableSharedFlowInt(3) //回放3个数据 launch { for(i in 1..5){ shared.emit(i) println(emmit$i) //如果这里不延迟虽然launch是并行执行发送是很快的这里才5个值 //相当于是没有订阅者的5个值都会被丢弃 //除非把下面的订阅launch写在这个launch上面 delay(1000) //1秒更新一个值 } } //订阅者甲 launch { shared.collect{ println(甲$it) } } //订阅者乙5秒后再订阅也就是值全都更新完了再订阅 delay(5000) launch { shared.collect{ println(乙$it) } } } 打印 emmit1 甲1 emmit2 甲2 emmit3 甲3 emmit4 甲4 emmit5 甲5 乙3 //回放了3个已更新过的数据 乙4 乙53.2 Flow 转 SharedFlowpublic fun T FlowT.shareIn(scope: CoroutineScope, //数据共享时所在的协程作用域started: SharingStarted, //启动策略replay: Int 0 //回放新订阅时得到几个之前已经发射过的旧值。): SharedFlowTstarted启动策略SharingStarted.Eagerly立即发送数据直到scope结束。SharingStarted.Lazily在首个订阅者观察时才开始发送数据当订阅者都没了还是活跃的直到scope结束。这保证了第一个订阅者能获得所有值后续订阅者获得最新replay数量的值。SharingStarted.WhileSubscribed在首个订阅者观察时才开始发送数据直到最后一个订阅者消失时停止当又有新订阅者时会再次启动避免引起资源浪费例如一直从数据库、传感器中读取数据。提供了两个配置stopTimeoutMillis超时时间。最后一个订阅者消失后数据流继续活跃多久用于等待新订阅者默认值0表示立刻停止。避免订阅者都消失后就马上关闭数据流例如不想UI有那么几秒不再监听就停止。replayExpirationMillis回放过期时间。数据流停止后保留回放数据的超时时间默认值 Long.MAX_VALUE 表示永久保存。四、状态流 StateFlowStateFlow继承自SharedFlow是一种特殊配置相当于MutableSharedFlow(1,0, BufferOverflow.DROP_OLDEST)可以使用value属性来访问值可以当作是用来取代LiveData。必须传入默认值Null安全。回放个数1额外缓冲区大小0只持有1个值。缓存策略DROP_OLDEST只持有最新值新值会替换旧值。数据防抖即仅在新值内容发生变化才会消费。StateFlowpublic interface StateFlowout T : SharedFlowT {// 当前值public val value: T}MutableStateFlowpublic interface MutableStateFlowT : StateFlowT, MutableSharedFlowT {// 当前值public override var value: T// 比较并设置通过 equals 对比如果值发生真实变化返回 truepublic fun compareAndSet(expect: T, update: T): Boolean}4.1 通过构造创建public fun T MutableStateFlow(value: T): MutableStateFlowT StateFlowImpl(value ?: NULL)形参value是默认值。val state MutableStateFlow(0) //默认值会被覆盖 launch { for(i in 1..5){ state.emit(i) println(emit$i) delay(1000) //1秒更新一个值 } } launch { delay(2000) //2秒后开始订阅 state.collect{ println(collect$it) } } 打印 emit1 emit2 collect2 //2秒后开始订阅不会收到之前已更新的数据 emit3 collect3 emit4 collect4 emit5 collect54.2 Flow 转 StateFlowpublic fun T FlowT.stateIn(scope: CoroutineScope, //数据共享时所在的协程作用域started: SharingStarted, //启动策略initialValue: T //默认值): StateFlowTpublic suspend fun T FlowT.stateIn(scope: CoroutineScope): StateFlowT {val config configureSharing(1)val result CompletableDeferredStateFlowT()scope.launchSharingDeferred(config.context, config.upstream, result)return result.await()}挂起函数版本不用指定默认值会挂起直到产出第一个值。val uiState: StateFlowUiState repository.getDataFlow() .onStart { emit(UiState.Loading) } .stateIn( scope viewModelScope, started SharingStarted.WhileSubscribed(5_000), initialValue UiState.Loading )四、区别及选择4.1 通信(Channel) or 广播(SharedFlow)ChannelSharedFlow必须性事件不会被丢弃且必须执行。时效性过期的事件没有意义且不应该被延迟消费就像广播电台不管有没有人收听都在播放内容当你开始收听的时候只能听到后续新内容之前的内容就是错过。分配性挨个对订阅者发送多个订阅者互斥收到的值不是同一个。共享性同时对全体订阅者发送多个订阅者共享接收到的值是同一个。同步性事件只能消费一次。并发性事件会被多处消费。4.2 事件(Event) or 状态(State)事件按顺序都要执行到状态只关心最新值。SharedFlowStateFlow类型回放和额外缓存默认为0无订阅者直接丢弃数据符合时效性事件特点。仅持有单个且最新的数据。初始值无事件发生后才处理不需要默认值有UI组件应当一直有一个值来表明其状态回放默认0可配置新的订阅者不重复处理已发生过的事情或者需要知道之前发生过的事件1新的订阅者也应该知道当前状态若用作事件处理会出现粘性事件额外缓冲区默认0可配置0UI组件只显示最新的值缓存模式默认SUSPEND等待消费DROP_OLDESTUI组件只显示最新的值发送重复的值会消费事件都应该被处理。不消费防抖无变化不用处理。收集方式只能调用方法 collect() 在协程中收集。既可以通过调用属性 value 随处收集也可以调用方法 collect() 在协程中收集。

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询