【问题标题】:Emit items in RxJava in an interval which depends on the emitted item itself在 RxJava 中以取决于发射项目本身的间隔发射项目
【发布时间】:2014-10-24 17:18:53
【问题描述】:

在 RxJava for Android 中,我想在取决于项目本身的间隔内发出项目:在 Observable 中,我从队列中拉出一个项目,处理它并发出它。根据项目的类型,我想调整下一个项目的发射时间(减慢或加快间隔)。

@a.bertucci 在此处Emit objects for drawing in the UI in a regular interval using RxJava on Android 提出的以下代码演示了如何定期发射项目。

private void drawPath(final String chars) {
    Observable.zip(
        Observable.create(new Observable.OnSubscribe<Path>() {
            // all the drawing stuff here
            ...
        }),
        Observable.timer(0, 50, TimeUnit.MILLISECONDS),
        new Func2<Path, Long, Path>() {
            @Override
            public Path call(Path path, Long aLong) {
                return path;
            }
        }
    )
    .subscribeOn(Schedulers.newThread())
    .observeOn(AndroidSchedulers.mainThread())
    ...
}

我现在的问题是,是否有可能在可观察对象发射时修改发射频率,以及使用 RxJava 的首选实现是什么。

【问题讨论】:

    标签: java android rx-java


    【解决方案1】:

    你可以使用

     public final <U, V> Observable<T> delay(
                Func0<? extends Observable<U>> subscriptionDelay,
                Func1<? super T, ? extends Observable<V>> itemDelay)
    

    public final <U> Observable<T> delay(Func1<? super T, ? extends Observable<U>> itemDelay)
    

    你可以使用itemDelay来控制速度。

    【讨论】:

    • 这不是将整个序列移动了一定的延迟吗?这不是我需要的。我不想改变顺序。我希望立即发出第一项。然后(或之前)发射该项目,我想查看我发射的项目并决定何时发射下一个项目(取决于项目)。然后发出下一个项目,并为之后的项目再次更改时间。
    【解决方案2】:

    根据您的评论,我认为您应该为此建立一个新的运营商。

    该运算符将采用一个函数来计算应用到发出下一个项目的延迟

     Observable.range(1, 1000).lift(new ConfigurableDelay((item) -> 3 SECONDS)
                              .subscribe();
    

    你可以试试这样的:

    public class ConfDelay {
    
    public static void main(String[] args) {
        Observable.range(1, 1000).lift(new ConfigurableDelay(ConfDelay::delayPerItem, Schedulers.immediate()))
                .map((i) -> "|")
                .subscribe(System.out::print);
    }
    
    
    public static TimeConf delayPerItem(Object item) {
        long value = ((Integer) item).longValue();
        return new TimeConf(value * value, TimeUnit.MILLISECONDS);
    }
    
    private static class TimeConf {
        private final long time;
        private final TimeUnit unit;
    
        private TimeConf(final long time, final TimeUnit unit) {
            this.time = time;
            this.unit = unit;
        }
    }
    
    private static class ConfigurableDelay<T> implements Observable.Operator<T, T> {
        private final Func1<T, TimeConf> itemToTime;
        private final Scheduler scheduler;
    
        public ConfigurableDelay(final Func1<T, TimeConf> itemToTime) {
            this(itemToTime, Schedulers.computation());
        }
    
        public ConfigurableDelay(final Func1<T, TimeConf> itemToTime, final Scheduler scheulder) {
            this.itemToTime = itemToTime;
            this.scheduler = scheulder;
        }
    
        @Override
        public Subscriber<? super T> call(final Subscriber<? super T> subscriber) {
            return new Subscriber<T>(subscriber) {
    
                private TimeConf nextTime = null;
    
                @Override
                public void onCompleted() {
                    subscriber.onCompleted();
                }
    
                @Override
                public void onError(final Throwable e) {
                    subscriber.onError(e);
                }
    
                @Override
                public void onNext(final T t) {
                    TimeConf previousNextTime = nextTime;
                    this.nextTime = itemToTime.call(t);
                    if (previousNextTime == null) {
                        subscriber.onNext(t);
                    } else {
                        scheduler.createWorker().schedule(() -> subscriber.onNext(t), previousNextTime.time, previousNextTime.unit);
                    }
                }
            };
        }
    }
    }
    

    【讨论】:

    • 我认为我的代码可以用这样的 zip / flatmap 代替:obs.map(e -> new Pair(e, 3 SECONDS)) .flatMap((conf) -> Observable。只是(conf.event).delay(conf.time)) .subscribe();
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-02-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多