diff --git a/reactor-core/src/main/java/com/baeldung/reactor/creation/CharacterCreator.java b/reactor-core/src/main/java/com/baeldung/reactor/creation/CharacterCreator.java new file mode 100644 index 0000000000..382615b01c --- /dev/null +++ b/reactor-core/src/main/java/com/baeldung/reactor/creation/CharacterCreator.java @@ -0,0 +1,15 @@ +package com.baeldung.reactor.creation; + +import java.util.List; +import java.util.function.Consumer; + +import reactor.core.publisher.Flux; + +public class CharacterCreator { + + public Consumer> consumer; + + public Flux createCharacterSequence() { + return Flux.create(sink -> CharacterCreator.this.consumer = items -> items.forEach(sink::next)); + } +} diff --git a/reactor-core/src/main/java/com/baeldung/reactor/creation/CharacterGenerator.java b/reactor-core/src/main/java/com/baeldung/reactor/creation/CharacterGenerator.java new file mode 100644 index 0000000000..a22ea02bba --- /dev/null +++ b/reactor-core/src/main/java/com/baeldung/reactor/creation/CharacterGenerator.java @@ -0,0 +1,18 @@ +package com.baeldung.reactor.creation; + +import reactor.core.publisher.Flux; + +public class CharacterGenerator { + + public Flux generateCharacters() { + + return Flux.generate(() -> 97, (state, sink) -> { + char value = (char) state.intValue(); + sink.next(value); + if (value == 'z') { + sink.complete(); + } + return state + 1; + }); + } +} diff --git a/reactor-core/src/test/java/com/baeldung/reactor/creation/CharacterUnitTest.java b/reactor-core/src/test/java/com/baeldung/reactor/creation/CharacterUnitTest.java new file mode 100644 index 0000000000..2609f1f403 --- /dev/null +++ b/reactor-core/src/test/java/com/baeldung/reactor/creation/CharacterUnitTest.java @@ -0,0 +1,50 @@ +package com.baeldung.reactor.creation; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.ArrayList; +import java.util.List; + +import org.junit.Test; + +import reactor.core.publisher.Flux; +import reactor.test.StepVerifier; + +public class CharacterUnitTest { + + @Test + public void whenGeneratingCharacters_thenCharactersAreProduced() { + CharacterGenerator characterGenerator = new CharacterGenerator(); + Flux characterFlux = characterGenerator.generateCharacters().take(3); + + StepVerifier.create(characterFlux) + .expectNext('a', 'b', 'c') + .expectComplete() + .verify(); + } + + @Test + public void whenCreatingCharactersWithMultipleThreads_thenSequenceIsProducedAsynchronously() throws InterruptedException { + CharacterGenerator characterGenerator = new CharacterGenerator(); + List sequence1 = characterGenerator.generateCharacters().take(3).collectList().block(); + List sequence2 = characterGenerator.generateCharacters().take(2).collectList().block(); + + CharacterCreator characterCreator = new CharacterCreator(); + Thread producingThread1 = new Thread( + () -> characterCreator.consumer.accept(sequence1) + ); + Thread producingThread2 = new Thread( + () -> characterCreator.consumer.accept(sequence2) + ); + + List consolidated = new ArrayList<>(); + characterCreator.createCharacterSequence().subscribe(consolidated::add); + + producingThread1.start(); + producingThread2.start(); + producingThread1.join(); + producingThread2.join(); + + assertThat(consolidated).containsExactlyInAnyOrder('a', 'b', 'c', 'a', 'b'); + } +}