RxJava 2 - Omit interval emissions while processing

Viewed 112

Suppose I have the following task chain to execute periodically.

Observable.interval(0, 10, TimeUnit.SECONDS)
    .flatMapIterable(l -> provider.getAssets())
    .map(renderer::render)
    .map(converter::convert)
    .subscribe(provider::publish)

Although each task can take much more then 10 seconds to complete. I want to call getAssets() every 10 seconds, process, but don't want to catch up with emissions that occurred since processing started. How do I omit them? Thanks.

2 Answers

I think you can solve this by adding "in progress flag". In your case it would look something like this:

AtomicBoolean inProgress = new AtomicBoolean(false);
Observable.interval(0, 10, TimeUnit.SECONDS)
    .filter(emission -> !inProgress.get())
    .flatMapIterable(l -> provider.getAssets())
    .map(renderer -> {
        // set to true when processing started
        inProgress.set(true);
        return renderer.render();
    })
    .map(converter::convert)
    .subscribe(provider -> { 
        // set to false when processing finished
        provider.publish();
        inProgress.set(false);
    });

It turns out I was after backpressure operators, specifically onBackpressureDrop to omit the latest values.

Flowable.interval(1, TimeUnit.MINUTES)
     .onBackpressureDrop()
     .observeOn(Schedulers.io())
     .doOnNext(e -> networkCall.doStuff())
     .subscribe(v -> { }, Throwable::printStackTrace);
Related