What is the difference between Flux.create and Flux.push? I am looking--ideally with an example use case--to understand when I should use one or the other.
What is the difference between Flux.create and Flux.push? I am looking--ideally with an example use case--to understand when I should use one or the other.
As the documentation says:
They are useful if the case you want to adapt some other async external API and not worry about cancellation and backpressure (which is handled automatically by this to method).
Here an example:
@Test
void testCreateToWrapMultiThreadsAsyncExternalAPI() {
SequenceCreator sequenceCreator = new SequenceCreator();
int numberOfElements = 10000;
StepVerifier.create(sequenceCreator.createNumberSequence(numberOfElements))
.expectNextCount(numberOfElements)
.verifyComplete();
}
@Slf4j
class SequenceCreator {
public Flux<Integer> createNumberSequence(Integer elementsToEmit) {
return Flux.create(sharedSink -> multiThreadSource(elementsToEmit, sharedSink));
}
void multiThreadSource(Integer elementsToEmit, FluxSink<Integer> sharedSink) {
Thread producingThread1 = new Thread(() -> emitElements(sharedSink, elementsToEmit / 2), "Thread_1");
Thread producingThread2 = new Thread(() -> emitElements(sharedSink, elementsToEmit / 2), "Thread_2");
producingThread1.start(); // Start to emit elements
producingThread2.start();
try {
producingThread1.join(); // Wait that thread finishes
producingThread2.join();
} catch (InterruptedException e) {
e.printStackTrace();
}
sharedSink.complete();
}
public void emitElements(FluxSink<Integer> sink, Integer count) {
IntStream.range(1, count + 1).boxed().forEach(n -> {
log.info("onNext {}", n);
sink.next(n);
});
}
}
Here you have a source that emits in elements in parallel. The source is composed of 2 threads and each thread emit numberOfElements / 2 elements that corresponds to 5000 elements emitted by Thread 1 and Thread 2 in this example. This source is wrapped with the create method and here we test that in total 10'000 elements are emitted. The test passed.
Now try to replace create with push. The test won't pass (if it is passing use a higher number for numberOfElements). That's because push expect that only one producing thread may invoke next, complete or error at a time since it won't manage the use on the FluxSink API in a concurrently way.
Hoping this toy example can help you to understand when to use one than another.
From the documentation at https://projectreactor.io/docs/core/release/api/reactor/core/publisher/Flux.html
create() Programmatically create a Flux with the capability of emitting multiple elements in a synchronous or asynchronous manner through the FluxSink API.
push() Programmatically create a Flux with the capability of emitting multiple elements from a single-threaded producer through the FluxSink API.
With create() you can produce items from multiple threads. Use push() only if you do not intend to use multiple threads.