一、概念
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): ChannelResult<Unit> 普通函数式发送,立即返回结果,成功返回Success,失败返回Closed或队列已满,适用于非关键事件,允许丢失,避免协程挂起。 | |
| close() | public fun close(cause: Throwable? = null): Boolean 调用后表示关闭发送功能,此时 isClosedForSend() 会返回true,继续发送元素会报错ClosedSendChannelException。缓冲区里的元素可以继续被消费,等所有元素被消费后 isClosedForReceive() 会返回true,for循环会自动结束,继续消费会报错ClosedReceiveChannelException。具有原子性。 | |
| isClosedForSend() | public val isClosedForSend: Boolean 判断通道是否关闭了发送功能,若为true,调用 send() 继续发送元素会报错ClosedSendChannelException。 | |
ReceiveChannel 消费者通道 | receive() | public suspend fun receive(): E 从通道中消费一个元素并移除,通道中没有元素时将被挂起,直到有新元素被发送进来。 |
| tryReceive() | public fun tryReceive(): ChannelResult<E> 当通道中有值时就从中消费,并返回成功结果。否则返回失败或关闭的结果。 | |
| receiveCatching() | public suspend fun receiveCatching(): ChannelResult<E> 如果此通道不为空,则从中检索并删除元素,返回成功结果;如果通道为空,则返回失败结果;如果通道关闭,则返回关闭的原因。比 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(): ChannelIterator<E> 返回通道的迭代器。 |
- 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 = Channel<String>() val receiveChannel: ReceiveChannel<String> = _channel val sendChannel: SendChannel<String> = _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.0(2018年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 ActorScope<E>.() -> Unit ): SendChannel<E> |
fun main() = runBlocking { //创建生产者通道 val send: SendChannel<Int> = 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( |
fun main() = runBlocking { //创建消费者通道 val receiveChannel: ReceiveChannel<Int> = produce { (1..3).forEach { send(it) } //定义了生产方式 } //其它协程拿来接收数据 launch { for (i in receiveChannel) println(i) } }2.1.3 构造创建 Channel
| public fun <E> Channel( capacity: Int = RENDEZVOUS, //缓冲区容量 onBufferOverflow: BufferOverflow = BufferOverflow.SUSPEND, //缓冲区溢出策略 onUndeliveredElement: ((E) -> Unit)? = null //异常回调 ): Channel<E> |
capacity 缓冲区容量 | 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 = Channel<String>() //约会类型 val bufferedChannel = Channel<String>(10) //指定缓存大小类型 val conflatedChannel = Channel<String>(Channel.CONFLATED) //混合类型 val unlimitedChannel = Channel<String>(Channel.UNLIMITED) //无限缓存大小类型fun main() = runBlocking { val channel = Channel<Int>(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 = Channel<String>(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 = Channel<String>(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> ReceiveChannel<E>.consume(block: ReceiveChannel<E>.() -> R): R 执行完后关闭通道。 |
| consumeEach() | public suspend inline fun <E> ReceiveChannel<E>.consumeEach(action: (E) -> Unit): Unit 将 action 应用于每一个元素,执行完后会关闭通道。 |
//只消费第一个元素并关闭通道 suspend fun <E> ReceiveChannel<E>.consumeFirst(): E = consume { return receive() }2.2.4 转换成 Flow
底层通过 ChannelFlow 实现,Channel 转成 Flow 是为了使用那些方便的操作符,同时 Flow 的很多功能扩展底层由 Channel 实现(如flowOn、buffer)。
| consumeAsFlow() | public fun <T> ReceiveChannel<T>.consumeAsFlow(): Flow<T> = ChannelAsFlow(this, consume = true) 只能被收集一次(只能有一个消费者),多次收集抛异常 IllegalStateException。 |
| receiveAsFlow() | public fun <T> ReceiveChannel<T>.receiveAsFlow(): Flow<T> = ChannelAsFlow(this, consume = false) 采用扇出模式(可以有多个消费者),一个元素只能被消费一次,该元素的消费者不确定是哪个。 |
三、事件流SharedFlow
BroadcastChannel 在协程1.4版本已被废弃,取而代之的是 SharedFlow(StateFlow是它的特定配置版本)。
SharedFlow只能订阅(消费),FlowCollector只能发送(生产),MutableSharedFlow继承了它俩既能发送也能订阅,根据实际需求暴露不同类型收窄功能。接收会一直监听,通过取消协程来关闭它的收集。
| SharedFlow | public interface SharedFlow<out T> : Flow<T> { //缓存的回放元素的快照 //收集元素 override suspend fun collect(collector: FlowCollector<T>): Nothing |
| FlowCollector | public fun interface FlowCollector<in T> { public suspend fun emit(value: T) //发送元素 } |
| MutableSharedFlow | public interface MutableSharedFlow<T> : SharedFlow<T>, FlowCollector<T> { // 发射元素(注意这是个挂起函数) override suspend fun emit(value: T) // 发射元素(注意这是个普通函数,如果缓存溢出策略是 SUSPEND,溢出时就不会挂起了而是直接返回 false) public fun tryEmit(value: T): Boolean // 活跃订阅者数量,将它设为0生产元素就会停止用来释放资源。 public val subscriptionCount: StateFlow<Int> //清空当前回放里的历史记录 public fun resetReplayCache() } |
fun main() = runBlocking { //利用多态暴露不同父类限制功能给外部使用 private val _mutableSharedFlow = MutableSharedFlow<String>() val sharedFlow: SharedFlow<String> = _mutableSharedFlow //或调用 asSharedFlow() val flowCollector: FlowCollector<String> = _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( |
3.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 = MutableSharedFlow<Int>(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") } } } 打印: emmit:1 甲:1 emmit:2 甲:2 emmit:3 甲:3 emmit:4 甲:4 emmit:5 甲:5 乙:3 //回放了3个已更新过的数据 乙:4 乙:53.2 Flow 转 SharedFlow
| public fun <T> Flow<T>.shareIn( scope: CoroutineScope, //数据共享时所在的协程作用域 started: SharingStarted, //启动策略 replay: Int = 0 //回放,新订阅时得到几个之前已经发射过的旧值。 ): SharedFlow<T> |
started 启动策略 | SharingStarted.Eagerly | 立即发送数据(直到scope结束)。 |
| SharingStarted.Lazily | 在首个订阅者观察时才开始发送数据(当订阅者都没了还是活跃的,直到scope结束)。这保证了第一个订阅者能获得所有值,后续订阅者获得最新replay数量的值。 | |
| SharingStarted.WhileSubscribed | 在首个订阅者观察时才开始发送数据,直到最后一个订阅者消失时停止,当又有新订阅者时会再次启动,避免引起资源浪费(例如一直从数据库、传感器中读取数据)。提供了两个配置:
|
四、状态流 StateFlow
StateFlow继承自SharedFlow是一种特殊配置,相当于MutableSharedFlow(1,0, BufferOverflow.DROP_OLDEST),可以使用value属性来访问值,可以当作是用来取代LiveData。
- 必须传入默认值:Null安全。
- 回放个数1+额外缓冲区大小0:只持有1个值。
- 缓存策略DROP_OLDEST:只持有最新值(新值会替换旧值)。
- 数据防抖:即仅在新值内容发生变化才会消费。
| StateFlow | public interface StateFlow<out T> : SharedFlow<T> { // 当前值 public val value: T } |
| MutableStateFlow | public interface MutableStateFlow<T> : StateFlow<T>, MutableSharedFlow<T> { // 当前值 public override var value: T // 比较并设置(通过 equals 对比,如果值发生真实变化返回 true) public fun compareAndSet(expect: T, update: T): Boolean } |
4.1 通过构造创建
public fun <T> MutableStateFlow(value: T): MutableStateFlow<T> = 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") } } 打印: emit:1 emit:2 collect:2 //2秒后开始订阅,不会收到之前已更新的数据 emit:3 collect:3 emit:4 collect:4 emit:5 collect:54.2 Flow 转 StateFlow
public fun <T> Flow<T>.stateIn( |
public suspend fun <T> Flow<T>.stateIn(scope: CoroutineScope): StateFlow<T> { 挂起函数版本,不用指定默认值,会挂起直到产出第一个值。 |
val uiState: StateFlow<UiState> = repository.getDataFlow() .onStart { emit(UiState.Loading) } .stateIn( scope = viewModelScope, started = SharingStarted.WhileSubscribed(5_000), initialValue = UiState.Loading )四、区别及选择
4.1 通信(Channel) or 广播(SharedFlow)
| Channel | SharedFlow |
| 必须性:事件不会被丢弃且必须执行。 | 时效性:过期的事件没有意义且不应该被延迟消费(就像广播电台,不管有没有人收听都在播放内容,当你开始收听的时候只能听到后续新内容,之前的内容就是错过)。 |
| 分配性:挨个对订阅者发送,多个订阅者互斥,收到的值不是同一个。 | 共享性:同时对全体订阅者发送,多个订阅者共享,接收到的值是同一个。 |
| 同步性:事件只能消费一次。 | 并发性:事件会被多处消费。 |
4.2 事件(Event) or 状态(State)
事件按顺序都要执行到,状态只关心最新值。
| SharedFlow | StateFlow | |
| 类型 | 回放和额外缓存默认为0,无订阅者直接丢弃数据,符合时效性事件特点。 | 仅持有单个且最新的数据。 |
| 初始值 | 无(事件发生后才处理,不需要默认值) | 有(UI组件应当一直有一个值来表明其状态) |
| 回放 | 默认0可配置(新的订阅者不重复处理已发生过的事情,或者需要知道之前发生过的事件) | 1(新的订阅者也应该知道当前状态,若用作事件处理会出现粘性事件) |
| 额外缓冲区 | 默认0可配置 | 0(UI组件只显示最新的值) |
| 缓存模式 | 默认SUSPEND(等待消费) | DROP_OLDEST(UI组件只显示最新的值) |
| 发送重复的值 | 会消费(事件都应该被处理)。 | 不消费(防抖,无变化不用处理)。 |
| 收集方式 | 只能调用方法 collect() 在协程中收集。 | 既可以通过调用属性 value 随处收集,也可以调用方法 collect() 在协程中收集。 |