【问题标题】:RxJava: buffer items until some condition is true for current itemRxJava:缓冲项目,直到当前项目的某些条件为真
【发布时间】:2016-03-09 01:21:07
【问题描述】:

这是我想弄清楚的一个 sn-p:

class RaceCondition {

    Subject<Integer, Integer> subject = PublishSubject.create();

    public void entryPoint(Integer data) {
        subject.onNext(data);
    }

    public void client() {
        subject /*some operations*/
                .buffer(getClosingSelector())
                .subscribe(/*handle results*/);
    }

    private Observable<Integer> getClosingSelector() {
        return subject /* some filtering */;
    }
}

有一个Subject 接受来自外部的事件。有一个客户订阅了这个主题,它处理事件以及buffers 它们。这里的主要思想是每次都应该根据使用流中的项目计算的某些条件发出缓冲的项目。

为此,缓冲区边界本身会监听主题。

一个重要的期望行为:每当边界发射项目时,它也应该包含在buffer 的以下发射中。当前配置并非如此,因为项目(至少我认为是这样)是从关闭选择器到达buffer之前发出的,因此它不包含在当前发射中,但是留在后面等待下一个。

有没有办法让关闭选择器等待项目首先被缓冲?如果没有,是否有另一种方法可以根据下一个传入的项目来缓冲和释放项目?

【问题讨论】:

    标签: java rx-java


    【解决方案1】:

    如果我理解正确,您希望缓冲直到某个谓词允许它基于项目。您可以使用一组复杂的运算符来做到这一点,但编写自定义运算符可能更容易:

    public final class BufferUntil<T> 
    implements Operator<List<T>, T>{
    
        final Func1<T, Boolean> boundaryPredicate;
    
        public BufferUntil(Func1<T, Boolean> boundaryPredicate) {
            this.boundaryPredicate = boundaryPredicate;
        }
    
        @Override
        public Subscriber<? super T> call(
                Subscriber<? super List<T>> child) {
            BufferWhileSubscriber parent = 
                    new BufferWhileSubscriber(child);
            child.add(parent);
            return parent;
        }
    
        final class BufferWhileSubscriber extends Subscriber<T> {
            final Subscriber<? super List<T>> actual;
    
            List<T> buffer = new ArrayList<>();
    
            /**
             * @param actual
             */
            public BufferWhileSubscriber(
                    Subscriber<? super List<T>> actual) {
                this.actual = actual;
            }
    
            @Override
            public void onNext(T t) {
                buffer.add(t);
                if (boundaryPredicate.call(t)) {
                    actual.onNext(buffer);
                    buffer = new ArrayList<>();
                }
            }
    
            @Override
            public void onError(Throwable e) {
                buffer = null;
                actual.onError(e);
            }
    
            @Override
            public void onCompleted() {
                List<T> b = buffer;
                buffer = null;
                if (!b.isEmpty()) {
                    actual.onNext(b);
                }
                actual.onCompleted();
            }
        }
    }
    

    【讨论】:

    • 非常感谢,大卫!这正是我想要的。
    • 您能否展示如何以 Observable 作为条件实现BufferUntil?缓冲直到另一个可观察对象完成?我试图让它工作,但就是不知道该怎么做......
    猜你喜欢
    • 1970-01-01
    • 2019-10-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多