How to create blocking backpressure with rxjava Flowables?

Viewed 817

I have a Flowable that we are returning in a function that will continually read from a database and add it to a Flowable.

public void scan() {
    Flowable<String> flow = Flowable.create((FlowableOnSubscribe<String>) emitter -> {
        Result result = new Result();
        while (!result.hasData()) {
            result = request.query(skip, limit);
            partialResult.getResult()
                .getFeatures().forEach(feature -> emmitter.emit(feature));
        }
    }, BackpressureStrategy.BUFFER)
        .subscribeOn(Schedulers.io());
    return flow;
}

Then I have another object that can call this method.

myObj.scan()
        .parallel()
        .runOn(Schedulers.computation())
        .map(feature -> {
            //Heavy Computation
        })
        .sequential()
        .blockingSubscribe(msg -> {
            logger.debug("Successfully processed " + msg);
        }, (e) -> {
            logger.error("Failed to process features because of error with scan", e);
        });

My heavy computation section could potentially take a very long time. So long in fact that there is a good chance that the database requests will load the whole database into memory before the consumer finishes the first couple entries.

I have read up on backpressure with rxjava but the only 4 options essentially make me drop data or replace it with the last.

Is there a way to make it so that when I call emmitter.emit(feature) the call blocks until there is more room in the Flowable?

I.E I want to treat the Flowable as a blocking queue where push will sleep if the queue is past the capacity.

0 Answers
Related