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

Kotlin Flow实现按事件类型切换collect/collectLatest行为

问题描述

存在两种事件类型:

enum class Event {
  REGULAR,
  URGENT
}

事件会发送至val eventsFlow = SharedFlow<Event>(或Channel),当前的收集逻辑如下:

scope.launch {
  eventsFlow.collect { event ->
     delay(100) // 统一的事件处理逻辑
  }
}

需要实现的效果:

  • REGULAR事件:不中断当前处理,若新的REGULAR事件到达时收集器正忙碌,保持标准collect的串行排队行为
  • URGENT事件:到达时立即中断当前正在处理的事件(无论当前处理的是REGULAR还是URGENT),类似collectLatest的抢占行为
解决方案

可以通过分离REGULAR事件的串行处理队列,结合collectLatest实现URGENT事件的抢占逻辑,具体代码如下:

1. 定义统一事件处理函数

suspend fun handleEvent(event: Event) {
    delay(100) // 替换为你的实际业务处理逻辑
}

2. 实现带优先级的流收集逻辑

scope.launch {
    // 用Channel缓存REGULAR事件,保证串行处理
    val regularChannel = Channel<Event>(Channel.UNLIMITED)
    var regularProcessingJob: Job? = null

    // 启动REGULAR事件的串行处理任务
    regularProcessingJob = launch {
        regularChannel.consumeEach { event ->
            handleEvent(event)
        }
    }

    // 用collectLatest处理原流,实现URGENT事件的抢占
    eventsFlow.collectLatest { event ->
        when (event) {
            Event.REGULAR -> {
                // REGULAR事件加入串行处理队列
                regularChannel.send(event)
            }
            Event.URGENT -> {
                // 中断当前正在进行的REGULAR事件处理
                regularProcessingJob?.cancel()
                // 立即处理URGENT事件
                handleEvent(event)
                // 重启REGULAR事件的串行处理任务,后续REGULAR事件继续排队处理
                regularProcessingJob = launch {
                    regularChannel.consumeEach { event ->
                        handleEvent(event)
                    }
                }
            }
        }
    }
}

3. 封装为可复用的Flow运算符(可选)

如果需要在多个地方复用该逻辑,可以封装成自定义Flow运算符:

fun <T> Flow<T>.handleWithPriority(
    isUrgent: (T) -> Boolean,
    handler: suspend (T) -> Unit
): Flow<Unit> = transformLatest { event ->
    val regularChannel = Channel<T>(Channel.UNLIMITED)
    var regularJob = launch {
        regularChannel.consumeEach { handler(it) }
    }

    if (isUrgent(event)) {
        regularJob.cancel()
        handler(event)
        regularJob = launch {
            regularChannel.consumeEach { handler(it) }
        }
    } else {
        regularChannel.send(event)
    }
    emit(Unit)
}

使用方式:

scope.launch {
    eventsFlow.handleWithPriority(
        isUrgent = { it == Event.URGENT },
        handler = ::handleEvent
    ).collect()
}

逻辑说明

  • REGULAR事件:通过Channel实现串行排队,新的REGULAR事件不会中断当前正在处理的REGULAR事件,符合标准collect的行为
  • URGENT事件:借助collectLatest的特性,新的URGENT会中断当前正在执行的任何处理(包括之前的URGENT处理);同时主动取消REGULAR事件的处理任务,立即执行URGENT逻辑,处理完成后重启REGULAR的串行处理队列

内容的提问来源于stack exchange,提问作者Sergey Borotynsky

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 21:25:32