【问题标题】:RXJava PausableBufferRXJava PausableBuffer
【发布时间】:2016-04-03 12:07:21
【问题描述】:

您好,我是 RXJava 的新手。 我正在寻找一种可观察的解决方案,该解决方案将根据收到的项目继续并暂停发射项目。

假设我们的条件是这个整数谓词:

Func1<Integer, Boolean> isOdd = number -> number%2==1;

当我们像myNumberSubject.onNext(someInt);这样向主题添加数字时,在添加奇数之前添加的所有数字都存储在缓冲区中,但是添加第一个奇数的数字,缓冲区中的所有数字都会发出一次爆发(包括奇怪的项目)。

之后,只要是奇数,就会一个一个地发出每个数字。当添加偶数时再次将其放入缓冲区。

我已经能够找到任何实际的工作示例,但是这个 pausableBuffer 的弹珠示例可能有可能完全按照我的意愿去做。 http://rxmarbles.com/#pausableBuffered 我希望有一些预先存在的 RXJava 解决方案可以解决问题。这是我自己的 hack 工作解决方案。

public class PausableBuffer<R> {

    private boolean isPaused;
    private List<R> buffer;
    private ReplaySubject<R> regulatedSubject;

    private PausableBuffer(){
        regulatedSubject = ReplaySubject.create();
        buffer=new LinkedList<>();
    }

    public static <R>Observable<R> create(Observable<R> observable, Func1<R, Boolean> continueCondition){
        PausableBuffer<R> pausableBuffer = new PausableBuffer<>();
        observable.subscribe(value -> {
            synchronized(pausableBuffer) {
                if(pausableBuffer.isPaused){
                    pausableBuffer.buffer.add(value);
                    if(continueCondition.call(value)){
                        pausableBuffer.isPaused=false;
                        for (R r : pausableBuffer.buffer) {
                            pausableBuffer.regulatedSubject.onNext(r);
                        }
                        pausableBuffer.buffer.clear();
                    }
                }else{
                    if(continueCondition.call(value)){
                        pausableBuffer.regulatedSubject.onNext(value);
                    }else{
                        pausableBuffer.isPaused=true;
                        pausableBuffer.buffer.add(value);
                    }
                }
            }
        });
        return pausableBuffer.regulatedSubject.asObservable();
    }

    public static void main(String[] args) {
        BehaviorSubject<Integer> behaviorSubject = BehaviorSubject.create();
        Observable<Integer> observable = PausableBuffer.<Integer>create(
                behaviorSubject.asObservable(),                    
                intValue -> intValue==5 || 6<intValue);//continueCondition
        observable.subscribe(v -> System.out.print(v+", "));
        for (int i = 0; i <= 8; i++) {
            System.out.print("adding " + i + " : ");
            behaviorSubject.onNext(i);
            System.out.println();
        }
    }
}

打印出来:

  • 加0:
  • 加1:
  • 添加2:
  • 添加 3:
  • 添加4:
  • 加 5 : 0, 1, 2, 3, 4, 5,
  • 加6:
  • 添加 7 : 6, 7,
  • 加 8 : 8,

【问题讨论】:

标签: java rx-java


【解决方案1】:

第二版:

public final class ContinueWhile<T> implements Observable.Operator<T, T> {

final Func1<T, Boolean> continuePredicate;

private ContinueWhile(Func1<T, Boolean> continuePredicate) {
    this.continuePredicate = continuePredicate;
}

public static <T>ContinueWhile<T> create(Func1<T, Boolean> whileTrue){
    return new ContinueWhile<>(whileTrue);
}

@Override
public Subscriber<? super T> call(Subscriber<? super T> child) {
    ContinueWhileSubscriber parent = new ContinueWhileSubscriber(child);
    child.add(parent);
    return parent;
}

final class ContinueWhileSubscriber extends Subscriber<T> {

    final Subscriber<? super T> actual;
    Deque<T> buffer = new ConcurrentLinkedDeque<>();

    public ContinueWhileSubscriber(Subscriber<? super T> actual) {
        this.actual = actual;
    }

    @Override
    public void onNext(T t) {
        buffer.add(t);
        if (continuePredicate.call(t)) {
            while(!buffer.isEmpty())
                actual.onNext(buffer.poll());
        }
    }

    @Override
    public void onError(Throwable e) {
        buffer = null;
        actual.onError(e);
    }

    @Override
    public void onCompleted() {
        while (!buffer.isEmpty())
            actual.onNext(buffer.poll());
        buffer=null;
        actual.onCompleted();
    }
}
}




public static void main(String[] args) {
    BehaviorSubject<Integer> subject = BehaviorSubject.create();
    subject.asObservable()
            .doOnNext(v -> System.out.print("next "))
            .lift(ContinueWhile.create(i -> i%3==0))
            .subscribe(v -> System.out.print(v + ", "));
    for (int i = 0; i < 10; i++) {
        subject.onNext(i);
    }
}
}

感谢 AndroidEx。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2014-10-27
    • 1970-01-01
    • 1970-01-01
    • 2016-02-14
    • 2021-03-21
    • 2019-08-13
    • 2015-09-30
    • 1970-01-01
    相关资源
    最近更新 更多