I'm using reactive programming where the client receives a flux of event streams of Server side events then those events are consumed. Functionality wise it works. I've a problem when I try to log the count of total records in the the flux stream. Below are the code snippets.
Let's create an instance connected to the server
final WebClient client = WebClient
.builder()
.baseUrl(url)
.build();
And then we start the connection, subscribing to its topic
final Flux<ServerSentEvent<SomeEvent>> eventStream = client.get()
.uri("/bus/sse?id=" + subscription)
.accept(MediaType.TEXT_EVENT_STREAM)
.exchange()
.flatMapMany(response -> response.bodyToFlux(type))
.repeat();
log(eventStream, "connectSSE");
eventStream
.doOnError(throwable -> this.onError(throwable, eventStream))
.doOnComplete(() -> this.onComplete(eventStream));
subscribe = eventStream.subscribe(someServerSideEvent -> this.onEvent(someServerSideEvent , eventStream));
The below method handles the event
private void onEvent(final ServerSentEvent<SomeEvent> content, Flux<ServerSentEvent<SomeEvent>> eventStream) {
log(eventStream, "onEvent");
//Code for handling event
}
I've a issue with the below piece of code. Actually I want to log the count of records in the stream and I was expecting it will print some numbers but it prints something like below. Need some solution without using .block(). Any help is welcome.
"Counted values in this Flux: MonoMapFuseable, caller onEvent"
private void log1(Flux<ServerSentEvent<SomeEvent>> eventStream, final String caller) {
try {
eventStream.count().map(count -> {
LOGGER.info("Counted values in this Flux: {}, caller {}", count.longValue(), caller);
return count;
});
} catch (final Exception e) {
LOGGER.info("Counted values in this Flux failed", e);
}
}