如何以可读方式整合相互依赖的Kotlin Coroutine Flow链式调用?
问题描述
我遇到了一个复杂的协程Flow依赖场景,现有一组相互关联的链式调用代码如下:
viewModelScope.launch { repository.cacheAccount(person) .flatMapConcat { it-> Log.d(App.TAG, "[2] create account call (server)") repository.createAccount(person) } .flatMapConcat { it -> if (it is Response.Data) { repository.cacheAccount(it.data) .collect { it -> // no op, just execute the command Log.d(App.TAG, "account has been cached") } } flow { emit(it) } } .catch { e -> Log.d(App.TAG, "[3] get an exception in catch block") Log.e(App.TAG, "Got an exception during network call", e) state.update { state -> val errors = state.errors + getErrorMessage(PersonRepository.Response.Error.Exception(e)) state.copy(errors = errors, isLoading = false) } } .collect { it -> Log.d(App.TAG, "[4] collect the result") updateStateProfile(it) } }
当前流程步骤:
- 在本地磁盘缓存账户
- 调用后端接口创建账户
- 创建成功时,把新账户数据缓存到本地磁盘
现在需要新增以太坊链相关的API调用,流程变得更复杂:
- 4a. 账户创建成功时,将初始化的事务缓存到本地磁盘(调用
cacheRepository.createChainTx()) - 4b. 账户创建失败时,直接转发后端返回的响应
- 事务缓存成功后,调用第二个端点注册用户(
repository.registerUser())
- 事务缓存成功后,调用第二个端点注册用户(
- 无论第二个端点调用成功还是失败(异常除外),都要将响应更新到缓存中,以更新事务状态
新增流程的单独实现代码如下:
viewModelScope.launch { lateinit var newTx: ITransaction cacheRepository.createChainTxAsFlow(RegisterUserTransaction(userWalletAddress = userWalletAddress)) .map { it -> newTx= it repository.registerUserOnSwapMarket(userWalletAddress) } .onEach { it -> preProcessResponse(it, newTx) } .flowOn(backgroundDispatcher) .collect { it -> processResponse(it) } }
我需要把这个新增场景整合到第一个Flow链中,但不清楚怎么用清晰的链式写法实现;如果放弃链式写法,会产生大量if else语句。请问如何以易读的方式实现这个整合后的场景?
解决方案
要保持Flow链式的易读性,核心是把复杂分支逻辑拆成独立函数,让主流程保持线性结构,同时通过flatMapConcat、onEach等操作符串联业务步骤。
步骤1:拆分独立逻辑函数
先把缓存事务、处理注册响应等分支逻辑抽成单独函数,避免主流程嵌套混乱:
// 处理账户创建成功后的链上注册全流程 private fun handleChainRegistrationFlow( accountData: Response.Data<Account>, userWalletAddress: String ): Flow<Response<RegisterResult>> { return cacheRepository.createChainTxAsFlow(RegisterUserTransaction(userWalletAddress)) .flatMapConcat { newTx -> repository.registerUserOnSwapMarket(userWalletAddress) .onEach { registerResponse -> // 无论成功失败(异常除外),更新缓存事务状态 preProcessResponse(registerResponse, newTx) } } } // 转发失败响应的通用函数 private fun <T> forwardFailureResponse(response: Response.Error): Flow<Response<T>> { return flow { emit(response as Response<T>) } }
步骤2:重构主Flow链
将新增逻辑整合到原流程中,通过flatMapConcat处理分支,主流程保持线性清晰:
viewModelScope.launch { repository.cacheAccount(person) // 步骤1:缓存本地账户后,调用后端创建账户 .flatMapConcat { Log.d(App.TAG, "[2] create account call (server)") repository.createAccount(person) } // 步骤2:根据账户创建结果分支处理 .flatMapConcat { createAccountResponse -> when (createAccountResponse) { is Response.Data -> { // 步骤3:缓存新账户,完成后启动链上注册流程 repository.cacheAccount(createAccountResponse.data) .onEach { Log.d(App.TAG, "account has been cached") } .flatMapConcat { handleChainRegistrationFlow(createAccountResponse, userWalletAddress) } // 可选:如果需要保留原账户创建响应,用onStart转发 .onStart { emit(createAccountResponse) } } is Response.Error -> { // 步骤4b:直接转发失败响应 forwardFailureResponse(createAccountResponse) } } } // 全局异常捕获 .catch { e -> Log.d(App.TAG, "[3] get an exception in catch block") Log.e(App.TAG, "Got an exception during network call", e) state.update { val errors = it.errors + getErrorMessage(PersonRepository.Response.Error.Exception(e)) it.copy(errors = errors, isLoading = false) } } // 统一指定后台线程 .flowOn(backgroundDispatcher) // 最终收集结果更新状态 .collect { finalResponse -> Log.d(App.TAG, "[4] collect the result") updateStateProfile(finalResponse) } }
关键优化点
- 逻辑解耦:把分支逻辑抽成独立函数,主流程只关注步骤串联,可读性大幅提升
- 链式保持:用
flatMapConcat处理异步分支,避免嵌套if else堆砌 - 统一异常处理:全局
catch块捕获所有流程异常,不用分散处理 - 节点清晰:每个操作符对应明确业务步骤,注释标注核心动作
如果需要同时保留账户创建结果和注册结果,可以用Pair或自定义数据类封装:
// 修改handleChainRegistrationFlow返回封装结果 private fun handleChainRegistrationFlow( accountData: Response.Data<Account>, userWalletAddress: String ): Flow<Pair<Response.Data<Account>, Response<RegisterResult>>> { return cacheRepository.createChainTxAsFlow(RegisterUserTransaction(userWalletAddress)) .flatMapConcat { newTx -> repository.registerUserOnSwapMarket(userWalletAddress) .onEach { preProcessResponse(it, newTx) } .map { registerResponse -> accountData to registerResponse } } }
内容的提问来源于stack exchange,提问作者Gleichmut
相关产品推荐
相关产品推荐

