【问题标题】:Transform observable to have a minimum delay between values转换 observable 以使值之间的延迟最小
【发布时间】:2016-05-25 19:25:15
【问题描述】:

我正在寻找this question 的一个很好的解决方案,但在 RxJava 中实现。这个问题也超过五年了,所以我想知道 - 有没有更好的方法来实现这个输出?

我想要实现的是缓冲来自某些的传入事件 IObservable(它们会爆发)并进一步释放它们,但是一个 一个,以均匀的间隔。像这样:

-oo-ooo-oo------------------oooo-oo-o-------------->

-o--o--o--o--o--o--o--------o--o--o--o--o--o--o---->

对我来说最大的要求是不丢失任何可观察对象,并且事件的顺序保持不变。

【问题讨论】:

    标签: java rx-java observable


    【解决方案1】:

    此特定模式需要记住上一个事件的安排时间,因此如果下一个事件在一段时间之后出现,它可以立即发出并开始新的周期性发出。也许更简单有效的方法是编写自定义运算符:

    import java.util.concurrent.TimeUnit;
    
    import rx.*;
    import rx.Observable.Operator;
    import rx.schedulers.Schedulers;
    
    public class SpanOut<T> implements Operator<T, T> {
        final long time;
    
        final TimeUnit unit;
    
        final Scheduler scheduler;
    
        public SpanOut(long time, TimeUnit unit, Scheduler scheduler) {
            this.time = time;
            this.unit = unit;
            this.scheduler = scheduler;
        }
    
        @Override
        public Subscriber<? super T> call(Subscriber<? super T> t) {
            Scheduler.Worker w = scheduler.createWorker();
    
            SpanSubscriber<T> parent = new SpanSubscriber<>(t, unit.toMillis(time), w);
    
            t.add(w);
            t.add(parent);
    
            return parent;
        }
    
        static final class SpanSubscriber<T> extends Subscriber<T> {
            final Subscriber<? super T> actual;
    
            final long spanMillis;
    
            final Scheduler.Worker worker;
    
            long lastTime;
    
            public SpanSubscriber(Subscriber<? super T> actual, 
                    long spanMillis, Scheduler.Worker worker) {
                this.actual = actual;
                this.spanMillis = spanMillis;
                this.worker = worker;
            }
    
            @Override
            public void onNext(T t) {
                long now = worker.now();
                if (now >= lastTime + spanMillis) {
                    lastTime = now + spanMillis;
                    worker.schedule(() -> {
                        actual.onNext(t);
                    });
                } else {
                    long next = lastTime - now;
                    lastTime += spanMillis;
                    worker.schedule(() -> {
                        actual.onNext(t);
                    }, next, TimeUnit.MILLISECONDS);
                }
            }
    
            @Override
            public void onError(Throwable e) {
                worker.schedule(() -> {
                    actual.onError(e);
                    unsubscribe();
                });
            }
    
            @Override
            public void onCompleted() {
                long next = lastTime - worker.now();
                worker.schedule(() -> {
                    actual.onCompleted();
                    unsubscribe();
                }, next, TimeUnit.MILLISECONDS);
            }
    
            @Override
            public void setProducer(Producer p) {
                actual.setProducer(p);
            }
        }
    
        public static void main(String[] args) {
            Observable.range(1, 5)
            .concatWith(Observable.just(6).delay(6500, TimeUnit.MILLISECONDS))
            .concatWith(Observable.range(7, 4))
            .lift(new SpanOut<>(1, TimeUnit.SECONDS, Schedulers.computation()))
            .timeInterval()
            .toBlocking()
            .subscribe(System.out::println);
        }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2013-03-19
      • 2020-11-17
      • 1970-01-01
      • 2020-10-02
      • 1970-01-01
      • 2020-12-22
      • 2015-04-18
      相关资源
      最近更新 更多