diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java index 7a4c9600c6..7249f10881 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/IntegrationFlowDefinition.java @@ -25,9 +25,8 @@ import java.util.concurrent.Executor; import java.util.function.Consumer; import java.util.function.Function; -import org.reactivestreams.Processor; import org.reactivestreams.Publisher; -import org.reactivestreams.Subscriber; + import org.springframework.aop.framework.Advised; import org.springframework.aop.support.AopUtils; import org.springframework.beans.factory.BeanCreationException; @@ -96,7 +95,6 @@ import org.springframework.util.Assert; import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; -import reactor.core.publisher.EmitterProcessor; import reactor.util.function.Tuple2; /** @@ -2728,10 +2726,8 @@ public abstract class IntegrationFlowDefinition processor = EmitterProcessor.create(false); - publisher = (Publisher>) processor; - Subscriber> subscriber = (Subscriber>) processor; - addComponent(new ReactiveConsumer(channelForPublisher, subscriber)); + Publisher messagePublisher = ReactiveConsumer.adaptToPublisher(channelForPublisher); + publisher = (Publisher>) messagePublisher; } else { MessageChannel reactiveChannel = new ReactiveChannel(); @@ -2742,7 +2738,7 @@ public abstract class IntegrationFlowDefinition(this.integrationComponents, publisher); + return new PublisherIntegrationFlow<>(this.integrationComponents, publisher); } private > B register(S endpointSpec, diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveConsumer.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveConsumer.java index 45b2a589c6..a32174cd5e 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveConsumer.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveConsumer.java @@ -111,7 +111,7 @@ public class ReactiveConsumer extends AbstractEndpoint { this.subscriber.cancel(); } - private Publisher> adaptToPublisher(MessageChannel inputChannel) { + public static Publisher> adaptToPublisher(MessageChannel inputChannel) { if (inputChannel instanceof SubscribableChannel) { return adaptSubscribableChannelToPublisher((SubscribableChannel) inputChannel); } @@ -124,11 +124,11 @@ public class ReactiveConsumer extends AbstractEndpoint { } } - private Publisher> adaptSubscribableChannelToPublisher(SubscribableChannel inputChannel) { + private static Publisher> adaptSubscribableChannelToPublisher(SubscribableChannel inputChannel) { return new SubscribableChannelPublisherAdapter(inputChannel); } - private Publisher> adaptPollableChannelToPublisher(PollableChannel inputChannel) { + private static Publisher> adaptPollableChannelToPublisher(PollableChannel inputChannel) { return new PollableChannelPublisherAdapter(inputChannel); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java index 4a7f37e583..eb5ac55f8b 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java @@ -52,6 +52,7 @@ import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.support.GenericMessage; import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.annotation.Repeat; import org.springframework.test.context.junit4.SpringRunner; import reactor.core.publisher.Flux; @@ -103,6 +104,7 @@ public class ReactiveStreamsTests { } @Test + @Repeat(10) public void testPollableReactiveFlow() throws Exception { this.inputChannel.send(new GenericMessage<>("1,2,3,4,5")); @@ -115,8 +117,6 @@ public class ReactiveStreamsTests { .doOnNext(p -> latch.countDown()) .subscribe(); - CountDownLatch secondSubscriberLatch = new CountDownLatch(1); - Future> future = Executors.newSingleThreadExecutor().submit(() -> Flux.just("11,12,13") @@ -129,12 +129,9 @@ public class ReactiveStreamsTests { .map(Message::getPayload) .log("org.springframework.integration.flux") .collectList() - .doOnRequest(s -> secondSubscriberLatch.countDown()) .block(Duration.ofSeconds(10)) ); - assertTrue(secondSubscriberLatch.await(10, TimeUnit.SECONDS)); - this.inputChannel.send(new GenericMessage<>("6,7,8,9,10")); assertTrue(latch.await(10, TimeUnit.SECONDS));