Reactive Programming in Java
by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
[Link] GmbH
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Agenda
• Reactive Programming in general
• Reactive Streams and JDK 9 Flow API
• RxJava 2
• Spring Reactor 3
• Demo of reactive application with Spring 5, Spring Boot 2, Netty, MongoDB
and Thymeleaf technology stack
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Reactive Streams
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Subscriber/Publisher
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Subscription
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Publisher Implementations
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
RxJava 2
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Classic vs. Reactive (Cold Publisher)
public List<Photo> getPhotos() {
List<Photo> result = … for (Item item : getPhoto()) {
while (...) { [Link]([Link]());
[Link](...); }
}
return result;
}
Observable<Photo> publisher = [Link](subscriber -> {
while (...) {
Photo photo = ...
[Link](photo); [Link](item->{
} [Link]([Link]());
[Link](); });
});
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
RxJava 2 Basics
Observable<Customer> customers = ...
Observable<Customer> adults = [Link](c -> [Link] > 18);
Observable<Address> addresses = [Link](c -> [Link]());
Observable<Order> orders = [Link](c ->
[Link]([Link]())
);
Observable<Picture> uploadedPhotos = [Link](c ->
[Link]([Link]([Link]()))
);
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Example. Print Photo Book
Fetch photo IDs from Book 0 ms
Load meta data from Database (ID, Dimentions, Source etc.) 10 x 100ms
Load images from server HDD if avaliable
90 x 50ms
Load images from source (for example Cloud Service)
10 x 500ms
Validate Image
100 x 50ms
+
User authentication & authorisation max 500ms
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Loading Meta Data. Buffer.
Observable<Integer> ids = [Link]([Link]());
Observable<List<Integer>> idBatches = [Link](10);
Observable<MetaData> metadata =
[Link](ids -> [Link](ids));
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Loading Picture from HDD. FlatMap with at most a single result
Observable<PictureFile> fromDisk = [Link](id -> loadingFromDisk(id));
public Observable<PictureFile> loadingFromDisk(Integer id) {
Observable<PictureFile> result = [Link](s -> {
try {
[Link]([Link](id));
} catch (PictureNotFound e) {
[Link]("Not found on disk " + id);
}
[Link]();
}
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Loading from Cloud. Filter
Observable<PictureFile> fromCloud =
[Link](m->[Link]).
flatMap(id -> loadingFromCloud([Link]));
metadata 1 disk-file 1 cloud-file 3
metadata 2 disk-file 2 cloud-file 5
metadata 3 disk-file 4
metadata 4
metadata 5
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Merging Cloud and HDD
Observable<PictureFile> allPictures = [Link](fromCloud);
metadata 1 disk-file 1
metadata 2 cloud-file 3
metadata 3 disk-file 2
metadata 4 disk-file 4
metadata 5 cloud-file 5
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Joining Metadata with Picture Files
Observable<Pair<MetaData, PictureFile>> pairs =
[Link](allPictures, (m)->lastMeta, (f)->lastFile, Pair::of);
metadata 1 + disk-file 1
metadata 2 + disk-file 1
metadata 3 + disk-file 1
metadata 2 + disk-file 2
metadata 3 + cloud-file 3
pairs = [Link](p->[Link] == [Link])
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Joining Metadata with Picture Files
Observable<Picture> pictures =
[Link](p->[Link] == [Link]).map(p->new Picture([Link],
[Link]));
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Validating Result
Observable<Picture> validated = [Link](p-
>validation(p));
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Adding Authentication
Observable<Boolean> auth = auth();
Observable<Pair<Boolean, Picture>> result =
[Link](auth, validated, Pair::of);
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Back in the Real World
BlockingObservableNext<Pair<Boolean, Picture>> b =
new BlockingObservableNext<>(result);
for (Pair<Boolean, Picture> item : b) {
[Link]("Finished " + [Link]);
}
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Compare Results
Observable<Integer> ids = loadingIds();
Observable<MetaData> metadata = [Link](10).flatMap(id -> loadingMetadata(id));
Observable<PictureFile> fromDisk = [Link](id -> loadingFromDisk(id));
Classic -> 20 sec
Observable<PictureFile> fromCloud = [Link](m->[Link]).flatMap(id ->
loadingFromCloud([Link]));
Observable<PictureFile> allPictures = [Link](fromCloud);
Reactor -> 1 sec
ObservableSource<MetaData> lastMeta = (l) -> [Link](1);
ObservableSource<PictureFile> lastFile = (l) -> [Link](1);
Observable<Pair<MetaData, PictureFile>> pairs = [Link](allPictures, (m)->lastMeta, (f)-
>lastFile, Pair::of);
Observable<Picture> pictures = [Link](p->[Link] == [Link]).map(p->new Picture([Link], [Link]));
Observable<Picture> validated = [Link](p->validation(p));
Observable<Boolean> auth = auth();
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor Timeline
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Big Picture
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Reactive Stream Interfaces
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Main Types
Flux implements Publisher
is capable of emitting of 0 or more items
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Main Types
Mono implements Publisher
can emit at most once item
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Publisher creation
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
Mono<String> productTitles = [Link]("Print 9x13");
Flux<String> productTitles = [Link]([Link]("Print 9x13",
"Photobook A4", "Calendar A4"));
Flux<String> productTitles = [Link](new String[]{"Print 9x13",
"Photobook A4", "Calendar A4"});
Flux<String> productTitles = [Link]([Link]("Print 9x13",
"Photobook A4", "Calendar A4“));
[Link](), [Link](), [Link]()
[Link](), [Link](), [Link](), [Link]()
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Event subscription
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
[Link]([Link]::println);
Output:
Print 9x13
Photobook A4
Calendar A4
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Logging
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
[Link]().subscribe([Link]::println);
Output:
INFO [Link].2 - | onSubscribe()
INFO [Link].2 - | request(unbounded)
INFO [Link].2 - | onNext(Print 9x13)
Print 9x13
INFO [Link].2 - | onNext(Photobook A4)
Photobook A4
INFO [Link].2 - | onNext(Calendar A4)
Calendar A4
INFO [Link].2 - | onComplete()
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Event subscription with own Subscriber
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
[Link](new Subscriber<String>() {
@Override
public void onSubscribe(Subscription s) {
[Link](Long.MAX_VALUE);
}
@Override
public void onNext(String t) {
[Link](t);
}
@Override
public void onError(Throwable t) {
}
@Override
public void onComplete() {
}
});
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Event subscription with custom Subscriber with back-pressure
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
[Link](new Subscriber<String>() {
private long count = 0;
private Subscription subscription;
public void onSubscribe(Subscription subscription) {
[Link] = subscription;
[Link](2);
}
public void onNext(String t) {
count++;
if (count>=2) {
count = 0;
[Link](2);
}
}
...
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Event subscription with custom Subscriber with back-pressure
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
[Link]().subscribe(new Subscriber<String>()
{[Link](2);..}
Output:
INFO [Link].2 - | onSubscribe()
INFO [Link].2 - | request(2)
INFO [Link].2 - | onNext(Print 9x13)
Print 9x13
INFO [Link].2 - | onNext(Photobook A4)
Photobook A4
INFO [Link].2 - | request(2)
INFO [Link].2 - | onNext(Calendar A4)
Calendar A4
INFO [Link].2 - | onComplete()
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Operations
Transforming (map, scan)
Combining (merge, startWith)
Filtering (last, skip)
Mathematical (count, average, max)
Boolean (every, some, includes)
Source: [Link]
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Elements filtering
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
Flux<String> productTitlesStartingWithP =
[Link](productTitle-> [Link]("P"));
[Link]([Link]::println);
Output:
Print 9x13
Photobook A4
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Elements counter
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
Mono<Long> productTitlesCount = [Link]();
[Link]([Link]::println);
Output:
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Check if all elements match a condition
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
Mono<Boolean> allProductTitlesLengthBiggerThan5=
[Link](productTitle-> [Link]() > 5);
[Link]([Link]::println);
Output:
true
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Element mapping
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
Flux<Integer> productTitlesLength =
[Link](productTitle-> [Link]()) ;
[Link]([Link]::println);
Output:
10
12
11
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Concatenation of 2 Publishers
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
Mono<String> anotherProductTitle=[Link]("Teddy");
Flux<String> concatProductTitles =
[Link](anotherProductTitle);
[Link]([Link]::println);
Output:
Print 9x13
Photobook A4
Calendar A4
Teddy
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Zipping of 2 Publishers
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
Flux<Double> productPrices = [Link]([Link](0.09, 29.99, 15.99));
Flux<Tuple2<String, Double>> zippedFlux = [Link](productTitles,
productPrices);
[Link]([Link]::println);
Output:
[Print 9x13,0.09]
[Photobook A4,29.99]
[Calendar A4,15.99]
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Parallel processing
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
Flux<Double> productPrices = [Link]([Link](0.09, 29.99, 15.99));
[Link](productTitles, productPrices)
.parallel() //returns ParallelFlux, uses all available CPUs or call
parallel(numberOfCPUs)
.runOn([Link]())
.sequential()
.subscribe([Link]::println);
Output:
[Print 9x13,0.09]
[Photobook A4,29.99]
[Calendar A4,15.99]
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Blocking Publisher
Blocking Flux:
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
Iterable<String> blockingConcatProductTitles =
[Link](anotherProductTitle) .toIterable();
or
Stream<String> blockingConcatProductTitles =
[Link](anotherProductTitle) .toStream();
Blocking Mono:
String blockingProductTitles = [Link]("Print 9x13") .block();
or
CompletableFuture blockingProductTitles = [Link]("Print 9x13").toFuture();
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
LMAX Disruptor
RingBuffer with multiple producers and consumers
Source: [Link]
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Test the publishers
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
Duration verificationDuration = [Link](productTitles).
expectNextMatches(productTitle -> [Link]("Print 9x13")).
expectNextMatches(productTitle -> [Link]("Photobook A4")).
expectNextMatches(productTitle -> [Link]("Calendar A4")).
expectComplete().
verify());
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor 3 Examples
Test the publishers
Flux<String> productTitles = [Link]("Print 9x13", "Photobook A4", "Calendar A4");
[Link](productTitles).
expectNextMatches(productTitle -> [Link]("Print 9x13")).
expectNextMatches(productTitle -> [Link]("Photobook A4")).
expectNextMatches(productTitle -> [Link]("Calendar A4")).
expectComplete().
verify())
Output:
Exception in thread "main" [Link]: expectation
"expectComplete" failed (expected: onComplete(); actual: onNext(Calendar
A4));
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring Reactor Compatibility to Java 9 Flow API
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
RxJava 2 vs Spring Reactor 3
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Comparing reactive and streaming implementations
Source: [Link]
Source: [Link]
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
The future of Reactive StreamAPIs
4th generation already
5th generation will heavily make use of advanced „operator
fusion“
Reactive Programming in Java
Source by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
[Link]
Spring 5 / Spring Boot 2 / Netty /
MongoDB / Thymeleaf
Demo
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring 5 / Spring Web Reactive
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Spring 5 / Spring Web Reactive
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Dependencies
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Non-blocking Netty Web Client
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
MongoDB Reactive Database Template
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Model
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Repository
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Thymeleaf Reactive View Resolver
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Controller
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
View
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Non-blocking JSON Parsing with Jackson
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Testing
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Benefits of the Reactive Applications
• Efficient resource utilization (spending less money on
servers and data centres)
• Processing higher loads with fewer threads
Source: [Link]
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Uses Cases for Reactive Applications
• External service calls
• Highly concurrent message consumers
Source: [Link]
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Pitfalls of reactive programming
• For the wrong problem, this makes things worse
• Hard to debug (no control over executing thread)
• Mistakenly blocking a single request leads to increased
latency for all requests -> blocking all requests brings a
server to its knees
Source: [Link]
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Questions?
Reactive Programming in Java by Vadym Kazulkin and Rodion Alukhanov, [Link] GmbH
Contact
Vadym Kazulkin :
Email : [Link]@[Link]
Xing : [Link]
Rodion Alukhanov :
Email : [Link]@[Link]
Xing : [Link]
Th an k You !
w w w.i p l ab s.d e