From c2ae2dbba8cb8b822d42507062ad7ec242e47a95 Mon Sep 17 00:00:00 2001 From: David Turanski Date: Wed, 7 Oct 2020 10:56:41 -0400 Subject: [PATCH] Change exception to warning message and ignore @PollableBean for imperative functions --- .../cloud/stream/function/FunctionConfiguration.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) 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 115576b4d..01e2c4387 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 @@ -242,8 +242,9 @@ public class FunctionConfiguration { boolean splittable = pollable != null && (boolean) AnnotationUtils.getAnnotationAttributes(pollable).get("splittable"); + boolean reactive = FunctionTypeUtils.isReactive(FunctionTypeUtils.getInputType(functionType, 0)); - if (pollable == null && FunctionTypeUtils.isReactive(FunctionTypeUtils.getInputType(functionType, 0))) { + if (pollable == null && reactive) { Publisher publisher = (Publisher) supplier.get(); publisher = publisher instanceof Mono ? ((Mono) publisher).delaySubscription(beginPublishingTrigger).map(this::wrapToMessageIfNecessary) @@ -256,7 +257,8 @@ public class FunctionConfiguration { } else { // implies pollable integrationFlowBuilder = IntegrationFlows.fromSupplier(supplier); - if (splittable) { + //only apply the PollableBean attributes if this is a reactive function. + if (splittable && reactive) { integrationFlowBuilder = integrationFlowBuilder.split(); } } @@ -712,7 +714,7 @@ public class FunctionConfiguration { registry.registerBeanDefinition(functionDefinition + "_binding", functionBindableProxyDefinition); } else { - throw new IllegalArgumentException("The function definition '" + streamFunctionProperties.getDefinition() + + logger.warn("The function definition '" + streamFunctionProperties.getDefinition() + "' is not valid. The referenced function bean or one of its components does not exist"); } }