diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java index 71655bc21..7d16f4645 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java @@ -100,6 +100,11 @@ public class GlobalKTableBinder extends .getExtendedConsumerProperties(channelName); } + public void setKafkaStreamsExtendedBindingProperties( + KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties) { + this.kafkaStreamsExtendedBindingProperties = kafkaStreamsExtendedBindingProperties; + } + @Override public KafkaStreamsProducerProperties getExtendedProducerProperties( String channelName) { 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 91387e3b2..b81e5c5b6 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 @@ -36,7 +36,8 @@ import org.springframework.context.annotation.Configuration; public class KafkaStreamsApplicationSupportAutoConfiguration { @Bean - @ConditionalOnProperty("spring.cloud.stream.kafka.streams.timeWindow.length") + @ConditionalOnProperty("spring.cloud.strea" + + "m.kafka.streams.timeWindow.length") public TimeWindows configuredTimeWindow( KafkaStreamsApplicationSupportProperties processorProperties) { return processorProperties.getTimeWindow().getAdvanceBy() > 0 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 584f720a0..08f6a9968 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 @@ -39,16 +39,17 @@ 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.function.KafkaStreamsBindableProxyFactory; 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; -import org.springframework.cloud.stream.binding.BindableProxyFactory; import org.springframework.cloud.stream.binding.BindingService; import org.springframework.cloud.stream.binding.StreamListenerResultAdapter; import org.springframework.cloud.stream.config.BinderProperties; import org.springframework.cloud.stream.config.BindingServiceConfiguration; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; +import org.springframework.cloud.stream.function.StreamFunctionProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Conditional; import org.springframework.core.env.ConfigurableEnvironment; @@ -257,10 +258,11 @@ public class KafkaStreamsBinderSupportAutoConfiguration { KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate, ObjectProvider cleanupConfig, - FunctionCatalog functionCatalog, BindableProxyFactory bindableProxyFactory) { + FunctionCatalog functionCatalog, KafkaStreamsBindableProxyFactory bindableProxyFactory, + StreamFunctionProperties streamFunctionProperties) { return new KafkaStreamsFunctionProcessor(bindingServiceProperties, kafkaStreamsExtendedBindingProperties, keyValueSerdeResolver, kafkaStreamsBindingInformationCatalogue, kafkaStreamsMessageConversionDelegate, - cleanupConfig.getIfUnique(), functionCatalog, bindableProxyFactory); + cleanupConfig.getIfUnique(), functionCatalog, bindableProxyFactory, streamFunctionProperties); } @Bean 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 2b63b7c8f..d9b1e6e5e 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 @@ -16,10 +16,13 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.util.ArrayList; import java.util.Arrays; import java.util.HashMap; import java.util.Iterator; import java.util.LinkedHashMap; +import java.util.LinkedHashSet; +import java.util.List; import java.util.Map; import java.util.Set; import java.util.TreeSet; @@ -36,16 +39,20 @@ import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.state.StoreBuilder; +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.BeanFactoryAware; import org.springframework.beans.factory.BeanInitializationException; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; +import org.springframework.beans.factory.support.RootBeanDefinition; import org.springframework.cloud.function.context.FunctionCatalog; -import org.springframework.cloud.function.core.FluxedConsumer; -import org.springframework.cloud.function.core.FluxedFunction; +import org.springframework.cloud.stream.binder.kafka.streams.function.KafkaStreamsBindableProxyFactory; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsConsumerProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties; -import org.springframework.cloud.stream.binding.BindableProxyFactory; import org.springframework.cloud.stream.binding.StreamListenerErrorMessages; import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.BindingServiceProperties; +import org.springframework.cloud.stream.function.StreamFunctionProperties; import org.springframework.core.ResolvableType; import org.springframework.kafka.config.StreamsBuilderFactoryBean; import org.springframework.kafka.core.CleanupConfig; @@ -57,7 +64,7 @@ import org.springframework.util.StringUtils; * @author Soby Chacko * @since 2.2.0 */ -public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderProcessor { +public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderProcessor implements BeanFactoryAware { private static final Log LOG = LogFactory.getLog(KafkaStreamsFunctionProcessor.class); @@ -69,10 +76,13 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro private final KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate; private final FunctionCatalog functionCatalog; - private Set origInputs = new TreeSet<>(); - private Set origOutputs = new TreeSet<>(); + private Set origInputs = new LinkedHashSet<>(); + private Set origOutputs = new LinkedHashSet<>(); private ResolvableType outboundResolvableType; + private KafkaStreamsBindableProxyFactory kafkaStreamsBindableProxyFactory; + private BeanFactory beanFactory; + private StreamFunctionProperties streamFunctionProperties; public KafkaStreamsFunctionProcessor(BindingServiceProperties bindingServiceProperties, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, @@ -81,7 +91,8 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate, CleanupConfig cleanupConfig, FunctionCatalog functionCatalog, - BindableProxyFactory bindableProxyFactory) { + KafkaStreamsBindableProxyFactory bindableProxyFactory, + StreamFunctionProperties streamFunctionProperties) { super(bindingServiceProperties, kafkaStreamsBindingInformationCatalogue, kafkaStreamsExtendedBindingProperties, keyValueSerdeResolver, cleanupConfig); this.bindingServiceProperties = bindingServiceProperties; @@ -90,8 +101,10 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro this.kafkaStreamsBindingInformationCatalogue = kafkaStreamsBindingInformationCatalogue; this.kafkaStreamsMessageConversionDelegate = kafkaStreamsMessageConversionDelegate; this.functionCatalog = functionCatalog; + this.kafkaStreamsBindableProxyFactory = bindableProxyFactory; this.origInputs.addAll(bindableProxyFactory.getInputs()); this.origOutputs.addAll(bindableProxyFactory.getOutputs()); + this.streamFunctionProperties = streamFunctionProperties; } private Map buildTypeMap(ResolvableType resolvableType) { @@ -103,7 +116,7 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro resolvableTypeGeneric = resolvableTypeGeneric.getGeneric(1); } - final Set inputs = new TreeSet<>(origInputs); + final Set inputs = new LinkedHashSet<>(origInputs); Map resolvableTypeMap = new LinkedHashMap<>(); final Iterator iterator = inputs.iterator(); @@ -140,24 +153,13 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro Object[] adaptedInboundArguments = adaptAndRetrieveInboundArguments(stringResolvableTypeMap, functionName); try { if (resolvableType.getRawClass() != null && resolvableType.getRawClass().equals(Consumer.class)) { - FluxedConsumer fluxedConsumer = functionCatalog.lookup(FluxedConsumer.class, functionName); - Assert.isTrue(fluxedConsumer != null, + Consumer consumer = (Consumer) this.beanFactory.getBean(functionName); + Assert.isTrue(consumer != null, "No corresponding consumer beans found in the catalog"); - Object target = fluxedConsumer.getTarget(); - - Consumer consumer = Consumer.class.isAssignableFrom(target.getClass()) ? (Consumer) target : null; - - if (consumer != null) { - consumer.accept(adaptedInboundArguments[0]); - } + consumer.accept(adaptedInboundArguments[0]); } else { - Function function = functionCatalog.lookup(Function.class, functionName); - Object target = null; - if (function instanceof FluxedFunction) { - target = ((FluxedFunction) function).getTarget(); - } - function = (Function) target; + Function function = (Function) beanFactory.getBean(functionName); Assert.isTrue(function != null, "Function bean cannot be null"); Object result = function.apply(adaptedInboundArguments[0]); int i = 1; @@ -178,24 +180,30 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro final Iterator outboundDefinitionIterator = outputs.iterator(); if (result.getClass().isArray()) { + // Binding target as the output bindings were deffered in the KafkaStreamsBindableProxyFacotyr + // due to the fact that it didn't know the returned array size. At this point in the execution, + // we know exactly the number of outbound components (from the array length), so do the binding. final int length = ((Object[]) result).length; - String[] methodAnnotatedOutboundNames = new String[length]; - for (int j = 0; j < length; j++) { - if (outboundDefinitionIterator.hasNext()) { - final String next = outboundDefinitionIterator.next(); - methodAnnotatedOutboundNames[j] = next; - this.origOutputs.remove(next); - } - } + List outputBindings = getOutputBindings(functionName, length); + Iterator iterator = outputBindings.iterator(); + BeanDefinitionRegistry registry = (BeanDefinitionRegistry) beanFactory; Object[] outboundKStreams = (Object[]) result; - int k = 0; - for (Object outboundKStream : outboundKStreams) { - Object targetBean = this.applicationContext.getBean(methodAnnotatedOutboundNames[k++]); + + for (int ij = 0; ij < length; ij++) { + + String next = iterator.next(); + this.kafkaStreamsBindableProxyFactory.addOutputBinding(next, KStream.class); + RootBeanDefinition rootBeanDefinition1 = new RootBeanDefinition(); + rootBeanDefinition1.setInstanceSupplier(() -> kafkaStreamsBindableProxyFactory.getOutputHolders().get(next).getBoundTarget()); + registry.registerBeanDefinition(next, rootBeanDefinition1); + + Object targetBean = this.applicationContext.getBean(next); KStreamBoundElementFactory.KStreamWrapper boundElement = (KStreamBoundElementFactory.KStreamWrapper) targetBean; - boundElement.wrap((KStream) outboundKStream); + boundElement.wrap((KStream) outboundKStreams[ij]); + } } else { @@ -217,6 +225,22 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro } } + private List getOutputBindings(String functionName, int outputs) { + List outputBindings = this.streamFunctionProperties.getOutputBindings().get(functionName); + List outputBindingNames = new ArrayList<>(); + if (!CollectionUtils.isEmpty(outputBindings)) { + outputBindingNames.addAll(outputBindings); + return outputBindingNames; + } + else { + for (int i = 0; i < outputs; i++) { + outputBindingNames.add(functionName + "-" + "output" + "-" + i); + } + } + return outputBindingNames; + + } + @SuppressWarnings({"unchecked"}) private Object[] adaptAndRetrieveInboundArguments(Map stringResolvableTypeMap, String functionName) { @@ -333,4 +357,8 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro return getkStream(bindingProperties, stream, nativeDecoding); } + @Override + public void setBeanFactory(BeanFactory beanFactory) throws BeansException { + this.beanFactory = beanFactory; + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBindableProxyFactory.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBindableProxyFactory.java new file mode 100644 index 000000000..48516b354 --- /dev/null +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBindableProxyFactory.java @@ -0,0 +1,241 @@ +/* + * 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.util.ArrayList; +import java.util.Iterator; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.function.Consumer; +import java.util.function.Function; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +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.BeansException; +import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.BeanFactoryAware; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; +import org.springframework.beans.factory.support.RootBeanDefinition; +import org.springframework.cloud.stream.binding.AbstractBindableProxyFactory; +import org.springframework.cloud.stream.binding.BindableProxyFactory; +import org.springframework.cloud.stream.binding.BoundTargetHolder; +import org.springframework.cloud.stream.function.StreamFunctionProperties; +import org.springframework.core.ResolvableType; +import org.springframework.util.Assert; +import org.springframework.util.CollectionUtils; + +/** + * Kafka Streams specific target bindings proxy factory. See {@link AbstractBindableProxyFactory} for more details. + * + * Targets bound by this factory: + * + * {@link KStream} + * {@link KTable} + * {@link GlobalKTable} + * + * This class looks at the Function bean's return signature as {@link ResolvableType} and introspect the individual types, + * binding them on the way. + * + * All types on the {@link ResolvableType} are bound except for KStream[] array types on the outbound, which will be + * deferred for binding at a later stage. The reason for doing that is because in this class, we don't have any way to know + * the actual size in the returned array. That has to wait until the function is invoked and we get a result. + * + * @author Soby Chacko + * @since 3.0.0 + */ +public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFactory implements InitializingBean, BeanFactoryAware { + + private static Log log = LogFactory.getLog(BindableProxyFactory.class); + + @Autowired + private StreamFunctionProperties streamFunctionProperties; + + private final ResolvableType type; + + private final String functionName; + + private BeanFactory beanFactory; + + + public KafkaStreamsBindableProxyFactory(ResolvableType type, String functionName) { + super(type.getType().getClass()); + this.type = type; + this.functionName = functionName; + } + + @Override + public void afterPropertiesSet() { + Assert.notEmpty(KafkaStreamsBindableProxyFactory.this.bindingTargetFactories, + "'bindingTargetFactories' cannot be empty"); + + ResolvableType arg0 = this.type.getGeneric(0); + List inputBindings = buildInputBindings(); + Iterator iterator = inputBindings.iterator(); + String next = iterator.next(); + bindInput(arg0, next); + BeanDefinitionRegistry registry = (BeanDefinitionRegistry) beanFactory; + + RootBeanDefinition rootBeanDefinition = new RootBeanDefinition(); + rootBeanDefinition.setInstanceSupplier(() -> inputHolders.get(next).getBoundTarget()); + registry.registerBeanDefinition(next, rootBeanDefinition); + + ResolvableType arg1 = this.type.getGeneric(1); + + while (isAnotherFunctionOrConsumerFound(arg1)) { + arg0 = arg1.getGeneric(0); + String next1 = iterator.next(); + bindInput(arg0, next1); + RootBeanDefinition rootBeanDefinition1 = new RootBeanDefinition(); + rootBeanDefinition1.setInstanceSupplier(() -> inputHolders.get(next1).getBoundTarget()); + registry.registerBeanDefinition(next1, rootBeanDefinition1); + + arg1 = arg1.getGeneric(1); + } + + //Introspect output for binding. + if (arg1 != null && arg1.getRawClass() != null && (arg1.isArray() || arg1.getRawClass().isAssignableFrom(KStream.class))) { + // if the type is array, we need to do a late binding as we don't know the number of + // output bindings at this point in the flow. + if (!arg1.isArray()) { + List outputBindings = streamFunctionProperties.getOutputBindings().get(this.functionName); + String outputBinding = null; + + if (!CollectionUtils.isEmpty(outputBindings)) { + Iterator outputBindingsIter = outputBindings.iterator(); + if (outputBindingsIter.hasNext()) { + outputBinding = outputBindingsIter.next(); + } + + } + else { + outputBinding = this.functionName + "-" + "output"; + } + Assert.isTrue(outputBinding != null, "output binding is not inferred."); + KafkaStreamsBindableProxyFactory.this.outputHolders.put(outputBinding, + new BoundTargetHolder(getBindingTargetFactory(KStream.class) + .createOutput(outputBinding), true)); + String outputBinding1 = outputBinding; + RootBeanDefinition rootBeanDefinition1 = new RootBeanDefinition(); + rootBeanDefinition1.setInstanceSupplier(() -> outputHolders.get(outputBinding1).getBoundTarget()); + registry.registerBeanDefinition(outputBinding1, rootBeanDefinition1); + } + } + } + + private boolean isAnotherFunctionOrConsumerFound(ResolvableType arg1) { + return arg1 != null && !arg1.isArray() && arg1.getRawClass() != null && + (arg1.getRawClass().isAssignableFrom(Function.class) || arg1.getRawClass().isAssignableFrom(Consumer.class)); + } + + /** + * If the application provides the property spring.cloud.stream.function.inputBindings.functionName, + * that gets precedence. Otherwise, use functionName-input or functionName-input-0, functionName-input-1 and so on + * for multiple inputs. + * + * @return an ordered collection of input bindings to use + */ + private List buildInputBindings() { + List inputs = new ArrayList<>(); + List inputBindings = streamFunctionProperties.getInputBindings().get(this.functionName); + if (!CollectionUtils.isEmpty(inputBindings)) { + inputs.addAll(inputBindings); + return inputs; + } + int numberOfInputs = getNumberOfInputs(); + if (numberOfInputs == 1) { + inputs.add(this.functionName + "-" + "input"); + return inputs; + } + else { + int i = 0; + while (i < numberOfInputs) { + inputs.add(this.functionName + "-" + "input" + "-" + i++); + } + return inputs; + } + } + + private int getNumberOfInputs() { + int numberOfInputs = 1; + ResolvableType arg1 = this.type.getGeneric(1); + + while (isAnotherFunctionOrConsumerFound(arg1)) { + arg1 = arg1.getGeneric(1); + numberOfInputs++; + } + return numberOfInputs; + + } + + private void bindInput(ResolvableType arg0, String inputName) { + if (arg0.getRawClass() != null) { + if (arg0.getRawClass().isAssignableFrom(KStream.class)) { + KafkaStreamsBindableProxyFactory.this.inputHolders.put(inputName, + new BoundTargetHolder(getBindingTargetFactory(KStream.class) + .createInput(inputName), true)); + } + else if (arg0.getRawClass().isAssignableFrom(KTable.class)) { + KafkaStreamsBindableProxyFactory.this.inputHolders.put(inputName, + new BoundTargetHolder(getBindingTargetFactory(KTable.class) + .createInput(inputName), true)); + } + else { + KafkaStreamsBindableProxyFactory.this.inputHolders.put(inputName, + new BoundTargetHolder(getBindingTargetFactory(GlobalKTable.class) + .createInput(inputName), true)); + } + } + } + + @Override + public Set getInputs() { + Set ins = new LinkedHashSet<>(); + this.inputHolders.forEach((s, BoundTargetHolder) -> ins.add(s)); + return ins; + } + + @Override + public Set getOutputs() { + Set outs = new LinkedHashSet<>(); + this.outputHolders.forEach((s, BoundTargetHolder) -> outs.add(s)); + return outs; + } + + @Override + public void setBeanFactory(BeanFactory beanFactory) throws BeansException { + this.beanFactory = beanFactory; + } + + public void addOutputBinding(String output, Class clazz) { + KafkaStreamsBindableProxyFactory.this.outputHolders.put(output, + new BoundTargetHolder(getBindingTargetFactory(clazz) + .createOutput(output), true)); + } + + public Map getOutputHolders() { + return outputHolders; + } +} + 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 24c25cff7..a453598da 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,8 +16,15 @@ package org.springframework.cloud.stream.binder.kafka.streams.function; +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.config.BeanFactoryPostProcessor; +import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; +import org.springframework.beans.factory.support.RootBeanDefinition; +import org.springframework.boot.autoconfigure.AutoConfigureBefore; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.streams.KafkaStreamsFunctionProcessor; +import org.springframework.cloud.stream.config.BinderFactoryAutoConfiguration; import org.springframework.cloud.stream.function.StreamFunctionProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Conditional; @@ -29,6 +36,7 @@ import org.springframework.context.annotation.Configuration; */ @Configuration @EnableConfigurationProperties(StreamFunctionProperties.class) +@AutoConfigureBefore(BinderFactoryAutoConfiguration.class) public class KafkaStreamsFunctionAutoConfiguration { @Bean @@ -41,7 +49,29 @@ public class KafkaStreamsFunctionAutoConfiguration { } @Bean + @Conditional(FunctionDetectorCondition.class) public KafkaStreamsFunctionBeanPostProcessor kafkaStreamsFunctionBeanPostProcessor() { return new KafkaStreamsFunctionBeanPostProcessor(); } + + @Bean + @Conditional(FunctionDetectorCondition.class) + public BeanFactoryPostProcessor implicitFunctionBinderhello(KafkaStreamsFunctionBeanPostProcessor kafkaStreamsFunctionBeanPostProcessor) { + return new BeanFactoryPostProcessor() { + @Override + public void postProcessBeanFactory(ConfigurableListableBeanFactory beanFactory) throws BeansException { + BeanDefinitionRegistry registry = (BeanDefinitionRegistry) beanFactory; + + for (String s : kafkaStreamsFunctionBeanPostProcessor.getResolvableTypes().keySet()) { + RootBeanDefinition rootBeanDefinition = new RootBeanDefinition( + KafkaStreamsBindableProxyFactory.class); + rootBeanDefinition.getConstructorArgumentValues() + .addGenericArgumentValue(kafkaStreamsFunctionBeanPostProcessor.getResolvableTypes().get(s)); + rootBeanDefinition.getConstructorArgumentValues() + .addGenericArgumentValue(s); + registry.registerBeanDefinition("kafkaStreamsBindableProxyFactory", rootBeanDefinition); + } + } + }; + } } 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 d13f4c6c2..9730d47ab 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 @@ -23,12 +23,19 @@ import java.util.function.Consumer; import java.util.function.Function; import java.util.stream.Stream; +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.BeansException; import org.springframework.beans.factory.BeanFactory; 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.beans.factory.support.BeanDefinitionRegistry; +import org.springframework.beans.factory.support.RootBeanDefinition; +import org.springframework.cloud.stream.config.BindableProvider; import org.springframework.core.ResolvableType; import org.springframework.util.ClassUtils; @@ -38,7 +45,7 @@ import org.springframework.util.ClassUtils; * @since 2.2.0 * */ -class KafkaStreamsFunctionBeanPostProcessor implements InitializingBean, BeanFactoryAware { +public class KafkaStreamsFunctionBeanPostProcessor implements InitializingBean, BeanFactoryAware { private ConfigurableListableBeanFactory beanFactory; private Map resolvableTypeMap = new TreeMap<>(); @@ -54,6 +61,18 @@ class KafkaStreamsFunctionBeanPostProcessor implements InitializingBean, BeanFac String[] consumerNames = this.beanFactory.getBeanNamesForType(Consumer.class); Stream.concat(Stream.of(functionNames), Stream.of(consumerNames)).forEach(this::extractResolvableTypes); + + BindableProvider bindableProvider = + clazz -> clazz.isAssignableFrom(KStream.class) || clazz.isAssignableFrom(KTable.class) + || clazz.isAssignableFrom(GlobalKTable.class); + + RootBeanDefinition rb = new RootBeanDefinition(); + rb.setInstanceSupplier(() -> bindableProvider); + rb.setAutowireCandidate(true); + BeanDefinitionRegistry registry = (BeanDefinitionRegistry) beanFactory; + registry.registerBeanDefinition("kafkaStreamsBindableProvider", rb); + //Forcing the bean to be created. + beanFactory.getBean("kafkaStreamsBindableProvider"); } private void extractResolvableTypes(String key) { 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 86b762e59..77bb18367 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 @@ -28,12 +28,12 @@ import org.springframework.core.ResolvableType; * @author Soby Chacko * @since 2.1.0 */ -class KafkaStreamsFunctionProcessorInvoker { +public class KafkaStreamsFunctionProcessorInvoker { private final KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor; private final Map resolvableTypeMap; - KafkaStreamsFunctionProcessorInvoker(Map resolvableTypeMap, + public KafkaStreamsFunctionProcessorInvoker(Map resolvableTypeMap, KafkaStreamsFunctionProcessor kafkaStreamsFunctionProcessor) { this.kafkaStreamsFunctionProcessor = kafkaStreamsFunctionProcessor; this.resolvableTypeMap = resolvableTypeMap; 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 7d2ba85d0..aa2a90d40 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 @@ -22,12 +22,17 @@ import org.apache.kafka.streams.kstream.GlobalKTable; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.KTable; +import org.springframework.boot.autoconfigure.AutoConfigureBefore; import org.springframework.cloud.function.context.WrapperDetector; +import org.springframework.cloud.stream.config.BinderFactoryAutoConfiguration; +import org.springframework.context.annotation.Configuration; /** * @author Soby Chacko * @since 2.2.0 */ +@Configuration +@AutoConfigureBefore(BinderFactoryAutoConfiguration.class) public class KafkaStreamsFunctionWrapperDetector implements WrapperDetector { @Override public boolean isWrapper(Type type) { 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 995b7c50a..690fb7bb5 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 @@ -38,9 +38,6 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.annotation.Input; -import org.springframework.cloud.stream.annotation.Output; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; @@ -85,6 +82,8 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.jmx.enabled=false", + "--spring.cloud.stream.function.inputBindings.process=input", + "--spring.cloud.stream.function.outputBindings.process=output1,output2,output3", "--spring.cloud.stream.bindings.input.destination=words", "--spring.cloud.stream.bindings.output1.destination=counts", "--spring.cloud.stream.bindings.output2.destination=foo", @@ -94,15 +93,11 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", -// "--spring.cloud.stream.kafka.streams.bindings.output1.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", -// "--spring.cloud.stream.kafka.streams.bindings.output2.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", -// "--spring.cloud.stream.kafka.streams.bindings.output3.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", "--spring.cloud.stream.kafka.streams.timeWindow.length=5000", "--spring.cloud.stream.kafka.streams.timeWindow.advanceBy=0", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + "=KafkaStreamsBinderWordCountBranchesFunctionTests-abc", - "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString()); + "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString()); try { receiveAndValidate(context); } @@ -182,7 +177,6 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { } } - @EnableBinding(KStreamProcessorX.class) @EnableAutoConfiguration @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class WordCountProcessorApplication { @@ -207,18 +201,4 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { } } - interface KStreamProcessorX { - - @Input("input") - KStream input(); - - @Output("output1") - KStream output1(); - - @Output("output2") - KStream output2(); - - @Output("output3") - KStream output3(); - } } 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 b7095e722..e22c0aa2b 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 @@ -39,8 +39,6 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; @@ -64,7 +62,7 @@ public class KafkaStreamsBinderWordCountFunctionTests { private static Consumer consumer; @BeforeClass - public static void setUp() throws Exception { + public static void setUp() { Map consumerProps = KafkaTestUtils.consumerProps("group", "false", embeddedKafka); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); @@ -83,7 +81,10 @@ public class KafkaStreamsBinderWordCountFunctionTests { SpringApplication app = new SpringApplication(WordCountProcessorApplication.class); app.setWebApplicationType(WebApplicationType.NONE); - try (ConfigurableApplicationContext context = app.run("--server.port=0", + try (ConfigurableApplicationContext context = app.run( + "--spring.cloud.stream.function.inputBindings.process=input", + "--spring.cloud.stream.function.outputBindings.process=output", + "--server.port=0", "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=words", "--spring.cloud.stream.bindings.output.destination=counts", @@ -93,7 +94,6 @@ public class KafkaStreamsBinderWordCountFunctionTests { "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", - //"--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) { receiveAndValidate(context); } @@ -164,10 +164,9 @@ public class KafkaStreamsBinderWordCountFunctionTests { } } - @EnableBinding(KafkaStreamsProcessor.class) @EnableAutoConfiguration @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) - static class WordCountProcessorApplication { + public static class WordCountProcessorApplication { @Bean public Function, KStream> process() { diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java index 6880f6380..033ef4491 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java @@ -33,8 +33,6 @@ import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.annotation.Input; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.DefaultKafkaProducerFactory; @@ -60,7 +58,7 @@ public class KafkaStreamsFunctionStateStoreTests { try (ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.jmx.enabled=false", - "--spring.cloud.stream.bindings.input.destination=words", + "--spring.cloud.stream.bindings.process-input.destination=words", "--spring.cloud.stream.kafka.streams.default.consumer.application-id=basic-word-count-1", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + @@ -97,9 +95,8 @@ public class KafkaStreamsFunctionStateStoreTests { } } - @EnableBinding(KStreamProcessorX.class) @EnableAutoConfiguration - static class StateStoreTestApplication { + public static class StateStoreTestApplication { KeyValueStore state1; WindowStore state2; @@ -150,9 +147,4 @@ public class KafkaStreamsFunctionStateStoreTests { } } - interface KStreamProcessorX { - @Input("input") - KStream input(); - } - } 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 594e46975..dc935f3db 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 @@ -39,9 +39,6 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.annotation.Input; -import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; @@ -72,19 +69,20 @@ public class StreamToGlobalKTableFunctionTests { app.setWebApplicationType(WebApplicationType.NONE); try (ConfigurableApplicationContext ignored = app.run("--server.port=0", "--spring.jmx.enabled=false", - "--spring.cloud.stream.bindings.input.destination=orders", - "--spring.cloud.stream.bindings.input-x.destination=customers", - "--spring.cloud.stream.bindings.input-y.destination=products", - "--spring.cloud.stream.bindings.output.destination=enriched-order", + "--spring.cloud.stream.function.inputBindings.process=order,customer,product", + "--spring.cloud.stream.function.outputBindings.process=enriched-order", + "--spring.cloud.stream.bindings.order.destination=orders", + "--spring.cloud.stream.bindings.customer.destination=customers", + "--spring.cloud.stream.bindings.product.destination=products", + "--spring.cloud.stream.bindings.enriched-order.destination=enriched-order", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=10000", - "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId=" + + "--spring.cloud.stream.kafka.streams.bindings.order.consumer.applicationId=" + "StreamToGlobalKTableJoinFunctionTests-abc", - "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString())) { + "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) { Map senderPropsCustomer = KafkaTestUtils.producerProps(embeddedKafka); senderPropsCustomer.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); senderPropsCustomer.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, @@ -176,16 +174,6 @@ public class StreamToGlobalKTableFunctionTests { } } - interface CustomGlobalKTableProcessor extends KafkaStreamsProcessor { - - @Input("input-x") - GlobalKTable inputX(); - - @Input("input-y") - GlobalKTable inputY(); - } - - @EnableBinding(CustomGlobalKTableProcessor.class) @EnableAutoConfiguration @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class OrderEnricherApplication { 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 d48e83d9c..01381aa47 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 @@ -44,9 +44,6 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.annotation.Input; -import org.springframework.cloud.stream.annotation.Output; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; @@ -84,19 +81,17 @@ 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-1", - "--spring.cloud.stream.bindings.input-2.destination=user-regions-1", - "--spring.cloud.stream.bindings.output.destination=output-topic-1", + "--spring.cloud.stream.bindings.process-input-0.destination=user-clicks-1", + "--spring.cloud.stream.bindings.process-input-1.destination=user-regions-1", + "--spring.cloud.stream.bindings.process-output.destination=output-topic-1", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=10000", - "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.applicationId" + + "--spring.cloud.stream.kafka.streams.bindings.process-input-0.consumer.applicationId" + "=StreamToTableJoinFunctionTests-abc", - "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString())) { + "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) { // Input 1: Region per user (multiple records allowed per user). List> userRegions = Arrays.asList( @@ -214,11 +209,12 @@ public class StreamToTableJoinFunctionTests { template.sendDefault(keyValue.key, keyValue.value); } - try (ConfigurableApplicationContext ignored = app.run("--server.port=0", + try (ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.jmx.enabled=false", + "--spring.cloud.stream.function.inputBindings.process=input-1,input-2", "--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", + "--spring.cloud.stream.bindings.process-output.destination=output-topic-2", "--spring.cloud.stream.bindings.input-1.consumer.useNativeDecoding=true", "--spring.cloud.stream.bindings.input-2.consumer.useNativeDecoding=true", "--spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", @@ -298,8 +294,17 @@ public class StreamToTableJoinFunctionTests { assertThat(count).isEqualTo(expectedClicksPerRegion.size()); assertThat(actualClicksPerRegion).hasSameElementsAs(expectedClicksPerRegion); + //the following removal is a code smell. Check with Oleg to see why this is happening. + //culprit is BinderFactoryAutoConfiguration line 309 with the following code: + //if (StringUtils.hasText(name)) { + // ((StandardEnvironment) environment).getSystemProperties() + // .putIfAbsent("spring.cloud.stream.function.definition", name); + // } + context.getEnvironment().getSystemProperties() + .remove("spring.cloud.stream.function.definition"); } finally { + consumer.close(); } } @@ -333,13 +338,12 @@ public class StreamToTableJoinFunctionTests { } - @EnableBinding(KStreamKTableProcessor.class) @EnableAutoConfiguration @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class CountClicksPerRegionApplication { @Bean - public Function, Function, KStream>> process1() { + public Function, Function, KStream>> process() { return userClicksStream -> (userRegionsTable -> (userClicksStream .leftJoin(userRegionsTable, (clicks, region) -> new RegionWithClicks(region == null ? "UNKNOWN" : region, clicks), @@ -347,36 +351,9 @@ public class StreamToTableJoinFunctionTests { .map((user, regionWithClicks) -> new KeyValue<>(regionWithClicks.getRegion(), regionWithClicks.getClicks())) .groupByKey(Serialized.with(Serdes.String(), Serdes.Long())) - .reduce((firstClicks, secondClicks) -> firstClicks + secondClicks) + .reduce(Long::sum) .toStream())); } } - interface KStreamKTableProcessor { - - /** - * Input binding. - * - * @return {@link Input} binding for {@link KStream} type. - */ - @Input("input-1") - KStream input1(); - - /** - * Input binding. - * - * @return {@link Input} binding for {@link KStream} type. - */ - @Input("input-2") - KTable input2(); - - /** - * Output binding. - * - * @return {@link Output} binding for {@link KStream} type. - */ - @Output("output") - KStream output(); - - } }