Commit 1ff05ac1 authored by Bhaskar Padala's avatar Bhaskar Padala

webflux with different use cases

parent 443bf42e
HELP.md
.gradle
build/
!gradle/wrapper/gradle-wrapper.jar
!**/src/main/**
!**/src/test/**
### STS ###
.apt_generated
.classpath
.factorypath
.project
.settings
.springBeans
.sts4-cache
### IntelliJ IDEA ###
.idea
*.iws
*.iml
*.ipr
out/
### NetBeans ###
/nbproject/private/
/nbbuild/
/dist/
/nbdist/
/.nb-gradle/
### VS Code ###
.vscode/
plugins {
id 'org.springframework.boot' version '2.2.7.RELEASE'
id 'io.spring.dependency-management' version '1.0.9.RELEASE'
id 'java'
}
group = 'com.nisum'
version = '0.0.1-SNAPSHOT'
sourceCompatibility = '1.8'
repositories {
mavenCentral()
}
dependencies {
implementation('org.springframework.boot:spring-boot-starter-data-mongodb-reactive')
implementation('org.springframework.boot:spring-boot-starter-webflux')
compileOnly('org.projectlombok:lombok')
testImplementation('org.springframework.boot:spring-boot-starter-test')
testImplementation('de.flapdoodle.embed:de.flapdoodle.embed.mongo')
testImplementation('io.projectreactor:reactor-test')
}
test {
useJUnitPlatform()
}
distributionBase=GRADLE_USER_HOME
distributionPath=wrapper/dists
distributionUrl=https\://services.gradle.org/distributions/gradle-6.3-bin.zip
zipStoreBase=GRADLE_USER_HOME
zipStorePath=wrapper/dists
#!/usr/bin/env sh
#
# Copyright 2015 the original author or authors.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# https://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
#
##############################################################################
##
## Gradle start up script for UN*X
##
##############################################################################
# Attempt to set APP_HOME
# Resolve links: $0 may be a link
PRG="$0"
# Need this for relative symlinks.
while [ -h "$PRG" ] ; do
ls=`ls -ld "$PRG"`
link=`expr "$ls" : '.*-> \(.*\)$'`
if expr "$link" : '/.*' > /dev/null; then
PRG="$link"
else
PRG=`dirname "$PRG"`"/$link"
fi
done
SAVED="`pwd`"
cd "`dirname \"$PRG\"`/" >/dev/null
APP_HOME="`pwd -P`"
cd "$SAVED" >/dev/null
APP_NAME="Gradle"
APP_BASE_NAME=`basename "$0"`
# Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to pass JVM options to this script.
DEFAULT_JVM_OPTS='"-Xmx64m" "-Xms64m"'
# Use the maximum available, or set MAX_FD != -1 to use that value.
MAX_FD="maximum"
warn () {
echo "$*"
}
die () {
echo
echo "$*"
echo
exit 1
}
# OS specific support (must be 'true' or 'false').
cygwin=false
msys=false
darwin=false
nonstop=false
case "`uname`" in
CYGWIN* )
cygwin=true
;;
Darwin* )
darwin=true
;;
MINGW* )
msys=true
;;
NONSTOP* )
nonstop=true
;;
esac
CLASSPATH=$APP_HOME/gradle/wrapper/gradle-wrapper.jar
# Determine the Java command to use to start the JVM.
if [ -n "$JAVA_HOME" ] ; then
if [ -x "$JAVA_HOME/jre/sh/java" ] ; then
# IBM's JDK on AIX uses strange locations for the executables
JAVACMD="$JAVA_HOME/jre/sh/java"
else
JAVACMD="$JAVA_HOME/bin/java"
fi
if [ ! -x "$JAVACMD" ] ; then
die "ERROR: JAVA_HOME is set to an invalid directory: $JAVA_HOME
Please set the JAVA_HOME variable in your environment to match the
location of your Java installation."
fi
else
JAVACMD="java"
which java >/dev/null 2>&1 || die "ERROR: JAVA_HOME is not set and no 'java' command could be found in your PATH.
Please set the JAVA_HOME variable in your environment to match the
location of your Java installation."
fi
# Increase the maximum file descriptors if we can.
if [ "$cygwin" = "false" -a "$darwin" = "false" -a "$nonstop" = "false" ] ; then
MAX_FD_LIMIT=`ulimit -H -n`
if [ $? -eq 0 ] ; then
if [ "$MAX_FD" = "maximum" -o "$MAX_FD" = "max" ] ; then
MAX_FD="$MAX_FD_LIMIT"
fi
ulimit -n $MAX_FD
if [ $? -ne 0 ] ; then
warn "Could not set maximum file descriptor limit: $MAX_FD"
fi
else
warn "Could not query maximum file descriptor limit: $MAX_FD_LIMIT"
fi
fi
# For Darwin, add options to specify how the application appears in the dock
if $darwin; then
GRADLE_OPTS="$GRADLE_OPTS \"-Xdock:name=$APP_NAME\" \"-Xdock:icon=$APP_HOME/media/gradle.icns\""
fi
# For Cygwin or MSYS, switch paths to Windows format before running java
if [ "$cygwin" = "true" -o "$msys" = "true" ] ; then
APP_HOME=`cygpath --path --mixed "$APP_HOME"`
CLASSPATH=`cygpath --path --mixed "$CLASSPATH"`
JAVACMD=`cygpath --unix "$JAVACMD"`
# We build the pattern for arguments to be converted via cygpath
ROOTDIRSRAW=`find -L / -maxdepth 1 -mindepth 1 -type d 2>/dev/null`
SEP=""
for dir in $ROOTDIRSRAW ; do
ROOTDIRS="$ROOTDIRS$SEP$dir"
SEP="|"
done
OURCYGPATTERN="(^($ROOTDIRS))"
# Add a user-defined pattern to the cygpath arguments
if [ "$GRADLE_CYGPATTERN" != "" ] ; then
OURCYGPATTERN="$OURCYGPATTERN|($GRADLE_CYGPATTERN)"
fi
# Now convert the arguments - kludge to limit ourselves to /bin/sh
i=0
for arg in "$@" ; do
CHECK=`echo "$arg"|egrep -c "$OURCYGPATTERN" -`
CHECK2=`echo "$arg"|egrep -c "^-"` ### Determine if an option
if [ $CHECK -ne 0 ] && [ $CHECK2 -eq 0 ] ; then ### Added a condition
eval `echo args$i`=`cygpath --path --ignore --mixed "$arg"`
else
eval `echo args$i`="\"$arg\""
fi
i=`expr $i + 1`
done
case $i in
0) set -- ;;
1) set -- "$args0" ;;
2) set -- "$args0" "$args1" ;;
3) set -- "$args0" "$args1" "$args2" ;;
4) set -- "$args0" "$args1" "$args2" "$args3" ;;
5) set -- "$args0" "$args1" "$args2" "$args3" "$args4" ;;
6) set -- "$args0" "$args1" "$args2" "$args3" "$args4" "$args5" ;;
7) set -- "$args0" "$args1" "$args2" "$args3" "$args4" "$args5" "$args6" ;;
8) set -- "$args0" "$args1" "$args2" "$args3" "$args4" "$args5" "$args6" "$args7" ;;
9) set -- "$args0" "$args1" "$args2" "$args3" "$args4" "$args5" "$args6" "$args7" "$args8" ;;
esac
fi
# Escape application args
save () {
for i do printf %s\\n "$i" | sed "s/'/'\\\\''/g;1s/^/'/;\$s/\$/' \\\\/" ; done
echo " "
}
APP_ARGS=`save "$@"`
# Collect all arguments for the java command, following the shell quoting and substitution rules
eval set -- $DEFAULT_JVM_OPTS $JAVA_OPTS $GRADLE_OPTS "\"-Dorg.gradle.appname=$APP_BASE_NAME\"" -classpath "\"$CLASSPATH\"" org.gradle.wrapper.GradleWrapperMain "$APP_ARGS"
exec "$JAVACMD" "$@"
@rem
@rem Copyright 2015 the original author or authors.
@rem
@rem Licensed under the Apache License, Version 2.0 (the "License");
@rem you may not use this file except in compliance with the License.
@rem You may obtain a copy of the License at
@rem
@rem https://www.apache.org/licenses/LICENSE-2.0
@rem
@rem Unless required by applicable law or agreed to in writing, software
@rem distributed under the License is distributed on an "AS IS" BASIS,
@rem WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
@rem See the License for the specific language governing permissions and
@rem limitations under the License.
@rem
@if "%DEBUG%" == "" @echo off
@rem ##########################################################################
@rem
@rem Gradle startup script for Windows
@rem
@rem ##########################################################################
@rem Set local scope for the variables with windows NT shell
if "%OS%"=="Windows_NT" setlocal
set DIRNAME=%~dp0
if "%DIRNAME%" == "" set DIRNAME=.
set APP_BASE_NAME=%~n0
set APP_HOME=%DIRNAME%
@rem Resolve any "." and ".." in APP_HOME to make it shorter.
for %%i in ("%APP_HOME%") do set APP_HOME=%%~fi
@rem Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to pass JVM options to this script.
set DEFAULT_JVM_OPTS="-Xmx64m" "-Xms64m"
@rem Find java.exe
if defined JAVA_HOME goto findJavaFromJavaHome
set JAVA_EXE=java.exe
%JAVA_EXE% -version >NUL 2>&1
if "%ERRORLEVEL%" == "0" goto init
echo.
echo ERROR: JAVA_HOME is not set and no 'java' command could be found in your PATH.
echo.
echo Please set the JAVA_HOME variable in your environment to match the
echo location of your Java installation.
goto fail
:findJavaFromJavaHome
set JAVA_HOME=%JAVA_HOME:"=%
set JAVA_EXE=%JAVA_HOME%/bin/java.exe
if exist "%JAVA_EXE%" goto init
echo.
echo ERROR: JAVA_HOME is set to an invalid directory: %JAVA_HOME%
echo.
echo Please set the JAVA_HOME variable in your environment to match the
echo location of your Java installation.
goto fail
:init
@rem Get command-line arguments, handling Windows variants
if not "%OS%" == "Windows_NT" goto win9xME_args
:win9xME_args
@rem Slurp the command line arguments.
set CMD_LINE_ARGS=
set _SKIP=2
:win9xME_args_slurp
if "x%~1" == "x" goto execute
set CMD_LINE_ARGS=%*
:execute
@rem Setup the command line
set CLASSPATH=%APP_HOME%\gradle\wrapper\gradle-wrapper.jar
@rem Execute Gradle
"%JAVA_EXE%" %DEFAULT_JVM_OPTS% %JAVA_OPTS% %GRADLE_OPTS% "-Dorg.gradle.appname=%APP_BASE_NAME%" -classpath "%CLASSPATH%" org.gradle.wrapper.GradleWrapperMain %CMD_LINE_ARGS%
:end
@rem End local scope for the variables with windows NT shell
if "%ERRORLEVEL%"=="0" goto mainEnd
:fail
rem Set variable GRADLE_EXIT_CONSOLE if you need the _script_ return code instead of
rem the _cmd.exe /c_ return code!
if not "" == "%GRADLE_EXIT_CONSOLE%" exit 1
exit /b 1
:mainEnd
if "%OS%"=="Windows_NT" endlocal
:omega
rootProject.name = 'webflux-usecases'
package com.nisum;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class WebfluxUsecasesApplication {
public static void main(String[] args) {
SpringApplication.run(WebfluxUsecasesApplication.class, args);
}
}
package com.nisum;
import org.junit.Test;
import reactor.core.publisher.ConnectableFlux;
import reactor.core.publisher.Flux;
import java.time.Duration;
public class ColdAndHotPublisherTest {
@Test
public void coldPublisherTest() throws InterruptedException {
Flux<String> stringFlux = Flux.just("A","B","C","D","E","F")
.delayElements(Duration.ofSeconds(1));
stringFlux.subscribe(s -> System.out.println("Subscriber 1 : " + s)); //emits the value from beginning
Thread.sleep(2000);
stringFlux.subscribe(s -> System.out.println("Subscriber 2 : " + s));//emits the value from beginning
Thread.sleep(4000);
}
@Test
public void hotPublisherTest() throws InterruptedException {
Flux<String> stringFlux = Flux.just("A","B","C","D","E","F")
.delayElements(Duration.ofSeconds(1));
ConnectableFlux<String> connectableFlux = stringFlux.publish();
connectableFlux.connect();
connectableFlux.subscribe(s -> System.out.println("Subscriber 1 : " + s));
Thread.sleep(3000);
connectableFlux.subscribe(s -> System.out.println("Subscriber 2 : " + s)); // does not emit the values from beginning
Thread.sleep(4000);
}
}
package com.nisum;
public class CustomException extends Throwable {
private String message;
@Override
public String getMessage() {
return message;
}
public void setMessage(String message) {
this.message = message;
}
public CustomException(Throwable e) {
this.message = e.getMessage();
}
}
package com.nisum;
import org.junit.Test;
import org.reactivestreams.Subscription;
import reactor.core.publisher.BaseSubscriber;
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
public class FluxAndMonoBackPressureTest {
@Test
public void backPressureTest() {
Flux<Integer> finiteFlux = Flux.range(1, 10)
.log();
StepVerifier.create(finiteFlux)
.expectSubscription()
.thenRequest(1)
.expectNext(1)
.thenRequest(1)
.expectNext(2)
.thenCancel()
.verify();
}
@Test
public void backPressure() {
Flux<Integer> finiteFlux = Flux.range(1, 10)
.log();
finiteFlux.subscribe((element) -> System.out.println("Element is : " + element)
, (e) -> System.err.println("Exception is : " + e)
, () -> System.out.println("Done")
, (subscription -> subscription.request(2)));
}
@Test
public void backPressure_cancel() {
Flux<Integer> finiteFlux = Flux.range(1, 10)
.log();
finiteFlux.subscribe((element) -> System.out.println("Element is : " + element)
, (e) -> System.err.println("Exception is : " + e)
, () -> System.out.println("Done")
, (subscription -> subscription.cancel()));
}
@Test
public void customized_backPressure() {
Flux<Integer> finiteFlux = Flux.range(1, 10)
.log();
finiteFlux.subscribe(new BaseSubscriber<Integer>() {
@Override
protected void hookOnNext(Integer value) {
request(1);
System.out.println("Value received is : " + value);
if(value == 4){
cancel();
}
}
});
}
}
package com.nisum;
import org.junit.Test;
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
import reactor.test.scheduler.VirtualTimeScheduler;
import java.time.Duration;
public class FluxAndMonoCombineTest {
@Test
public void combineUsingMerge(){
Flux<String> flux1 = Flux.just("A","B","C");
Flux<String> flux2 = Flux.just("D","E","F");
Flux<String> mergedFlux = Flux.merge(flux1,flux2);
StepVerifier.create(mergedFlux.log())
.expectSubscription()
.expectNext("A","B","C","D","E","F")
.verifyComplete();
}
@Test
public void combineUsingMerge_withDelay(){
Flux<String> flux1 = Flux.just("A","B","C").delayElements(Duration.ofSeconds(1));
Flux<String> flux2 = Flux.just("D","E","F").delayElements(Duration.ofSeconds(1));;
Flux<String> mergedFlux = Flux.merge(flux1,flux2);
StepVerifier.create(mergedFlux.log())
.expectSubscription()
.expectNextCount(6)
//.expectNext("A","B","C","D","E","F")
.verifyComplete();
}
@Test
public void combineUsingConcat(){
Flux<String> flux1 = Flux.just("A","B","C");
Flux<String> flux2 = Flux.just("D","E","F");
Flux<String> mergedFlux = Flux.concat(flux1,flux2);
StepVerifier.create(mergedFlux.log())
.expectSubscription()
.expectNext("A","B","C","D","E","F")
.verifyComplete();
}
@Test
public void combineUsingConcat_withDelay(){
VirtualTimeScheduler.getOrSet();
Flux<String> flux1 = Flux.just("A","B","C").delayElements(Duration.ofSeconds(1));
Flux<String> flux2 = Flux.just("D","E","F").delayElements(Duration.ofSeconds(1));
Flux<String> mergedFlux = Flux.concat(flux1,flux2);
StepVerifier.withVirtualTime(()->mergedFlux.log())
.expectSubscription()
.thenAwait(Duration.ofSeconds(6))
.expectNextCount(6)
.verifyComplete();
/* StepVerifier.create(mergedFlux.log())
.expectSubscription()
.expectNext("A","B","C","D","E","F")
.verifyComplete();*/
}
@Test
public void combineUsingZip(){
Flux<String> flux1 = Flux.just("A","B","C");
Flux<String> flux2 = Flux.just("D","E","F");
Flux<String> mergedFlux = Flux.zip(flux1,flux2, (t1,t2) -> {
return t1.concat(t2); // AD, BE, CF
}); //A,D : B,E : C:F
StepVerifier.create(mergedFlux.log())
.expectSubscription()
.expectNext("AD", "BE", "CF")
.verifyComplete();
}
}
package com.nisum;
import org.junit.Test;
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
import java.time.Duration;
public class FluxAndMonoErrorTest {
@Test
public void fluxErrorHandling() {
Flux<String> stringFlux = Flux.just("A", "B", "C")
.concatWith(Flux.error(new RuntimeException("Exception Occurred")))
.concatWith(Flux.just("D"))
.onErrorResume((e) -> { // this block gets executed
System.out.println("Exception is : " + e);
return Flux.just("default", "default1");
});
StepVerifier.create(stringFlux.log())
.expectSubscription()
.expectNext("A", "B", "C")
//.expectError(RuntimeException.class)
//.verify();
.expectNext("default", "default1")
.verifyComplete();
}
@Test
public void fluxErrorHandling_OnErrorReturn() {
Flux<String> stringFlux = Flux.just("A", "B", "C")
.concatWith(Flux.error(new RuntimeException("Exception Occurred")))
.concatWith(Flux.just("D"))
.onErrorReturn("default");
StepVerifier.create(stringFlux.log())
.expectSubscription()
.expectNext("A", "B", "C")
.expectNext("default")
.verifyComplete();
}
@Test
public void fluxErrorHandling_OnErrorMap() {
Flux<String> stringFlux = Flux.just("A", "B", "C")
.concatWith(Flux.error(new RuntimeException("Exception Occurred")))
.concatWith(Flux.just("D"))
.onErrorMap((e) -> new CustomException(e));
StepVerifier.create(stringFlux.log())
.expectSubscription()
.expectNext("A", "B", "C")
.expectError(CustomException.class)
.verify();
}
@Test
public void fluxErrorHandling_OnErrorMap_withRetry() {
Flux<String> stringFlux = Flux.just("A", "B", "C")
.concatWith(Flux.error(new RuntimeException("Exception Occurred")))
.concatWith(Flux.just("D"))
.onErrorMap((e) -> new CustomException(e))
.retry(2);
StepVerifier.create(stringFlux.log())
.expectSubscription()
.expectNext("A", "B", "C")
.expectNext("A", "B", "C")
.expectNext("A", "B", "C")
.expectError(CustomException.class)
.verify();
}
@Test
public void fluxErrorHandling_OnErrorMap_withRetryBackoff() {
Flux<String> stringFlux = Flux.just("A", "B", "C")
.concatWith(Flux.error(new RuntimeException("Exception Occurred")))
.concatWith(Flux.just("D"))
.onErrorMap((e) -> new CustomException(e))
.retryBackoff(2, Duration.ofSeconds(5));
StepVerifier.create(stringFlux.log())
.expectSubscription()
.expectNext("A", "B", "C")
.expectNext("A", "B", "C")
.expectNext("A", "B", "C")
.expectError(IllegalStateException.class)
.verify();
}
}
package com.nisum;
import org.junit.Test;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
import java.util.Arrays;
import java.util.List;
import java.util.function.Supplier;
public class FluxAndMonoFactoryTest {
List<String> names = Arrays.asList("adam","anna","jack","jenny");
@Test
public void fluxUsingIterable(){
Flux<String> namesFlux = Flux.fromIterable(names)
.log();
StepVerifier.create(namesFlux)
.expectNext("adam","anna","jack","jenny")
.verifyComplete();
}
@Test
public void fluxUsingArray(){
String[] names = new String[]{"adam","anna","jack","jenny"};
Flux<String> namesFlux = Flux.fromArray(names);
StepVerifier.create(namesFlux)
.expectNext("adam","anna","jack","jenny")
.verifyComplete();
}
@Test
public void fluxUsingStream(){
Flux<String> namesFlux = Flux.fromStream(names.stream());
StepVerifier.create(namesFlux)
.expectNext("adam","anna","jack","jenny")
.verifyComplete();
}
@Test
public void monoUsingJustOrEmpty(){
Mono<String> mono = Mono.justOrEmpty(null); //Mono.Empty();
StepVerifier.create(mono.log())
.verifyComplete();
}
@Test
public void monoUsingSupplier(){
Supplier<String> stringSupplier = () -> "adam";
Mono<String> stringMono = Mono.fromSupplier(stringSupplier);
System.out.println(stringSupplier.get());
StepVerifier.create(stringMono.log())
.expectNext("adam")
.verifyComplete();
}
@Test
public void fluxUsingRange(){
Flux<Integer> integerFlux = Flux.range(1,5).log();
StepVerifier.create(integerFlux)
.expectNext(1,2,3,4,5)
.verifyComplete();
}
}
package com.nisum;
import org.junit.Test;
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
import java.util.Arrays;
import java.util.List;
public class FluxAndMonoFilterTest {
List<String> names = Arrays.asList("adam","anna","jack","jenny");
@Test
public void filterTest(){
Flux<String> namesFlux = Flux.fromIterable(names) // adam, anna, jack,jenny
.filter(s->s.startsWith("a"))
.log(); //adam,anna
StepVerifier.create(namesFlux)
.expectNext("adam","anna")
.verifyComplete();
}
@Test
public void filterTestLength(){
Flux<String> namesFlux = Flux.fromIterable(names) // adam, anna, jack,jenny
.filter(s->s.length() >4)
.log(); //jenny
StepVerifier.create(namesFlux)
.expectNext("jenny")
.verifyComplete();
}
}
package com.nisum;
import org.junit.Test;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
public class FluxAndMonoTest {
@Test
public void fluxTest() {
Flux<String> stringFlux = Flux.just("Spring", "Spring Boot", "Reactive Spring")
/*.concatWith(Flux.error(new RuntimeException("Exception Occurred")))*/
.concatWith(Flux.just("After Error"))
.log();
stringFlux
.subscribe(System.out::println,
(e) -> System.err.println("Exception is " + e)
, () -> System.out.println("Completed"));
}
@Test
public void fluxTestElements_WithoutError() {
Flux<String> stringFlux = Flux.just("Spring", "Spring Boot", "Reactive Spring")
.log();
StepVerifier.create(stringFlux)
.expectNext("Spring")
.expectNext("Spring Boot")
.expectNext("Reactive Spring")
.verifyComplete();
}
@Test
public void fluxTestElements_WithError() {
Flux<String> stringFlux = Flux.just("Spring", "Spring Boot", "Reactive Spring")
.concatWith(Flux.error(new RuntimeException("Exception Occurred")))
.log();
StepVerifier.create(stringFlux)
.expectNext("Spring")
.expectNext("Spring Boot")
.expectNext("Reactive Spring")
//.expectError(RuntimeException.class)
.expectErrorMessage("Exception Occurred")
.verify();
}
@Test
public void fluxTestElements_WithError1() {
Flux<String> stringFlux = Flux.just("Spring", "Spring Boot", "Reactive Spring")
.concatWith(Flux.error(new RuntimeException("Exception Occurred")))
.log();
StepVerifier.create(stringFlux)
.expectNext("Spring","Spring Boot","Reactive Spring")
.expectErrorMessage("Exception Occurred")
.verify();
}
@Test
public void fluxTestElementsCount_WithError() {
Flux<String> stringFlux = Flux.just("Spring", "Spring Boot", "Reactive Spring")
.concatWith(Flux.error(new RuntimeException("Exception Occurred")))
.log();
StepVerifier.create(stringFlux)
.expectNextCount(3)
.expectErrorMessage("Exception Occurred")
.verify();
}
@Test
public void monoTest(){
Mono<String> stringMono = Mono.just("Spring");
StepVerifier.create(stringMono.log())
.expectNext("Spring")
.verifyComplete();
}
@Test
public void monoTest_Error(){
StepVerifier.create(Mono.error(new RuntimeException("Exception Occurred")).log())
.expectError(RuntimeException.class)
.verify();
}
}
package com.nisum;
import org.junit.Test;
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
import java.util.Arrays;
import java.util.List;
import static reactor.core.scheduler.Schedulers.parallel;
public class FluxAndMonoTransformTest {
List<String> names = Arrays.asList("adam", "anna", "jack", "jenny");
@Test
public void transformUsingMap() {
Flux<String> namesFlux = Flux.fromIterable(names)
.map(s -> s.toUpperCase()) //ADAM, ANNA, JACK, JENNY
.log();
StepVerifier.create(namesFlux)
.expectNext("ADAM", "ANNA", "JACK", "JENNY")
.verifyComplete();
}
@Test
public void transformUsingMap_Length() {
Flux<Integer> namesFlux = Flux.fromIterable(names)
.map(s -> s.length()) //ADAM, ANNA, JACK, JENNY
.log();
StepVerifier.create(namesFlux)
.expectNext(4,4,4,5)
.verifyComplete();
}
@Test
public void transformUsingMap_Length_repeat() {
Flux<Integer> namesFlux = Flux.fromIterable(names)
.map(s -> s.length()) //ADAM, ANNA, JACK, JENNY
.repeat(1)
.log();
StepVerifier.create(namesFlux)
.expectNext(4,4,4,5,4,4,4,5)
.verifyComplete();
}
@Test
public void transformUsingMap_Filter() {
Flux<String> namesFlux = Flux.fromIterable(names)
.filter(s -> s.length()>4)
.map(s -> s.toUpperCase()) // JENNY
.log();
StepVerifier.create(namesFlux)
.expectNext("JENNY")
.verifyComplete();
}
@Test
public void tranformUsingFlatMap(){
Flux<String> stringFlux = Flux.fromIterable(Arrays.asList("A","B","C","D","E","F")) // A, B, C, D, E, F
.flatMap(s -> {
return Flux.fromIterable(convertToList(s)); // A -> List[A, newValue] , B -> List[B, newValue]
})//db or external service call that returns a flux -> s -> Flux<String>
.log();
StepVerifier.create(stringFlux)
.expectNextCount(12)
.verifyComplete();
}
private List<String> convertToList(String s) {
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
e.printStackTrace();
}
return Arrays.asList(s, "newValue");
}
@Test
public void tranformUsingFlatMap_usingparallel(){
Flux<String> stringFlux = Flux.fromIterable(Arrays.asList("A","B","C","D","E","F")) // Flux<String>
.window(2) //Flux<Flux<String> -> (A,B), (C,D), (E,F)
.flatMap((s) ->
s.map(this::convertToList).subscribeOn(parallel())) // Flux<List<String>
.flatMap(s -> Flux.fromIterable(s)) //Flux<String>
.log();
StepVerifier.create(stringFlux)
.expectNextCount(12)
.verifyComplete();
}
@Test
public void tranformUsingFlatMap_parallel_maintain_order(){
Flux<String> stringFlux = Flux.fromIterable(Arrays.asList("A","B","C","D","E","F")) // Flux<String>
.window(2) //Flux<Flux<String> -> (A,B), (C,D), (E,F)
/* .concatMap((s) ->
s.map(this::convertToList).`(parallel())) */// Flux<List<String>
.flatMapSequential((s) ->
s.map(this::convertToList).subscribeOn(parallel()))
.flatMap(s -> Flux.fromIterable(s)) //Flux<String>
.log();
StepVerifier.create(stringFlux)
.expectNextCount(12)
.verifyComplete();
}
}
package com.nisum;
import org.junit.Test;
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
import java.time.Duration;
public class FluxAndMonoWithTimeTest {
@Test
public void infiniteSequence() throws InterruptedException {
Flux<Long> infiniteFlux = Flux.interval(Duration.ofMillis(100))
.log(); // starts from 0 --> ......
infiniteFlux
.subscribe((element) -> System.out.println("Value is : " + element));
Thread.sleep(3000);
}
@Test
public void infiniteSequenceTest() throws InterruptedException {
Flux<Long> finiteFlux = Flux.interval(Duration.ofMillis(100))
.take(3)
.log();
StepVerifier.create(finiteFlux)
.expectSubscription()
.expectNext(0L, 1L,2L)
.verifyComplete();
}
@Test
public void infiniteSequenceMap() throws InterruptedException {
Flux<Integer> finiteFlux = Flux.interval(Duration.ofMillis(100))
.map(l -> new Integer(l.intValue()))
.take(3)
.log();
StepVerifier.create(finiteFlux)
.expectSubscription()
.expectNext(0, 1,2)
.verifyComplete();
}
@Test
public void infiniteSequenceMap_withDelay() throws InterruptedException {
Flux<Integer> finiteFlux = Flux.interval(Duration.ofMillis(100))
.delayElements(Duration.ofSeconds(1))
.map(l -> new Integer(l.intValue()))
.take(3)
.log();
StepVerifier.create(finiteFlux)
.expectSubscription()
.expectNext(0, 1,2)
.verifyComplete();
}
}
package com.nisum;
import org.junit.Test;
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
import reactor.test.scheduler.VirtualTimeScheduler;
import java.time.Duration;
public class VirtualTimeTest {
@Test
public void testingWithoutVirtualTime(){
Flux<Long> longFlux = Flux.interval(Duration.ofSeconds(1))
.take(3);
StepVerifier.create(longFlux.log())
.expectSubscription()
.expectNext(0l,1l,2l)
.verifyComplete();
}
@Test
public void testingWithVirtualTime(){
VirtualTimeScheduler.getOrSet();
Flux<Long> longFlux = Flux.interval(Duration.ofSeconds(1))
.take(3);
StepVerifier.withVirtualTime(() -> longFlux.log())
.expectSubscription()
.thenAwait(Duration.ofSeconds(3))
.expectNext(0l,1l,2l)
.verifyComplete();
}
}
package com.nisum;
import org.junit.jupiter.api.Test;
import org.springframework.boot.test.context.SpringBootTest;
@SpringBootTest
class WebfluxUsecasesApplicationTests {
@Test
void contextLoads() {
}
}
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment