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

预期未抛出EmptyResultSetException的原因排查(Room+RxJava场景)

问题根源分析

你遇到的核心问题是Room的Observable查询行为特性:

  • 当@Query返回Observable<T>时,Room会创建一个持续订阅的流:
    1. 首次订阅时执行查询并发射结果;
    2. 当关联表数据更新时,自动重新执行查询并发射新结果;
    3. 如果查询结果为空,不会发射onError(包括EmptyResultSetException),也不会发射onComplete,而是保持订阅状态等待新数据。

你的流原本依赖EmptyResultSetException来判断所有邮件发送完毕,但这个异常只会在Single<T>类型的一次性查询无结果时抛出,Room自动触发的Observable查询不会抛出该异常,所以流会卡在最后一步,无法触发完成逻辑。


解决方案

我们需要调整DAO返回类型和流逻辑,让系统能明确检测到“无待发送消息”的状态,从而终止流。

步骤1:修改DAO查询返回类型

将getMessageToSend()改为返回Maybe<MessageToSend>(单次查询,有结果发射onSuccess,无结果发射onComplete),同时新增一个轻量级查询判断是否还有待发送消息:

// 原查询:获取单条待发送消息(确保SQL包含LIMIT 1)
@Query("SELECT * FROM your_message_table WHERE sent = false LIMIT 1")
fun getMessageToSend(): Maybe<MessageToSend>

// 新增:检查是否还有待发送消息(性能开销极小)
@Query("SELECT EXISTS(SELECT 1 FROM your_message_table WHERE sent = false)")
fun hasPendingMessages(): Single<Boolean>

步骤2:重构Rx流逻辑

改用Flowable配合repeatWhen实现循环获取+处理,直到无待发送消息:

Flowable.defer { 
    // 每次重复时重新查询待发送消息
    App.context.repository.getMessageToSend().toFlowable()
}
.flatMap { messageToSend ->
    // 发送邮件 + 延迟 + 更新状态的逻辑保持不变
    App.context.repository.sendMessage(messageToSend)
        .doOnError { messageToSend.failureSending = true }
        .zipWith(
            Flowable.interval(1, TimeUnit.SECONDS),
            BiFunction { item: MessageToSend, _: Long -> item }
        )
        .flatMap { msg ->
            App.context.repository.storeMessageSent(msg)
                .doOnError { msg.failureSending = true }
                .toFlowable()
        }
}
.repeatWhen { completed ->
    // 当一条消息处理完成后,检查是否还有待发送消息
    completed.delay(100, TimeUnit.MILLISECONDS) // 等待数据库更新生效
        .flatMapPublisher {
            App.context.repository.hasPendingMessages()
                .flatMapPublisher { hasPending ->
                    if (hasPending) {
                        // 还有消息,触发重复查询
                        Flowable.just(0)
                    } else {
                        // 无消息,终止流
                        Flowable.empty()
                    }
                }
        }
}
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe(
    { /* 单条消息处理完成的回调 */ },
    { ex -> 
        // 处理发送/更新过程中的错误
        when(ex) {
            is EmptyResultSetException -> {}
            else -> {}
        }
    },
    { 
        // 所有消息发送完成!这里执行收尾逻辑
    }
)

可选简化方案:主动循环获取

如果不需要严格依赖Room的自动更新,可以改用generate运算符实现主动循环:

Observable.generate<Boolean, MessageToSend> { emitter ->
    // 每次迭代查询一条待发送消息
    App.context.repository.getMessageToSend()
        .subscribe(
            { msg -> emitter.onNext(msg) },
            { ex -> emitter.onError(ex) },
            { 
                // 无消息,终止循环
                emitter.onComplete()
            }
        )
}
.flatMap { messageToSend ->
    // 发送+延迟+更新逻辑同前
    App.context.repository.sendMessage(messageToSend)
        .doOnError { messageToSend.failureSending = true }
        .zipWith(
            Observable.interval(1, TimeUnit.SECONDS),
            BiFunction { item: MessageToSend, _: Long -> item }
        )
        .flatMap { msg ->
            App.context.repository.storeMessageSent(msg)
                .doOnError { msg.failureSending = true }
        }
}
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe(
    {},
    { ex -> /* 错误处理 */ },
    { /* 所有消息完成 */ }
)

关键注意点

  • Room的Observable类型查询是持续订阅流,仅用于监听数据变化,不适合“一次性循环获取数据直到为空”的场景;
  • Maybe<T>类型更适合单次查询:有结果则发射数据,无结果则触发onComplete,能明确区分“有数据”和“无数据”状态;
  • 新增的hasPendingMessages()查询是轻量级的存在性检查,性能开销可以忽略。

内容的提问来源于stack exchange,提问作者Johann

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 09:14:33