From 72e96a3b2ea9179084481462a4b87d167cd035f4 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 22 Oct 2019 13:22:56 -0400 Subject: [PATCH] Provide MessageConverterConfigurer conditionally Check for conditions when MessageConverterConfigurer must be instantiated even in the presence of `spring.cloud.stream.function.definition`. --- .../config/BinderFactoryAutoConfiguration.java | 16 +++++++++++++--- 1 file changed, 13 insertions(+), 3 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java index cf890e15e..7c34a1634 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java @@ -212,11 +212,21 @@ public class BinderFactoryAutoConfiguration { BindingServiceProperties bindingServiceProperties, @Qualifier(IntegrationContextUtils.ARGUMENT_RESOLVER_MESSAGE_CONVERTER_BEAN_NAME) CompositeMessageConverter compositeMessageConverter, Environment environment) { + + MessageConverterConfigurer messageConverterConfigurer = null; if (StringUtils.hasText(environment.getProperty("spring.cloud.stream.function.definition"))) { - return null; + try { + ClassUtils.forName("org.apache.kafka.streams.kstream.KStream", ClassUtils.getDefaultClassLoader()); + messageConverterConfigurer = new MessageConverterConfigurer(bindingServiceProperties, compositeMessageConverter); + } + catch (Exception e) { + // ignore + } } - return new MessageConverterConfigurer(bindingServiceProperties, - compositeMessageConverter); + else { + messageConverterConfigurer = new MessageConverterConfigurer(bindingServiceProperties, compositeMessageConverter); + } + return messageConverterConfigurer; } @Bean