如何创建可响应多数据源变更的Flowable实现聊天数据实时更新
问题根因
你现有代码无法收到后续更新的核心原因有三点:
getLastUnreadMessage和getUnreadCount当前返回的是单次请求的Single类型,查询完成后流就终止,无法监听后续数据变更flatMapSingle+toList()的组合需要等待所有联系人的单次查询都完成后才会发射一次列表,之后整个流直接终止,不会保留监听- 原有逻辑只能响应联系人列表的变更,单个联系人的未读消息、未读计数变更无法触发流的重新发射
优化实现方案
前提要求
首先需要把三类数据源的DAO查询都改成返回Flowable类型,这样数据库数据变更时会自动发射最新值,示例DAO定义:
// 监听联系人列表变更 @Query("SELECT * FROM contact") fun queryContactsFlowable(): Flowable<List<Contact>> // 监听单个联系人的最新未读消息变更 @Query("SELECT * FROM messages WHERE contact_id = :contactId ORDER BY create_time DESC LIMIT 1") fun getLastUnreadMessage(contactId: Long): Flowable<Messages?> // 监听单个联系人的未读消息计数变更 @Query("SELECT COUNT(*) FROM messages WHERE contact_id = :contactId AND is_read = 0") fun getUnreadCount(contactId: Long): Flowable<Int>
核心流实现
fun queryAllChats(): Flowable<List<Chat>> = dao.queryContactsFlowable() // 联系人列表变更时,重建所有联系人的监听流,丢弃旧的无用监听 .switchMap { contacts -> if (contacts.isEmpty()) { Flowable.just(emptyList()) } else { // 为每个联系人生成独立的监听流,未读消息/计数任意变更就发射新的Chat对象 val perContactFlow = contacts.map { contact -> Flowable.combineLatest( getLastUnreadMessage(contact.id), getUnreadCount(contact.id) ) { lastMsg, unreadCount -> Chat(contact, lastMsg, unreadCount) } // 同一个Chat内容无变更时不重复发射,减少无用更新 .distinctUntilChanged() } // 合并所有联系人的流,任意一个联系人数据变更就发射最新的完整列表 Flowable.combineLatest(perContactFlow) { array -> array.map { it as Chat } } } } // 整个列表内容无变更时不重复发射 .distinctUntilChanged() // 所有数据库操作都放在IO线程,避免阻塞UI .subscribeOn(schedulers.io) // 最终结果切换到主线程再分发 .observeOn(schedulers.main)
ViewModel调用
原有调用逻辑不需要修改:
val chatListLiveData = LiveDataReactiveStreams.fromPublisher(queryAllChats())
额外优化建议
- 实现仅更新必要数据:UI层配合RecyclerView的
DiffUtil实现局部刷新,无需全局重绘列表,示例差异校验逻辑:
class ChatDiffCallback : DiffUtil.ItemCallback<Chat>() { override fun areItemsTheSame(oldItem: Chat, newItem: Chat): Boolean { return oldItem.contact.id == newItem.contact.id } override fun areContentsTheSame(oldItem: Chat, newItem: Chat): Boolean { return oldItem == newItem } }
- 避免频繁刷新:如果数据变更频率过高,可以在流中添加
.onBackpressureLatest()背压策略,只保留最新的列表变更,也可以加throttleLast(300, TimeUnit.MILLISECONDS)合并短时间内的多次变更,减少UI刷新次数。
内容的提问来源于stack exchange,提问作者Monan Rise
相关产品推荐
相关产品推荐

