如何启动可按需关闭的多Flow,实现消息列表全量查询与搜索切换
场景
List<Message>位于Composable中messages = mutableStateListOf<Message>定义在ViewModel中messages依赖getAll、search两类操作输出- 两类操作均监听SQLite数据,但同一时间仅允许一个Flow处于活跃状态(用户要么查看全量消息,要么处于搜索状态)
- SQLite搭配Flow使用是因为新消息可从网络侧同步更新
问题
- 如何在
getAll和search两类Flow之间切换? - 触发新搜索时,如何取消旧的
searchFlow并启动新的对应Flow?
解决思路
你当前的实现没有对Flow的生命周期做统一管理,会出现多个Flow同时活跃、并发修改messages的问题。可以借助flatMapLatest操作符实现Flow的自动切换与旧Flow销毁:
- 用一个状态流持有当前查询条件,
null代表查询全量消息,非空字符串代表搜索关键词 - 通过
flatMapLatest将查询条件映射为对应的getAll或searchFlow,该操作符会在新查询触发时自动取消前一个未完成的Flow收集 - 统一收集映射后的Flow更新页面状态即可
修正后代码
@Composable fun ListOfMessages() { val messages = viewModel.messages // 页面UI渲染逻辑 } // ----------------------------- // ViewModel // ----------------------------- @ExperimentalCoroutinesApi class MessageViewModel(application: Application) : AndroidViewModel(application) { val messages = mutableStateListOf<Message>() val isLoading = mutableStateOf(false) // 查询状态:null=查全量,非null=搜索对应关键词 private val queryState = MutableStateFlow<String?>(null) init { // 统一监听查询状态切换Flow queryState // 如果需要搜索防抖可加这行,单位毫秒,避免输入过程频繁触发请求 // .debounce(300) .flatMapLatest { query -> if (query == null) { MessageUseCase(messagesDB).getAll() } else { MessageUseCase(messagesDB).search(query) } } .onEach { dataState -> isLoading.value = dataState.loading dataState.data?.let { data -> messages.clear() messages.addAll(data) } dataState.error?.let { error -> // 处理错误逻辑 } } .launchIn(viewModelScope) } // 切换到全量列表 fun fetchMessages() { queryState.value = null } // 触发搜索 fun searchWithInMessage(q: String) { queryState.value = q } } // ----------------------------- // Use Cases(修正原代码多余的一层大括号) // ----------------------------- class MessageUseCase(private val messagesDB: messageDao) { @ExperimentalCoroutinesApi fun getAll(): Flow<DataState<List<Message>>> = channelFlow { send(DataState.loading()) try { fetchAndSaveLatestMessagesFromRemote() val messages = messagesDB.getAllStream() messages.collectLatest { list -> // 业务逻辑处理 send(DataState.success(list)) } } catch (e: Exception){ send(DataState.error<List<Message>>(e.message?: "Unknown Error")) } } @ExperimentalCoroutinesApi fun search(q: String): Flow<DataState<List<Message>>> = channelFlow { send(DataState.loading()) try { fetchAndSaveSearchedMessageFromRemote(q) val messages = messagesDB.searchStream(q) messages.collectLatest { list -> // 业务逻辑处理 send(DataState.success(list)) } } catch (e: Exception){ send(DataState.error<List<Message>>(e.message?: "Unknown Error")) } } }
内容的提问来源于stack exchange,提问作者clamentjohn
相关产品推荐
相关产品推荐

