如何在不取消父协程的前提下取消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
相关产品推荐
相关产品推荐

