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 c5222fd08..a9aa4a699 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 @@ -73,6 +73,7 @@ public class FunctionConfiguration { return new IntegrationFlowFunctionSupport(functionCatalog, functionInspector, messageConverterFactory, functionProperties, bindingServiceProperties, context); + } /** diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionInvoker.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionInvoker.java index 0fe7b9773..a18fa47a8 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionInvoker.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionInvoker.java @@ -27,11 +27,13 @@ import org.apache.commons.logging.LogFactory; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; +import org.springframework.beans.factory.BeanFactory; import org.springframework.cloud.function.context.FunctionCatalog; import org.springframework.cloud.function.context.FunctionType; import org.springframework.cloud.function.context.catalog.FunctionInspector; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.ProducerProperties; +import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; import org.springframework.integration.support.MessageBuilder; @@ -72,7 +74,7 @@ class FunctionInvoker implements Function>, Flux implements Function>, Flux implements Function>, Flux implements Function>, Flux originalMessage) { - if (this.errorChannel != null) { - ErrorMessage em = new ErrorMessage(t, (Message) originalMessage); + String inputDestinationName = functionProperties.getInputDestinationName(); + BindingProperties bindingProperties = functionProperties.getBindingServiceProperties().getBindings().get(inputDestinationName); + String destinationName = bindingProperties.getDestination(); + String groupName = bindingProperties.getGroup(); + String bindingErrorChannelName = destinationName + "." + groupName + ".errors"; + + if (beanFactory != null) { + MessageChannel errorChannel = beanFactory.containsBean(bindingErrorChannelName) + ? beanFactory.getBean(bindingErrorChannelName, MessageChannel.class) + : beanFactory.getBean("errorChannel", MessageChannel.class); + ErrorMessage em = new ErrorMessage(t, originalMessage.getHeaders(), (Message) originalMessage); logger.error(em); - this.errorChannel.send(em); + errorChannel.send(em); } else { logger.error(t); diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/IntegrationFlowFunctionSupport.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/IntegrationFlowFunctionSupport.java index 339df5b0d..55d22f048 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/IntegrationFlowFunctionSupport.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/IntegrationFlowFunctionSupport.java @@ -26,7 +26,6 @@ import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.core.publisher.MonoSink; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cloud.function.context.FunctionCatalog; import org.springframework.cloud.function.context.FunctionRegistration; import org.springframework.cloud.function.context.FunctionType; @@ -61,12 +60,12 @@ public class IntegrationFlowFunctionSupport { private final StreamFunctionProperties functionProperties; + private final AtomicReference> triggerRef = new AtomicReference<>(); private final Publisher trigger; - @Autowired - private MessageChannel errorChannel; + private final GenericApplicationContext context; IntegrationFlowFunctionSupport(FunctionCatalog functionCatalog, FunctionInspector functionInspector, @@ -74,6 +73,7 @@ public class IntegrationFlowFunctionSupport { StreamFunctionProperties functionProperties, BindingServiceProperties bindingServiceProperties, GenericApplicationContext context) { + Assert.notNull(functionCatalog, "'functionCatalog' must not be null"); Assert.notNull(functionInspector, "'functionInspector' must not be null"); Assert.notNull(messageConverterFactory, @@ -83,6 +83,7 @@ public class IntegrationFlowFunctionSupport { this.functionInspector = functionInspector; this.messageConverterFactory = messageConverterFactory; this.functionProperties = functionProperties; + this.context = context; this.functionProperties.setBindingServiceProperties(bindingServiceProperties); trigger = Mono.create(emmiter -> { triggerRef.set(emmiter); @@ -252,7 +253,7 @@ public class IntegrationFlowFunctionSupport { } FunctionInvoker functionInvoker = new FunctionInvoker<>(functionProperties, this.functionCatalog, this.functionInspector, - this.messageConverterFactory, this.errorChannel); + this.messageConverterFactory, this.context.getBeanFactory()); if (outputChannel != null) { subscribeToInput(functionInvoker, publisher, outputChannel::send);