diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublisherIntegrationFlow.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublisherIntegrationFlow.java index f37a419629..41461e3dae 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublisherIntegrationFlow.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/PublisherIntegrationFlow.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2022 the original author or authors. + * Copyright 2016-2024 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -22,6 +22,7 @@ import org.reactivestreams.Publisher; import org.reactivestreams.Subscriber; import reactor.core.publisher.Flux; +import org.springframework.integration.endpoint.AbstractEndpoint; import org.springframework.messaging.Message; /** @@ -48,8 +49,11 @@ class PublisherIntegrationFlow extends StandardIntegrationFlow implements Pub if (autoStartOnSubscribe) { flux = flux.doOnSubscribe((sub) -> start()); for (Object component : integrationComponents.keySet()) { - if (component instanceof EndpointSpec) { - ((EndpointSpec) component).autoStartup(false); + if (component instanceof EndpointSpec endpointSpec) { + endpointSpec.autoStartup(false); + } + else if (component instanceof AbstractEndpoint endpoint) { + endpoint.setAutoStartup(false); } } } 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 e267a1f5c0..20617e3acd 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 @@ -1,5 +1,5 @@ /* - * Copyright 2016-2023 the original author or authors. + * Copyright 2016-2024 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -34,6 +34,7 @@ import org.reactivestreams.Publisher; import reactor.core.Disposable; import reactor.core.publisher.Flux; import reactor.core.scheduler.Schedulers; +import reactor.test.StepVerifier; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; @@ -45,8 +46,11 @@ import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.config.EnableIntegration; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.MessageChannels; +import org.springframework.integration.dsl.MessageProducerSpec; import org.springframework.integration.dsl.context.IntegrationFlowContext; import org.springframework.integration.endpoint.AbstractEndpoint; +import org.springframework.integration.endpoint.MessageProducerSupport; +import org.springframework.integration.endpoint.ReactiveMessageSourceProducer; import org.springframework.integration.endpoint.ReactiveStreamsConsumer; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; @@ -240,6 +244,27 @@ public class ReactiveStreamsTests { assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); } + @Autowired + MessageProducerSupport testMessageProducer; + + @Autowired + Publisher> messageProducerFlow; + + @Test + void messageProducerIsNotStartedAutomatically() { + assertThat(this.testMessageProducer.isRunning()).isFalse(); + + Flux flux = + Flux.from(this.messageProducerFlow) + .map(Message::getPayload); + + StepVerifier.create(flux) + .expectNext("test") + .expectNext("test") + .thenCancel() + .verify(Duration.ofSeconds(10)); + } + @Configuration @EnableIntegration public static class ContextConfiguration { @@ -287,6 +312,26 @@ public class ReactiveStreamsTests { .toReactivePublisher(); } + @Bean + public Publisher> messageProducerFlow() { + TestMessageProducerSpec testMessageProducerSpec = + new TestMessageProducerSpec(new ReactiveMessageSourceProducer(() -> new GenericMessage<>("test"))) + .id("testMessageProducer"); + + return IntegrationFlow + .from(testMessageProducerSpec) + .toReactivePublisher(true); + } + + } + + private static class TestMessageProducerSpec + extends MessageProducerSpec { + + TestMessageProducerSpec(ReactiveMessageSourceProducer producer) { + super(producer); + } + } }