diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsApplicationSupportAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsApplicationSupportAutoConfiguration.java index daeecb8e0..91387e3b2 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsApplicationSupportAutoConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsApplicationSupportAutoConfiguration.java @@ -27,10 +27,12 @@ import org.springframework.context.annotation.Configuration; /** * Application support configuration for Kafka Streams binder. * + * @deprecated Features provided on this class can be directly configured in the application itself using Kafka Streams. * @author Soby Chacko */ @Configuration @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) +@Deprecated public class KafkaStreamsApplicationSupportAutoConfiguration { @Bean diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java index c2f38aeec..12da94b08 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java @@ -31,12 +31,12 @@ import org.springframework.beans.factory.ObjectProvider; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; -import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.function.context.FunctionCatalog; import org.springframework.cloud.stream.binder.BinderConfiguration; +import org.springframework.cloud.stream.binder.kafka.streams.function.FunctionDetectorCondition; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties; import org.springframework.cloud.stream.binder.kafka.streams.serde.CompositeNonNativeSerde; @@ -48,6 +48,7 @@ import org.springframework.cloud.stream.config.BindingServiceConfiguration; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Conditional; import org.springframework.core.env.ConfigurableEnvironment; import org.springframework.core.env.Environment; import org.springframework.core.env.MapPropertySource; @@ -240,8 +241,7 @@ public class KafkaStreamsBinderSupportAutoConfiguration { } @Bean -// @ConditionalOnProperty("spring.cloud.stream.kafka.streams.function.definition") - @ConditionalOnProperty("spring.cloud.stream.function.definition") + @Conditional(FunctionDetectorCondition.class) public KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor(BindingServiceProperties bindingServiceProperties, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, KeyValueSerdeResolver keyValueSerdeResolver, diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java index 16ffa11f4..0677de1e2 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java @@ -71,6 +71,7 @@ import org.springframework.util.StringUtils; /** * @author Soby Chacko + * @since 2.2.0 */ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { @@ -88,6 +89,8 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { private ConfigurableApplicationContext applicationContext; + private Set origInputs = new TreeSet<>(); + private Set origOutputs = new TreeSet<>(); public KafkaStreamsFunctionProcessor(BindingServiceProperties bindingServiceProperties, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, @@ -105,43 +108,47 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { this.cleanupConfig = cleanupConfig; this.functionCatalog = functionCatalog; this.bindableProxyFactory = bindableProxyFactory; + this.origInputs.addAll(this.bindableProxyFactory.getInputs()); + this.origOutputs.addAll(this.bindableProxyFactory.getOutputs()); } private Map buildTypeMap(ResolvableType resolvableType) { - final Set inputs = new TreeSet<>(this.bindableProxyFactory.getInputs()); + int inputCount = 1; - Map map = new LinkedHashMap<>(); + ResolvableType resolvableTypeGeneric = resolvableType.getGeneric(1); + while (resolvableTypeGeneric != null && resolvableTypeGeneric.getRawClass() != null && (resolvableTypeGeneric.getRawClass().equals(Function.class) || + resolvableTypeGeneric.getRawClass().equals(Consumer.class))) { + inputCount++; + resolvableTypeGeneric = resolvableTypeGeneric.getGeneric(1); + } + + final Set inputs = new TreeSet<>(origInputs); + Map resolvableTypeMap = new LinkedHashMap<>(); final Iterator iterator = inputs.iterator(); - if (iterator.hasNext()) { - map.put(iterator.next(), resolvableType.getGeneric(0)); - ResolvableType generic = resolvableType.getGeneric(1); + final String next = iterator.next(); + resolvableTypeMap.put(next, resolvableType.getGeneric(0)); + origInputs.remove(next); - while (iterator.hasNext() && generic != null) { + for (int i = 1; i < inputCount; i++) { + if (iterator.hasNext()) { + ResolvableType generic = resolvableType.getGeneric(1); if (generic.getRawClass() != null && (generic.getRawClass().equals(Function.class) || generic.getRawClass().equals(Consumer.class))) { - map.put(iterator.next(), generic.getGeneric(0)); + final String next1 = iterator.next(); + resolvableTypeMap.put(next1, generic.getGeneric(0)); + origInputs.remove(next1); } - generic = generic.getGeneric(1); } } - - return map; + return resolvableTypeMap; } @SuppressWarnings("unchecked") - public void orchestrateStreamListenerSetupMethod(ResolvableType resolvableType, String functionName) { - final Set outputs = new TreeSet<>(this.bindableProxyFactory.getOutputs()); - - String[] methodAnnotatedOutboundNames = new String[outputs.size()]; - int j = 0; - for (String output : outputs) { - methodAnnotatedOutboundNames[j++] = output; - } - + public void orchestrateFunctionInvoking(ResolvableType resolvableType, String functionName) { final Map stringResolvableTypeMap = buildTypeMap(resolvableType); - Object[] adaptedInboundArguments = adaptAndRetrieveInboundArguments(stringResolvableTypeMap, "foobar"); + Object[] adaptedInboundArguments = adaptAndRetrieveInboundArguments(stringResolvableTypeMap, functionName); try { if (resolvableType.getRawClass() != null && resolvableType.getRawClass().equals(Consumer.class)) { Consumer consumer = functionCatalog.lookup(Consumer.class, functionName); @@ -168,15 +175,22 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { i++; } if (result != null) { + final Set outputs = new TreeSet<>(origOutputs); + final Iterator iterator = outputs.iterator(); + if (result.getClass().isArray()) { - Assert.isTrue(methodAnnotatedOutboundNames.length == ((Object[]) result).length, - "Result does not match with the number of declared outbounds"); - } - else { - Assert.isTrue(methodAnnotatedOutboundNames.length == 1, - "Result does not match with the number of declared outbounds"); - } - if (result.getClass().isArray()) { + + final int length = ((Object[]) result).length; + String[] methodAnnotatedOutboundNames = new String[length]; + + + for (int j = 0; j < length; j++) { + if (iterator.hasNext()) { + final String next = iterator.next(); + methodAnnotatedOutboundNames[j] = next; + this.origOutputs.remove(next); + } + } Object[] outboundKStreams = (Object[]) result; int k = 0; for (Object outboundKStream : outboundKStreams) { @@ -188,11 +202,15 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { } } else { - Object targetBean = this.applicationContext.getBean(methodAnnotatedOutboundNames[0]); + if (iterator.hasNext()) { + final String next = iterator.next(); + Object targetBean = this.applicationContext.getBean(next); + this.origOutputs.remove(next); - KStreamBoundElementFactory.KStreamWrapper - boundElement = (KStreamBoundElementFactory.KStreamWrapper) targetBean; - boundElement.wrap((KStream) result); + KStreamBoundElementFactory.KStreamWrapper + boundElement = (KStreamBoundElementFactory.KStreamWrapper) targetBean; + boundElement.wrap((KStream) result); + } } } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/FunctionDetectorCondition.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/FunctionDetectorCondition.java new file mode 100644 index 000000000..f1f16b9c5 --- /dev/null +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/FunctionDetectorCondition.java @@ -0,0 +1,89 @@ +/* + * Copyright 2019-2019 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.stream.binder.kafka.streams.function; + +import java.lang.reflect.Method; +import java.util.HashMap; +import java.util.Map; +import java.util.function.Consumer; +import java.util.function.Function; + +import org.apache.kafka.streams.kstream.GlobalKTable; +import org.apache.kafka.streams.kstream.KStream; +import org.apache.kafka.streams.kstream.KTable; + +import org.springframework.beans.factory.annotation.AnnotatedBeanDefinition; +import org.springframework.boot.autoconfigure.condition.ConditionOutcome; +import org.springframework.boot.autoconfigure.condition.SpringBootCondition; +import org.springframework.context.annotation.ConditionContext; +import org.springframework.core.ResolvableType; +import org.springframework.core.type.AnnotatedTypeMetadata; +import org.springframework.util.ClassUtils; + +/** + * Custom {@link org.springframework.context.annotation.Condition} that detects the presence + * of java.util.Function|Consumer beans. Used for Kafka Streams function support. + * + * @author Soby Chakco + * @since 2.2.0 + */ +public class FunctionDetectorCondition extends SpringBootCondition { + + @Override + public ConditionOutcome getMatchOutcome(ConditionContext context, AnnotatedTypeMetadata metadata) { + if (context != null && context.getBeanFactory() != null) { + + final Map functionTypes = context.getBeanFactory().getBeansOfType(Function.class); + final Map consumerTypes = context.getBeanFactory().getBeansOfType(Consumer.class); + + final Map prunedFunctionMap = pruneFunctionBeansForKafkaStreams(functionTypes, context); + final Map prunedConsumerMap = pruneFunctionBeansForKafkaStreams(consumerTypes, context); + + if (!prunedFunctionMap.isEmpty() || !prunedConsumerMap.isEmpty()) { + return ConditionOutcome.match("Matched. Function/Consumer beans found"); + } + else { + return ConditionOutcome.noMatch("No match. No Function/Consumer beans found"); + } + } + return ConditionOutcome.noMatch("No match. No Function/Consumer beans found"); + } + + private static Map pruneFunctionBeansForKafkaStreams(Map originalFunctionBeans, + ConditionContext context) { + final Map prunedMap = new HashMap<>(); + + for (String key : originalFunctionBeans.keySet()) { + final Class classObj = ClassUtils.resolveClassName(((AnnotatedBeanDefinition) + context.getBeanFactory().getBeanDefinition(key)) + .getMetadata().getClassName(), + ClassUtils.getDefaultClassLoader()); + try { + Method method = classObj.getMethod(key); + ResolvableType resolvableType = ResolvableType.forMethodReturnType(method, classObj); + final Class rawClass = resolvableType.getGeneric(0).getRawClass(); + if (rawClass == KStream.class || rawClass == KTable.class || rawClass == GlobalKTable.class) { + prunedMap.put(key, originalFunctionBeans.get(key)); + } + } + catch (NoSuchMethodException e) { + //ignore + } + } + return prunedMap; + } +} diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionAutoConfiguration.java index 6661b9db5..24c25cff7 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionAutoConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionAutoConfiguration.java @@ -16,34 +16,32 @@ package org.springframework.cloud.stream.binder.kafka.streams.function; -import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.streams.KafkaStreamsFunctionProcessor; import org.springframework.cloud.stream.function.StreamFunctionProperties; import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Conditional; import org.springframework.context.annotation.Configuration; /** * @author Soby Chacko + * @since 2.2.0 */ @Configuration -@ConditionalOnProperty("spring.cloud.stream.function.definition") @EnableConfigurationProperties(StreamFunctionProperties.class) public class KafkaStreamsFunctionAutoConfiguration { @Bean + @Conditional(FunctionDetectorCondition.class) public KafkaStreamsFunctionProcessorInvoker kafkaStreamsFunctionProcessorInvoker( KafkaStreamsFunctionBeanPostProcessor kafkaStreamsFunctionBeanPostProcessor, - KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor, - StreamFunctionProperties properties) { - return new KafkaStreamsFunctionProcessorInvoker(kafkaStreamsFunctionBeanPostProcessor.getResolvableType(), - properties.getDefinition(), kafkaStreamsFunctionProcessor); + KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor) { + return new KafkaStreamsFunctionProcessorInvoker(kafkaStreamsFunctionBeanPostProcessor.getResolvableTypes(), + kafkaStreamsFunctionProcessor); } @Bean - public KafkaStreamsFunctionBeanPostProcessor kafkaStreamsFunctionBeanPostProcessor( - StreamFunctionProperties properties) { - return new KafkaStreamsFunctionBeanPostProcessor(properties); + public KafkaStreamsFunctionBeanPostProcessor kafkaStreamsFunctionBeanPostProcessor() { + return new KafkaStreamsFunctionBeanPostProcessor(); } - } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java index 1b5c38f14..914983b83 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionBeanPostProcessor.java @@ -17,6 +17,10 @@ package org.springframework.cloud.stream.binder.kafka.streams.function; import java.lang.reflect.Method; +import java.util.Map; +import java.util.TreeMap; +import java.util.function.Consumer; +import java.util.function.Function; import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; @@ -24,40 +28,43 @@ import org.springframework.beans.factory.BeanFactoryAware; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.AnnotatedBeanDefinition; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; -import org.springframework.cloud.stream.function.StreamFunctionProperties; import org.springframework.core.ResolvableType; import org.springframework.util.ClassUtils; /** * * @author Soby Chacko - * @since 2.1.0 + * @since 2.2.0 * */ class KafkaStreamsFunctionBeanPostProcessor implements InitializingBean, BeanFactoryAware { - private final StreamFunctionProperties kafkaStreamsFunctionProperties; private ConfigurableListableBeanFactory beanFactory; - private ResolvableType resolvableType; + private Map resolvableTypeMap = new TreeMap<>(); - KafkaStreamsFunctionBeanPostProcessor(StreamFunctionProperties properties) { - this.kafkaStreamsFunctionProperties = properties; - } - - public ResolvableType getResolvableType() { - return this.resolvableType; + public Map getResolvableTypes() { + return this.resolvableTypeMap; } @Override - public void afterPropertiesSet() throws Exception { + public void afterPropertiesSet() { + + final Map functionTypes = this.beanFactory.getBeansOfType(Function.class); + final Map consumerTypes = this.beanFactory.getBeansOfType(Consumer.class); + + functionTypes.keySet().forEach(this::extractResolvableTypes); + consumerTypes.keySet().forEach(this::extractResolvableTypes); + } + + private void extractResolvableTypes(String key) { final Class classObj = ClassUtils.resolveClassName(((AnnotatedBeanDefinition) - this.beanFactory.getBeanDefinition(kafkaStreamsFunctionProperties.getDefinition())) + this.beanFactory.getBeanDefinition(key)) .getMetadata().getClassName(), ClassUtils.getDefaultClassLoader()); - try { - Method method = classObj.getMethod(this.kafkaStreamsFunctionProperties.getDefinition()); - this.resolvableType = ResolvableType.forMethodReturnType(method, classObj); + Method method = classObj.getMethod(key); + ResolvableType resolvableType = ResolvableType.forMethodReturnType(method, classObj); + resolvableTypeMap.put(key, resolvableType); } catch (NoSuchMethodException e) { //ignore diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionProcessorInvoker.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionProcessorInvoker.java index 2f854dea2..fa491192b 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionProcessorInvoker.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionProcessorInvoker.java @@ -16,6 +16,8 @@ package org.springframework.cloud.stream.binder.kafka.streams.function; +import java.util.Map; + import javax.annotation.PostConstruct; import org.springframework.cloud.stream.binder.kafka.streams.KafkaStreamsFunctionProcessor; @@ -29,18 +31,17 @@ import org.springframework.core.ResolvableType; class KafkaStreamsFunctionProcessorInvoker { private final KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor; - private final ResolvableType resolvableType; - private final String functionName; + private final Map resolvableTypeMap; - KafkaStreamsFunctionProcessorInvoker(ResolvableType resolvableType, String functionName, - KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor) { + KafkaStreamsFunctionProcessorInvoker(Map resolvableTypeMap, + KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor) { this.kafkaStreamsFunctionProcessor = kafkaStreamsFunctionProcessor; - this.resolvableType = resolvableType; - this.functionName = functionName; + this.resolvableTypeMap = resolvableTypeMap; } @PostConstruct void invoke() { - this.kafkaStreamsFunctionProcessor.orchestrateStreamListenerSetupMethod(resolvableType, functionName); + resolvableTypeMap.forEach((key, value) -> + this.kafkaStreamsFunctionProcessor.orchestrateFunctionInvoking(value, key)); } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionWrapperDetector.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionWrapperDetector.java index c77a1cc8d..7d2ba85d0 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionWrapperDetector.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionWrapperDetector.java @@ -26,6 +26,7 @@ import org.springframework.cloud.function.context.WrapperDetector; /** * @author Soby Chacko + * @since 2.2.0 */ public class KafkaStreamsFunctionWrapperDetector implements WrapperDetector { @Override diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsApplicationSupportProperties.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsApplicationSupportProperties.java index a0254ad97..04eeaa4b1 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsApplicationSupportProperties.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsApplicationSupportProperties.java @@ -25,9 +25,11 @@ import org.springframework.boot.context.properties.ConfigurationProperties; * stream processing and one can provide window specific properties at runtime and use * those properties in the applications using this class. * + * @deprecated The properties exposed by this class can be used directly on Kafka Streams API in the application. * @author Soby Chacko */ @ConfigurationProperties("spring.cloud.stream.kafka.streams") +@Deprecated public class KafkaStreamsApplicationSupportProperties { private TimeWindow timeWindow; diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java index fbd0c1280..0766bf25f 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java @@ -85,7 +85,6 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.jmx.enabled=false", - "--spring.cloud.stream.function.definition=process", "--spring.cloud.stream.bindings.input.destination=words", "--spring.cloud.stream.bindings.output1.destination=counts", "--spring.cloud.stream.bindings.output1.contentType=application/json", diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java index 11b6ddf90..7a8eae7f7 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java @@ -85,7 +85,6 @@ public class KafkaStreamsBinderWordCountFunctionTests { try (ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.jmx.enabled=false", - "--spring.cloud.stream.function.definition=process", "--spring.cloud.stream.bindings.input.destination=words", "--spring.cloud.stream.bindings.output.destination=counts", "--spring.cloud.stream.bindings.output.contentType=application/json", diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java index 727f5bfbe..4fa4caf94 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java @@ -72,7 +72,6 @@ public class StreamToGlobalKTableFunctionTests { app.setWebApplicationType(WebApplicationType.NONE); try (ConfigurableApplicationContext ignored = app.run("--server.port=0", "--spring.jmx.enabled=false", - "--spring.cloud.stream.function.definition=process", "--spring.cloud.stream.bindings.input.destination=orders", "--spring.cloud.stream.bindings.input-x.destination=customers", "--spring.cloud.stream.bindings.input-y.destination=products", diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java index f30cc4603..ff7af6d00 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java @@ -231,7 +231,6 @@ public class StreamToTableJoinFunctionTests { try (ConfigurableApplicationContext ignored = app.run("--server.port=0", "--spring.jmx.enabled=false", - "--spring.cloud.stream.function.definition=process1", "--spring.cloud.stream.bindings.input-1.destination=user-clicks-2", "--spring.cloud.stream.bindings.input-2.destination=user-regions-2", "--spring.cloud.stream.bindings.output.destination=output-topic-2",