【问题标题】:Room RxJava observable triggered multiple times on insertRoom RxJava observable 在插入时触发多次
【发布时间】:2020-04-20 13:54:50
【问题描述】:

我的存储库实现有一个奇怪的问题。每次我调用应该从数据库中获取数据并通过网络调用更新数据库的函数时,我都会从数据库观察者那里收到多个结果。

override fun getApplianceControls(
    serialNumber: SerialNumber
): Flowable<ApplianceControlState> {
    val subject = BehaviorProcessor.create<ApplianceControlState>()

    controlsDao.get(serialNumber.serial)
        .map { controls ->
            ApplianceControlState.Loaded(controls.toDomainModel())
        }
        .subscribe(subject)

    controlApi.getApplianceControls(serialNumber.serial)
        .flatMapObservable<ApplianceControlState> { response ->
            val entities = response.toEntity(serialNumber)
            // Store the fetched controls on the database.
            controlsDao.insert(entities).andThen(
                // Return an empty observable because the db will take care of emitting latest values.
                Observable.create { }
            )
        }
        .onErrorResumeNext { error: Throwable ->
            Observable.create { emitter -> emitter.onNext(ApplianceControlState.Error(error)) }
        }
        .subscribeOn(backgroundScheduler)
        .subscribe()


    return subject.distinctUntilChanged()
}
@Dao
interface ApplianceControlsDao {

    @Insert(onConflict = OnConflictStrategy.REPLACE)
    fun insert(controls: List<TemperatureControlEntity>): Completable

    @Query("SELECT * FROM control_temperature WHERE serial = :serial")
    fun get(serial: String): Flowable<List<TemperatureControlEntity>>
}

基本上,如果我调用getApplianceControls 一次,我会得到想要的结果。然后我再次调用,另一个序列号是空的,我得到了空数组。 然后我第三次调用,但序列号与第一次相同,在插入调用后我得到正确结果和空数组的混合。

像这样:

第一次调用,序列号“123” -> 加载([control1,control2,control3])

第二次调用,序列号“000” -> Loaded([])

第三次调用,序列号“123” -> Loaded([control1, control2, control3]), Loaded([]), Loaded([control1, control2, control3])

如果我从 api 响应中删除 db insert,它可以正常工作。在调用insert 之后,所有奇怪的事情都会发生。

编辑:getApplianceControls() 是从 ViewModel 调用的。

fun loadApplianceControls(serialNumber: SerialNumber) {
    Log.i("Loading appliance controls")

    applianceControlRepository.getApplianceControls(serialNumber)
        .subscribeOn(backgroundScheduler)
        .observeOn(mainScheduler)
        .subscribeBy(
            onError = { error ->
                Log.e("Error $error")
            },
            onNext = { controlState ->
                _controlsLiveData.value = controlState  
            }
        ).addTo(disposeBag)
}

【问题讨论】:

  • 是否尝试用 Single 替换 Flowable?当您发出单个数据时
  • @MustafaKhaled 是的,我试过了,但从长远来看,我希望这是一个可流动/可观察的,因为数据集可能会因来自服务器的推送通知而改变
  • 首先,您有 2 个未在任何地方取消订阅的订阅。这可能会导致内存泄漏。你能说明你在哪里使用这个fun getApplianceControls()吗?
  • @borichellow 他们是使用主题订阅订阅的,我相信这会导致它们在主题也被处理时被处理(如果我错了,请纠正我)。我将编辑问题以添加调用函数的位置

标签: android rx-java rx-java2 android-room


【解决方案1】:

正如我在评论中提到的,您有 2 个未在任何地方取消订阅的订阅,这可能会导致内存泄漏(在释放主题时它不会释放),而且通过这种实现,您会忽略 API 错误。 我会尝试将其更改为:

override fun getApplianceControls(serialNumber: SerialNumber): Flowable<ApplianceControlState> {

    val dbObservable = controlsDao.get(serialNumber.serial)
        .map { controls ->
            ApplianceControlState.Loaded(controls.toDomainModel())
        }

    val apiObservable = controlApi.getApplianceControls(serialNumber.serial)
        .map { response ->
            val entities = response.toEntity(serialNumber)
           // Store the fetched controls on the database.
           controlsDao.insert(entities).andThen( Unit )
        }
        .toObservable()
        .startWith(Unit)

    return Observables.combineLatest(dbObservable, apiObservable) { dbData, _ -> dbData }
        // apiObservable emits are ignored, but it will by subscribed with dbObservable and Errors are not ignored 
        .onErrorResumeNext { error: Throwable ->
            Observable.create { emitter -> emitter.onNext(ApplianceControlState.Error(error)) }
        }
        .subscribeOn(backgroundScheduler)
        //observeOn main Thread
        .distinctUntilChanged()
}

我不确定它是否能解决原来的问题。但如果是这样 - 问题出在flatMapObservable 查看controlApi.getApplianceControls() 的实现也很有用。

【讨论】:

  • controlApi.getApplianceControls() 来自 Retrofit,它只返回一个 Single&lt;ApplianceControlData&gt;。我会试试你的实现。
  • 那么您需要添加.toObservable() 才能使我的代码正常工作。我编辑了答案并将startWith() 添加到 apiObservable,因为combineLatest 将等待两个 Observable 的第一次发射。
  • 通过这个实现,我只收到一次结果。这就像 db observable 被卡住了。我认为这与我原来的问题的根本原因相同。每次我调用 DAO 函数时,它都会与最后一次调用堆叠在一起,并且我不断收到多个结果。但是有了这个实现,它就像是在重用同一个 observable,但它被卡住了。
  • 那么我猜问题出在 DB 部分,我对 Room 不熟悉,抱歉
  • 对您的解决方案稍作调整后,它终于奏效了。我想我做得不对,你没有正确退订是对的。谢谢。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-06-11
  • 1970-01-01
  • 1970-01-01
  • 2023-03-30
  • 1970-01-01
相关资源
最近更新 更多