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

如何正确结合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)
    }
}

实际流程出现的问题:

  1. 刷新工作完成(获取用户并写入数据库)
  2. 加载状态被设为false(InvokeSuccess发射),这一操作早于Room数据库发射用户数据
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 18:40:56