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

如何在不取消父协程的前提下取消Kotlin子协程Flow

问题描述

我正尝试为基于协程Flow的MVI实现添加取消操作功能。统一数据流(UDF)以SharedFlow事件为起点,随后被flatMap转换为其他Flow,这些Flow会返回结果并传递给reducer。如何为每个flatMap处理的动作创建可通过其他事件取消的Job?以下是需重构以添加可取消特性的现有代码:

class MyViewModel : ViewModel() {

    val inputs: MutableSharedFlow<I> = MutableSharedFlow()
    private val viewModelListener: MutableStateFlow<Output> = MutableStateFlow(initialState)
    private lateinit var job: Job
    
    fun observe(): StateFlow<Output> = viewModelListener.asStateFlow()
    
    private fun activate() {
       job = viewModelScope.launch(dispatcher) {
            inputs.flatMapMerge { cancellableAction(it) } // how to make this flow cancellable?
            .flatMapConcat { flowOf(it, ProgressOutcome(false)) }
            .onStart { emit(ProgressOutcome(true)) }
            .catch { cause -> emit(ErrorOutcome(cause)) }
            .flowOn(dispatcher)
            .collect { viewModelListener.emit(it) }
    }

    final override fun onCleared(): Unit = job.cancel()
}
解决方案

针对这个需求,有两种常见的实现思路,可根据业务场景选择:

思路1:手动维护当前活跃Job(支持精准取消)

如果需要单独取消特定操作,或者区分不同类型任务的取消逻辑,可以手动维护跟踪活跃Job的变量:

步骤:

  • 在输入事件类型中添加取消事件(例如CancelCurrentAction),或设计带标识的取消事件
  • 用变量/Map存储当前活跃的Job,收到取消事件或新操作事件时,先取消前序Job
  • 执行新操作时,启动新Job并关联到对应的Flow

重构后代码:

// 定义输入事件类型示例
sealed class Input {
    data class DoAction(val data: String) : Input()
    object CancelCurrentAction : Input()
}

class MyViewModel : ViewModel() {

    val inputs: MutableSharedFlow<Input> = MutableSharedFlow()
    private val viewModelListener: MutableStateFlow<Output> = MutableStateFlow(initialState)
    private var currentActionJob: Job? = null

    fun observe(): StateFlow<Output> = viewModelListener.asStateFlow()

    private fun activate() {
        viewModelScope.launch(dispatcher) {
            inputs.collect { input ->
                when (input) {
                    Input.CancelCurrentAction -> {
                        // 取消当前活跃任务
                        currentActionJob?.cancel()
                        currentActionJob = null
                        viewModelListener.emit(ProgressOutcome(false))
                    }
                    is Input.DoAction -> {
                        // 先取消之前的操作,避免并发执行
                        currentActionJob?.cancel()
                        // 启动新任务并跟踪Job
                        currentActionJob = viewModelScope.launch(dispatcher) {
                            runCatching {
                                cancellableAction(input)
                                    .onStart { emit(ProgressOutcome(true)) }
                                    .catch { cause -> emit(ErrorOutcome(cause)) }
                                    .collect { result ->
                                        viewModelListener.emit(result)
                                        viewModelListener.emit(ProgressOutcome(false))
                                    }
                            }.onFailure {
                                // 过滤取消异常,不把取消当成错误处理
                                if (!it.isCancellationException) {
                                    viewModelListener.emit(ErrorOutcome(it))
                                    viewModelListener.emit(ProgressOutcome(false))
                                }
                            }
                        }
                    }
                }
            }
        }
    }

    // 示例耗时操作Flow,内部需支持协程取消(如使用delay、withContext等)
    private suspend fun cancellableAction(input: Input.DoAction): Flow<Output> {
        return flow {
            delay(2000) // 模拟耗时操作
            emit(SuccessOutcome(input.data))
        }
    }

    override fun onCleared() {
        currentActionJob?.cancel()
        super.onCleared()
    }
}

思路2:使用flatMapLatest自动取消前序操作

如果你的场景是“新操作进来时自动取消旧操作”,可以直接用flatMapLatest替代flatMapMerge,它的特性是上游发出新值时,自动取消前一个下游Flow的收集:

重构后代码:

class MyViewModel : ViewModel() {

    val inputs: MutableSharedFlow<I> = MutableSharedFlow()
    private val viewModelListener: MutableStateFlow<Output> = MutableStateFlow(initialState)
    private lateinit var job: Job
    
    fun observe(): StateFlow<Output> = viewModelListener.asStateFlow()
    
    private fun activate() {
       job = viewModelScope.launch(dispatcher) {
            inputs.flatMapLatest { input ->
                // 处理取消事件(如果有)
                if (input is CancelAction) {
                    flowOf(ProgressOutcome(false))
                } else {
                    cancellableAction(input)
                        .onStart { emit(ProgressOutcome(true)) }
                        .catch { cause -> emit(ErrorOutcome(cause)) }
                        .onCompletion {
                            // 操作完成或取消时更新进度状态
                            if (it?.isCancellationException != true) {
                                emit(ProgressOutcome(false))
                            }
                        }
                }
            }
            .catch { cause -> emit(ErrorOutcome(cause)) }
            .flowOn(dispatcher)
            .collect { viewModelListener.emit(it) }
    }

    final override fun onCleared(): Unit = job.cancel()
}

补充说明:

  • flatMapLatest适合“只关心最新操作结果”的场景,比如搜索输入,用户连续输入时仅保留最后一次搜索请求
  • 若需要同时执行多个操作并单独取消,建议选择思路1,用Map存储不同操作的Job(按操作ID映射)
  • 确保cancellableAction内部的Flow支持协程取消,自定义耗时操作需检查isActive状态主动响应取消

内容的提问来源于stack exchange,提问作者Zeyad Gasser

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 22:57:47