In the solution:
`
@test
public void event_processor() {
Flux eventStream = eventProcessor()
.parallel()
.runOn(Schedulers.parallel())
.filter(event -> event.metaData.length() > 0)
.doOnNext(event -> System.out.println("Mapping event: " + event.metaData))
.map(this::toJson)
.sequential()
.concatMap(n -> appendToStore(n).thenReturn(n));
StepVerifier.create(eventStream)
.expectNextCount(250)
.verifyComplete();
List<String> steps = Scannable.from(eventStream)
.parents()
.map(Object::toString)
.collect(Collectors.toList());
String last = Scannable.from(eventStream)
.steps()
.collect(Collectors.toCollection(LinkedList::new))
.getLast();
Assertions.assertEquals("concatMap", last);
Assertions.assertTrue(steps.contains("ParallelMap"), "Map operator not executed in parallel");
Assertions.assertTrue(steps.contains("ParallelPeek"), "doOnNext operator not executed in parallel");
Assertions.assertTrue(steps.contains("ParallelFilter"), "filter operator not executed in parallel");
Assertions.assertTrue(steps.contains("ParallelRunOn"), "runOn operator not used");
}
private String toJson(Event n) {
try {
return new ObjectMapper().writeValueAsString(n);
} catch (JsonProcessingException e) {
throw Exceptions.propagate(e);
}
}
`
Assert:
Assertions.assertEquals("concatMap", last);
Is checking that the last step is a concatMap but in reallity what the solution is giving is a concatMapNoPrefetch.
To solve this we can add on the eventStream:
Flux<String> eventStream = eventProcessor() .parallel() .runOn(Schedulers.parallel()) .filter(event -> event.metaData.length() > 0) .doOnNext(event -> System.out.println("Mapping event: " + event.metaData)) .map(this::toJson) .sequential() .concatMap(n -> appendToStore(n).thenReturn(n), 500);
or changing the assetion:
Assertions.assertEquals("concatMapNoPrefetch", last);
In the solution:
`
@test
public void event_processor() {
Flux eventStream = eventProcessor()
.parallel()
.runOn(Schedulers.parallel())
.filter(event -> event.metaData.length() > 0)
.doOnNext(event -> System.out.println("Mapping event: " + event.metaData))
.map(this::toJson)
.sequential()
.concatMap(n -> appendToStore(n).thenReturn(n));
`
Assert:
Assertions.assertEquals("concatMap", last);Is checking that the last step is a concatMap but in reallity what the solution is giving is a concatMapNoPrefetch.
To solve this we can add on the eventStream:
Flux<String> eventStream = eventProcessor() .parallel() .runOn(Schedulers.parallel()) .filter(event -> event.metaData.length() > 0) .doOnNext(event -> System.out.println("Mapping event: " + event.metaData)) .map(this::toJson) .sequential() .concatMap(n -> appendToStore(n).thenReturn(n), 500);or changing the assetion:
Assertions.assertEquals("concatMapNoPrefetch", last);