From 97b7db08656be100ac8c10d7c7968f54b02ed442 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 27 Oct 2021 17:20:06 +0200 Subject: [PATCH] GH-2212 Fix Lifecycle control over components that rely on SPCA This fixes the condition when startiing/stopping message producers which themselves are controlled by SourcsPollingChannelAdapter which we must also start/stop Resolves #2212 --- .../binder/AbstractMessageChannelBinder.java | 8 +++++++ .../cloud/stream/binder/DefaultBinding.java | 23 +++++++++++++++++++ .../stream/binder/ProducerProperties.java | 4 ++++ .../function/FunctionConfiguration.java | 11 +++++---- .../SourceToFunctionsSupportTests.java | 7 +++++- 5 files changed, 48 insertions(+), 5 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java index c6733a205..0d37a9fa5 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java @@ -255,6 +255,7 @@ public abstract class AbstractMessageChannelBinder binding = new DefaultBinding(destination, outputChannel, producerMessageHandler instanceof Lifecycle ? (Lifecycle) producerMessageHandler : null) { @@ -285,6 +286,13 @@ public abstract class AbstractMessageChannelBinder) binding).setCompanion(companion); + doPublishEvent(new BindingCreatedEvent(binding)); return binding; } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinding.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinding.java index 1b0e2805f..2d1bae052 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinding.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinding.java @@ -61,6 +61,8 @@ public class DefaultBinding implements Binding { private boolean restartable; + private Lifecycle companion; + /** * Creates an instance that associates a given name, group and binding target with an * optional {@link Lifecycle} component, which will be stopped during unbinding. @@ -84,10 +86,12 @@ public class DefaultBinding implements Binding { this.restartable = true; } + @Override public String getName() { return this.name; } + @Override public String getBindingName() { String resolvedName = (this.target instanceof IntegrationObjectSupport) ? ((IntegrationObjectSupport) this.target).getComponentName() : getName(); @@ -111,6 +115,7 @@ public class DefaultBinding implements Binding { return state; } + @Override public boolean isRunning() { return this.lifecycle != null && this.lifecycle.isRunning(); } @@ -129,6 +134,9 @@ public class DefaultBinding implements Binding { if (!this.isRunning()) { if (this.lifecycle != null && this.restartable) { this.lifecycle.start(); + if (this.companion != null) { + this.companion.start(); + } } else { this.logger.warn("Can not re-bind an anonymous binding"); @@ -140,6 +148,9 @@ public class DefaultBinding implements Binding { public synchronized void stop() { if (this.isRunning()) { this.lifecycle.stop(); + if (this.companion != null) { + this.companion.stop(); + } } } @@ -199,4 +210,16 @@ public class DefaultBinding implements Binding { return isRunning() ? "running" : "stopped"; } + /** + * Sets the companion Lifecycle. + * In most cases, when dealing with message producer (e.g., Supplier), performing + * any lifecycle operation on it does nothing as we may need to also perform the same operation on its companion + * object (e.g., SourcePollingChannelAdapter) + * + * @param companion instance of companion {@link Lifecycle} object + */ + public void setCompanion(Lifecycle companion) { + this.companion = companion; + } + } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/ProducerProperties.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/ProducerProperties.java index ec78451c3..ac0a0b8e9 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/ProducerProperties.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/ProducerProperties.java @@ -40,6 +40,10 @@ import org.springframework.expression.Expression; @JsonInclude(Include.NON_DEFAULT) public class ProducerProperties { + public ProducerProperties() { + System.out.println(); + } + /** * Signals if this producer needs to be started automatically. Default: true */ diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index f7fceaf04..46529f3a0 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -229,7 +229,7 @@ public class FunctionConfiguration { if (functionWrapper != null) { Type functionType = functionWrapper.getFunctionType(); IntegrationFlow integrationFlow = integrationFlowFromProvidedSupplier(new PartitionAwareFunctionWrapper(functionWrapper, context, producerProperties), - beginPublishingTrigger, pollable, context, taskScheduler, functionType) + beginPublishingTrigger, pollable, context, taskScheduler, functionType, producerProperties, outputName) .route(Message.class, message -> { if (message.getHeaders().get("spring.cloud.stream.sendto.destination") != null) { String destinationName = (String) message.getHeaders().get("spring.cloud.stream.sendto.destination"); @@ -244,7 +244,7 @@ public class FunctionConfiguration { else { Type functionType = ((FunctionInvocationWrapper) supplier).getFunctionType(); IntegrationFlow integrationFlow = integrationFlowFromProvidedSupplier(new PartitionAwareFunctionWrapper(supplier, context, producerProperties), - beginPublishingTrigger, pollable, context, taskScheduler, functionType) + beginPublishingTrigger, pollable, context, taskScheduler, functionType, producerProperties, outputName) .channel(c -> c.direct()) .fluxTransform((Function>, ? extends Publisher>) function) .route(Message.class, message -> { @@ -288,7 +288,7 @@ public class FunctionConfiguration { @SuppressWarnings({ "rawtypes", "unchecked" }) private IntegrationFlowBuilder integrationFlowFromProvidedSupplier(Supplier supplier, Publisher beginPublishingTrigger, PollableBean pollable, GenericApplicationContext context, - TaskScheduler taskScheduler, Type functionType) { + TaskScheduler taskScheduler, Type functionType, ProducerProperties producerProperties, String bindingName) { IntegrationFlowBuilder integrationFlowBuilder; @@ -308,7 +308,10 @@ public class FunctionConfiguration { taskScheduler.schedule(() -> { }, Instant.now()); // will keep AC alive } else { // implies pollable - integrationFlowBuilder = IntegrationFlows.fromSupplier(supplier); + + boolean autoStartup = producerProperties != null ? producerProperties.isAutoStartup() : true; + integrationFlowBuilder = IntegrationFlows + .fromSupplier(supplier, spca -> spca.id(bindingName + "_spca").autoStartup(autoStartup)); //only apply the PollableBean attributes if this is a reactive function. if (splittable && reactive) { integrationFlowBuilder = integrationFlowBuilder.split(); diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java index 41fede013..25f2e1867 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java @@ -36,6 +36,7 @@ import org.springframework.cloud.function.context.PollableBean; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.test.OutputDestination; import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; +import org.springframework.cloud.stream.binding.BindingsLifecycleController; import org.springframework.cloud.stream.messaging.Source; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; @@ -118,10 +119,14 @@ public class SourceToFunctionsSupportTests { FunctionsConfiguration.class, SupplierConfiguration.class)).web(WebApplicationType.NONE).run( "--spring.cloud.stream.function.definition=number", "--spring.jmx.enabled=false")) { - + BindingsLifecycleController lifecycle = context .getBean(BindingsLifecycleController.class); OutputDestination target = context.getBean(OutputDestination.class); String result = new String(target.receive(1000).getPayload(), StandardCharsets.UTF_8); assertThat(result).isEqualTo("1"); + result = new String(target.receive(1000).getPayload(), StandardCharsets.UTF_8); + lifecycle.stop("number-out-0"); + assertThat(result).isEqualTo("2"); + assertThat(target.receive(1000)).isNull(); } }