【问题标题】:Why is this happening in with onBackpressureDrop () in RxJava为什么在 RxJava 中使用 onBackpressureDrop () 会发生这种情况
【发布时间】:2020-10-21 04:12:59
【问题描述】:

这里我有一个 flowable,它每毫秒发射一次元素。

   Flowable<Long> source = Flowable.interval(1,TimeUnit.MILLISECONDS).take(14000);
        source.map(e->{
            Log.d("TAGBefore","before " + e);
            return e;
        })
        .onBackpressureDrop()
        .observeOn(Schedulers.computation())
        .subscribe(
                        e-> {
                            Log.d("TAGNext","onNext: " + e);
                            Thread.sleep(100);
                        },
                        e-> Log.d("TAGError","error: " + e),
                        ()-> Log.d("TAGComplete","onComplete")
        );

我使用 before 来知道 observable 发射元素的时刻,我怀疑这里从 127(当观察者满时)到 9688

   TAGNext: onNext: 125
   TAGNext: onNext: 126
   TAGNext: onNext: 127
   TAGNext: onNext: 9668
   TAGNext: onNext: 9669
   TAGNext: onNext: 9670

但是,当我更多地检查控制台(使用其他搜索过滤器)时,我意识到在发出 127 时它已经转到 12794,所以它不应该是 12794 或接近的数字而不是 9688? ,谢谢。

   TAGBefore: before: 12793
   TAGBefore: before: 12794
   TAGNext:   onNext: 127
   TAGBefore: before: 12795
   TAGBefore: before: 12796
 

但是,当我更多地检查控制台(使用其他搜索过滤器)时,我意识到当发出 127 时它已经为 12794,因此它不应该是 12794 或接近的数字,而不是 9688已经免费了?,我澄清一下我是 RxJava 的新手,以防我说错了,谢谢。

【问题讨论】:

    标签: rx-java reactive-programming rx-java2 rx-android rx-java3


    【解决方案1】:

    observeOn 有一个默认的 128 元素缓冲区,很快就会填满。

    它每 100 毫秒被耗尽,直到它只剩下 32 个元素,此时它又请求 96 个。因此,更多项目通过onBackpressureDrop 需要大约 9600 毫秒,因此您会看到TAGNext: onNext: 9668

    当第 127 个元素被耗尽时,运行大约需要 12700 毫秒,因此您从生产者端看到 TAGBefore: before: 12795

    【讨论】:

      猜你喜欢
      • 2013-08-02
      • 2018-12-19
      • 2021-01-24
      • 1970-01-01
      • 2021-07-19
      • 2013-06-09
      • 2021-04-30
      • 2010-09-22
      相关资源
      最近更新 更多