【问题标题】:RxJava: count emitted elements during emittingRxJava:在发射期间计算发射的元素
【发布时间】:2017-05-16 14:16:41
【问题描述】:

我想知道哪个项目刚刚被一个阻塞的长时间运行的 observable 发出,它将发出数千个项目。下面的代码有效,但它创建了来自range() 的巨大缓冲区。

sourceOservable
.zipWith(Observable.range(0, Integer.MAX_VALUE), (any, counter) -> counter)
.whatever(...)

有什么方法可以避免这种行为而不引入任何外部计数器字段?

【问题讨论】:

    标签: rx-java


    【解决方案1】:

    缓冲区归因于 Observable.range。它的产生速度可能比 sourceObservable 快。它必须缓冲所有值才能使用 sourceObservable 中的正确值进行压缩。

    请看一下我的实现:

    @Test
    void stackoverflow44004014() {
        Observable.just("i", "b", "c")
                .scan(0, (counter, sourceValue) -> {
                    return ++counter;
                })
                .skip(1)
                .test()
                .assertResult(1, 2, 3);
    }
    

    【讨论】:

    • 那是我一直在寻找的。谢谢你。我只需要将计数器更改为包含源值和 int 的对,以便能够将源值与计数器一起沿流进一步传播。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-05-23
    相关资源
    最近更新 更多