引起:rx.exceptions.MissingBackpressureException

Leo*_*Leo 1 android kotlin rx-java

我还有一个问题。这次我Caused by: rx.exceptions.MissingBackpressureException在执行这段代码时遇到了这个错误:

class UpdateHelper {
val numberOfFileToUpdate: PublishSubject<Int>

init {
    numberOfFileToUpdate = PublishSubject.create()
}

public fun startUpdate(): Observable<Int>{
    return getProducts().flatMap { products: ArrayList<Product> ->
            numberOfFileToUpdate.onNext(products.size)
            return@flatMap saveRows(products)
        }
}

private fun getProducts(): Observable<ArrayList<Product>> {
    return Observable.create {
        var products: ArrayList<Product> = ArrayList()
        var i = 0
        while (i++ < 100) {
            products.add(Product())
        }

        it.onNext(products)
        it.onCompleted()
    }
}


private fun saveRows(products: ArrayList<Product>): Observable<Int> {
    return Observable.create<Int> {
        var totalNumberOfRow = products.size

        while (totalNumberOfRow-- > 0){
            it.onNext(products.size - totalNumberOfRow)
            Thread.sleep(100)
        }
        it.onCompleted()
    }
}
Run Code Online (Sandbox Code Playgroud)

}

代码只是两个进程的测试代码。第一个过程Product从网络获取一个列表,然后将这些产品持久化到应用程序内的本地数据库中。这是主要思想。

该方法getProducts完成获取数据的工作,在本例中,我只创建一个包含 100 个产品的 ArrayList。该saveRows做的工作仍然存在。

这些saveRows方法发出Int代表已保存行的突出物。我这样做是因为在 UI 中我有一个进度条报告进度。

我从应用程序的另一个角度调用该方法startUpdate,在发出一些项目后,我得到了描述异常

at com.techbyflorin.rockapan.helpers.UpdateHelper$saveRows$1.call(UpdateHelper.kt:46)

at com.techbyflorin.rockapan.helpers.UpdateHelper$saveRows$1.call(UpdateHelper.kt:40)
Run Code Online (Sandbox Code Playgroud)

我明白为什么会发生这个异常,https://github.com/ReactiveX/RxJava/wiki/Backpressure但我不知道我做错了什么或如何解决它。

任何人都可以就此给我建议。

pt2*_*121 5

问题是您的 Observable 源发出的速度比消费者消耗的速度快。保存每个产品需要 100 毫秒。您可以添加 onBackpressureBuffer()。

UpdateHelper().startUpdate()
    .onBackpressureBuffer() // Add this
    .subscribeOn(Schedulers.io())
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe({
      Log.d(TAG, "next $it")
    }, {
      Log.d(TAG, it.message)
    }, {
    })
Run Code Online (Sandbox Code Playgroud)

此外,您可以尝试删除Thread.sleep(100).

flatmap 使用OperatorMerge ( merge(map(func))):您可以看到,在您的情况下,maponNexts 的发送速度比请求的要快。