如何正确结合Android Room与Flow?Room Flow延迟引发竞态问题
问题
我在探索Tivi项目并实现功能时遇到了竞态条件问题。当前通过结合加载状态Flow、Room数据Flow与消息Flow生成UI状态:
val uiState = combine( loadingState.observable .onEach { Log.d("HomeVM", "LOADING $it") }, observeUserDetails.flow .onEach { Log.d("HomeVM", "USER $it") }, uiMessageManager.message .onEach { Log.d("HomeVM", "MESSAGE ${it}") }, ) { refreshing, user, message -> Log.d("HomeVM", "--- EMISSION -- LOADING: $refreshing, USER: ${user}, MESSAGE: ${message}") if (refreshing) { HomeUiState.Loading } else if(user != null){ HomeUiState.Loaded(message) } else{ HomeUiState.Error(message) } } .stateIn( scope = viewModelScope, started = SharingStarted.WhileSubscribed(), initialValue = HomeUiState.Loading )
用户Flow的实现:
class ObserveUserDetails @Inject constructor( private val repository: UserDataRepository, ) : SubjectInteractor<ObserveUserDetails.Params, User?>() { override fun createObservable(params: Params): Flow<User?> { return repository .observeUser(params.id) } data class Params(val id: Int) }
abstract class SubjectInteractor<P : Any, T> { // Ideally this would be buffer = 0, since we use flatMapLatest below, BUT invoke is not // suspending. This means that we can't suspend while flatMapLatest cancels any // existing flows. The buffer of 1 means that we can use tryEmit() and buffer the value // instead, resulting in mostly the same result. private val paramState = MutableSharedFlow<P>( replay = 1, extraBufferCapacity = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST, ) val flow: Flow<T> = paramState .distinctUntilChanged() .flatMapLatest { createObservable(it) } .distinctUntilChanged() operator fun invoke(params: P) { paramState.tryEmit(params) } protected abstract fun createObservable(params: P): Flow<T> }
刷新逻辑由Interactor实现,在ViewModel的init块中调用:设置加载状态为true,从API获取用户数据写入数据库,再将加载状态设为false:
abstract class Interactor<in P> { operator fun invoke( params: P, timeoutMs: Long = defaultTimeoutMs, ): Flow<InvokeStatus> = flow { try { withTimeout(timeoutMs) { emit(InvokeStarted) doWork(params) emit(InvokeSuccess) } } catch (t: TimeoutCancellationException) { emit(InvokeError(t)) } }.catch { t -> emit(InvokeError(t)) } suspend fun executeSync(params: P) = doWork(params) protected abstract suspend fun doWork(params: P) companion object { private val defaultTimeoutMs = TimeUnit.MINUTES.toMillis(5) } }
实际流程出现的问题:
- 刷新工作完成(获取用户并写入数据库)
- 加载状态被设为false(InvokeSuccess发射),这一操作早于Room数据库发射用户数据
- Room随后发射用户数据
导致界面短暂出现加载结束但无数据的错误界面。Tivi项目中似乎通过Store库优化,保证先写入数据库发射数据,再设置加载状态为false。我目前的解决方案是过滤Room Flow仅在用户非空时发射,同时在onStart中emit(null),但这像hack且会大量重复。调整错误界面显示条件也没解决问题,求更好的方案。
解决方案
方案1:让刷新逻辑等待Room数据发射完成
修改刷新Interactor的doWork方法,在写入数据库后,等待Room的Flow发射最新非空数据再结束。这样能保证加载状态变为false前,用户数据已经准备就绪。
示例代码:
// 刷新用的Interactor实现 class RefreshUser @Inject constructor( private val api: UserApi, private val repository: UserDataRepository ) : Interactor<RefreshUser.Params>() { override suspend fun doWork(params: Params) { // 1. 从API拉取数据 val user = api.getUser(params.id) // 2. 写入本地数据库 repository.saveUser(user) // 3. 等待Room发射最新的非空用户数据,确认写入完成 repository.observeUser(params.id) .first { it != null } } data class Params(val id: Int) }
InvokeSuccess(即加载状态设为false)会在Room发射有效用户数据之后才触发,combine时就能拿到非空的user,直接进入Loaded状态,避免错误界面闪现。
方案2:调整UI状态判断逻辑
在combine的转换函数中,针对「加载结束但数据未就绪」的场景做特殊处理,不直接跳转Error状态,而是保留Loading或显示过渡状态,直到数据到达。
修改后的combine逻辑:
val uiState = combine( loadingState.observable, observeUserDetails.flow, uiMessageManager.message, ) { refreshing, user, message -> when { refreshing -> HomeUiState.Loading user != null -> HomeUiState.Loaded(message) // 加载结束但数据为空时,继续保持Loading(或新增过渡状态) !refreshing && user == null -> HomeUiState.Loading else -> HomeUiState.Error(message) } } .stateIn( scope = viewModelScope, started = SharingStarted.WhileSubscribed(), initialValue = HomeUiState.Loading )
如果需要区分「确实无数据」和「数据未就绪」,可以在ViewModel中添加一个Flow跟踪最近一次刷新的状态,比如记录刷新是否刚完成,再结合这个状态判断是否显示Error。
方案3:复用Tivi的Store模式(推荐)
Tivi中使用的Store封装会自动处理本地与远程数据的同步顺序,确保数据写入后再触发状态更新。可以参考这个思路,把用户数据的获取、存储逻辑封装到Store中:
示例思路:
class UserStore @Inject constructor( private val api: UserApi, private val repository: UserDataRepository ) { fun observeUser(id: Int): Flow<User?> { // 本地数据Flow + 刷新Flow合并,确保刷新后数据先更新再通知 return repository.observeUser(id) .mergeWith(refreshFlow(id)) } private fun refreshFlow(id: Int): Flow<User?> = flow { val user = api.getUser(id) repository.saveUser(user) emit(user) // 发射最新数据,确保本地Flow先更新 }.catch { emit(null) } }
在ViewModel中使用Store的Flow,结合加载状态时,能保证数据更新优先于加载状态结束,从根源避免竞态问题。
内容的提问来源于stack exchange,提问作者Daniel Wilson

