From b547c923d10e7e84bffc9dff340f94085084a7db Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 1 May 2018 15:50:10 -0400 Subject: [PATCH] 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 --- .../integration/channel/FluxMessageChannel.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/FluxMessageChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/FluxMessageChannel.java index 6b63d2c600..84a8861913 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/FluxMessageChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/FluxMessageChannel.java @@ -45,7 +45,7 @@ public class FluxMessageChannel extends AbstractMessageChannel private final List>> subscribers = new ArrayList<>(); - private final Map>, ConnectableFlux>> publishers = new ConcurrentHashMap<>(); + private final Map>, ConnectableFlux> publishers = new ConcurrentHashMap<>(); private final Flux> flux; @@ -79,10 +79,11 @@ public class FluxMessageChannel extends AbstractMessageChannel @Override public void subscribeTo(Publisher> publisher) { - ConnectableFlux> 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);