Controlling back pressure in Spring cloud stream

Viewed 97

Is there a way to control the amount data that can be processed by a flux. Ex: Lets say I get millions of messages in a stream and I want to process this data. I want to be able to meter each of my k8s pods saying that I want to process 1000 messages a second, how can this be done in Flux? I have already tried limit rate and concurrency in flatmap but its not helping

    Flux.fromIterable(a).limitRate(1000).publishOn(Schedulers.boundedElastic())
            .flatMap(local -> {
        System.out.println(System.currentTimeMillis() + ": time for : flatmap "+ local);

        return Mono.just(local);
    }, 1).subscribe();
1 Answers

You can not control back-pressure at the moment due to the fact that reactive support in s-c-stream via s-c-function depends on the actual implementation of source/target binders and at the moment those are non-reactive and because of it there is no way to control back-pressure. Once we have truly reactive binders, you'll get back-pressure functionality

Related