Use Flux.handle() instead of doOnNext() in FluxMC
For better back-pressure handling do not perform active operations directly on the current `Flux` (e.g. via `doOnNext()`). It's better to postpone such a handling to back-pressure ready operators Like in this `FluxMessageChannel` case, the `.handle()` operator is wrapping the "hard" `send()` operation. Also include `.errorStrategyContinue()` do not stop during message processing
This commit is contained in:
@@ -45,7 +45,7 @@ public class FluxMessageChannel extends AbstractMessageChannel
|
||||
|
||||
private final List<Subscriber<? super Message<?>>> subscribers = new ArrayList<>();
|
||||
|
||||
private final Map<Publisher<Message<?>>, ConnectableFlux<Message<?>>> publishers = new ConcurrentHashMap<>();
|
||||
private final Map<Publisher<Message<?>>, ConnectableFlux<?>> publishers = new ConcurrentHashMap<>();
|
||||
|
||||
private final Flux<Message<?>> flux;
|
||||
|
||||
@@ -79,10 +79,11 @@ public class FluxMessageChannel extends AbstractMessageChannel
|
||||
|
||||
@Override
|
||||
public void subscribeTo(Publisher<Message<?>> publisher) {
|
||||
ConnectableFlux<Message<?>> connectableFlux =
|
||||
ConnectableFlux<?> connectableFlux =
|
||||
Flux.from(publisher)
|
||||
.handle((message, sink) -> sink.next(send(message)))
|
||||
.errorStrategyContinue()
|
||||
.doOnComplete(() -> this.publishers.remove(publisher))
|
||||
.doOnNext(this::send)
|
||||
.publish();
|
||||
|
||||
this.publishers.put(publisher, connectableFlux);
|
||||
|
||||
Reference in New Issue
Block a user