diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java
index d8d484f666..ecc2710c11 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java
@@ -344,6 +344,10 @@ public abstract class AbstractMessageChannel extends IntegrationObjectSupport
}
}
+ protected boolean isApplicationRunning() {
+ return this.applicationRunning;
+ }
+
private void assertApplicationRunning(Message> message) {
if (!this.applicationRunning) {
ApplicationContext applicationContext = getApplicationContext();
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 4f49fee71e..b0a8b5e24b 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
@@ -17,6 +17,8 @@
package org.springframework.integration.channel;
import java.time.Duration;
+import java.util.ArrayList;
+import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.concurrent.locks.LockSupport;
@@ -30,6 +32,7 @@ import reactor.core.publisher.Mono;
import reactor.core.publisher.Sinks;
import reactor.util.context.ContextView;
+import org.springframework.context.Lifecycle;
import org.springframework.core.log.LogMessage;
import org.springframework.integration.IntegrationMessageHeaderAccessor;
import org.springframework.integration.StaticMessageHeaderAccessor;
@@ -38,10 +41,14 @@ import org.springframework.integration.util.IntegrationReactiveUtils;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageDeliveryException;
import org.springframework.util.Assert;
+import org.springframework.util.ReflectionUtils;
/**
* The {@link AbstractMessageChannel} implementation for the
* Reactive Streams {@link Publisher} based on the Project Reactor {@link Flux}.
+ *
+ * This class implements {@link Lifecycle} to control subscriptions to publishers
+ * attached via {@link #subscribeTo(Publisher)}, when this channel is restarted.
*
* @author Artem Bilan
* @author Gary Russell
@@ -50,11 +57,13 @@ import org.springframework.util.Assert;
* @since 5.0
*/
public class FluxMessageChannel extends AbstractMessageChannel
- implements Publisher>, ReactiveStreamsSubscribableChannel {
+ implements Publisher>, ReactiveStreamsSubscribableChannel, Lifecycle {
private final Sinks.Many> sink = Sinks.many().multicast().onBackpressureBuffer(1, false);
- private final Disposable.Composite upstreamSubscriptions = Disposables.composite();
+ private final List>> sourcePublishers = new ArrayList<>();
+
+ private volatile Disposable.Composite upstreamSubscriptions = Disposables.composite();
private volatile boolean active = true;
@@ -111,6 +120,56 @@ public class FluxMessageChannel extends AbstractMessageChannel
.subscribe(subscriber);
}
+ @Override
+ public void start() {
+ this.active = true;
+ this.upstreamSubscriptions = Disposables.composite();
+ this.sourcePublishers.forEach(this::doSubscribeTo);
+ }
+
+ @Override
+ public void stop() {
+ this.active = false;
+ this.upstreamSubscriptions.dispose();
+ }
+
+ @Override
+ public boolean isRunning() {
+ return this.active;
+ }
+
+ private void disposeUpstreamSubscription(AtomicReference disposableReference) {
+ Disposable disposable = disposableReference.get();
+ if (disposable != null) {
+ this.upstreamSubscriptions.remove(disposable);
+ disposable.dispose();
+ }
+ }
+
+ @Override
+ public void subscribeTo(Publisher extends Message>> publisher) {
+ this.sourcePublishers.add(publisher);
+ doSubscribeTo(publisher);
+ }
+
+ private void doSubscribeTo(Publisher extends Message>> publisher) {
+ Flux