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

如何创建可响应多数据源变更的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 08:30:04