WebFlux and Reactor 3.4.0 - Deprecated FluxProcessors - How to subscribe with a sink?

Viewed 2575

With Reactor 3.4.0 the different FluxProcessors like "DirectProcessor" is getting deprecated. I used such processor as subscribers, see example below.

Now I am wondering how I have to migrate my code to make use of the recommended Sinks.many() approach? Any ideas?

Old Code:

DirectProcessor<String> output = DirectProcessor.create();
output.subscribe(msg -> System.out.println(msg));

WebSocketClient client = new ReactorNettyWebSocketClient();
client.execute(uri, session -> 
    // send message
    session.send(Mono.just(session.textMessage(command)))
        .thenMany(session.receive()
        .map(message -> message.getPayloadAsText())
        .subscribeWith(output))
    .then()).block();

According the JavaDoc of the deprecated DirectProcessor I should make use of Sinks.many().multicast().directBestEffort(). But I am wondering how to use this in my WebSocketClient?

Migrated Code:

Many<String> sink = Sinks.many().multicast().directBestEffort();        
Flux<String> flux = sink.asFlux();
flux.subscribe(msg -> System.out.println(msg));

WebSocketClient client = new ReactorNettyWebSocketClient();
client.execute(uri, session -> 
    // send message
    session.send(Mono.just(session.textMessage(command)))
        .thenMany(session.receive()
        .map(message -> message.getPayloadAsText())
        .subscribe ...   // <-- how to do this with a Sink ??
    .then()).block();

Thank you for any suggestions in advance.

0 Answers
Related