如何将重复的Flow处理逻辑提取为公共函数?
提取重复Flow逻辑的解决方案
核心思路
把所有任务中重复的Flow操作链、Job生命周期管理、状态更新逻辑封装成公共函数或扩展方法,让每个任务函数只需要关注调用对应UseCase和传递参数。
方案1:Flow扩展函数+Job管理公共函数
这种方式解耦性更强,逻辑复用更灵活。
1. 封装通用Flow处理逻辑(扩展函数)
给Flow<Resource<MyType>>编写扩展函数,将重复的onEach、catch、retry逻辑封装进去:
// 替换YourStateType为myStateFlow实际的状态类型 fun Flow<Resource<MyType>>.applyCommonStateHandling( stateFlow: MutableStateFlow<YourStateType>, loadingMsg: String = "loading message", errorMsg: String = "error message", retryDelayMs: Long = 2500 ): Flow<Resource<MyType>> { return this .onEach { resource -> stateFlow.update { it.copy( isConsumed = false, resource = resource ) } } .catch { throwable -> val currentState = stateFlow.value val currentTryCount = currentState.currentTryCount val maxTryCount = currentState.maxTryCount if (throwable is BleReTryableError && currentTryCount < maxTryCount) { stateFlow.update { it.copy( isConsumed = false, resource = Resource.Loading(loadingMsg), currentTryCount = currentTryCount + 1 ) } // 抛出异常触发重试逻辑 throw BleReTryableError(throwable.message ?: "") } else { stateFlow.update { it.copy( isConsumed = false, resource = Resource.Error(errorMsg) ) } } } .retry { delay(retryDelayMs) it is BleReTryableError } }
2. 封装Job启动/取消逻辑
编写公共函数处理旧任务取消、新任务启动的逻辑:
private fun <ParamType> launchManagedJob( existingJob: Job?, param: ParamType, useCase: (ParamType) -> Flow<Resource<MyType>>, scope: CoroutineScope = viewModelScope ): Job? { // 取消旧任务 existingJob?.cancel() // 启动新任务并返回Job引用 return useCase(param) .applyCommonStateHandling(myStateFlow) .launchIn(scope) }
3. 简化原有任务函数
现在每个Job函数只需要一行代码:
private var job1Job: Job? = null private fun job1(param: Param) { job1Job = launchManagedJob(job1Job, param, useCases1) } private var job2Job: Job? = null private fun job2(param: Param) { job2Job = launchManagedJob(job2Job, param, useCases2) } private var job3Job: Job? = null private fun job3(param: Param) { job3Job = launchManagedJob(job3Job, param, useCases3) }
方案2:单一公共函数(类似你提到的myRun)
如果不需要灵活调整参数,也可以把所有逻辑合并到一个函数中:
private fun runJob( existingJob: Job?, flow: Flow<Resource<MyType>> ): Job? { existingJob?.cancel() return flow .onEach { resource -> myStateFlow.update { it.copy(isConsumed = false, resource = resource) } } .catch { throwable -> val currentState = myStateFlow.value val currentTryCount = currentState.currentTryCount val maxTryCount = currentState.maxTryCount if (throwable is BleReTryableError && currentTryCount < maxTryCount) { myStateFlow.update { it.copy( isConsumed = false, resource = Resource.Loading("loading message"), currentTryCount = currentTryCount + 1 ) } throw BleReTryableError(throwable.message ?: "") } else { myStateFlow.update { it.copy(isConsumed = false, resource = Resource.Error("error message")) } } } .retry { delay(2500) it is BleReTryableError } .launchIn(viewModelScope) }
任务函数简化为:
private fun job1(param: Param) { job1Job = runJob(job1Job, useCases1(param)) }
注意事项
- 确保
currentTryCount和maxTryCount从myStateFlow的当前状态中获取(原代码中这两个变量的初始化逻辑缺失,需要补充)。 - 如果不同任务需要自定义加载/错误提示、重试延迟,可以在扩展函数中保留可选参数,调用时传入差异值。
- 所有任务的
viewModelScope如果一致,可以写死在公共函数里,否则作为参数传入。
内容的提问来源于stack exchange,提问作者Metropol_Tilkisi
相关产品推荐
相关产品推荐

