From 4d0b62f8399eb6ca5feeb7be22e5492e25a773a5 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 25 Jun 2019 12:16:52 -0400 Subject: [PATCH] Functional enhancements in Kafka Streams binder This PR introduces some fundamental changes in the way functional model of Kafka Streams applications are supported in the binder. For the most part, this PR is comprised of the following changes. * Remove the usage of EnableBinding in Kafka Streams based Spring Cloud Stream applications where a bean of type java.util.function.Function or java.util.function.Consumer is provides (instead of a StreamListener). For StreamListener based Kafka Streams applications, EnableBinding is still necessary and required. * Target types (KStream, KTable, GlobalKTable) will be inferred and bound by the binder through a new binder specific bindable proxy factory. * By deault input bindings are named as -input for functions with single input and -input-0...-input-n for functions with n-1 number of inputs. * By deault output bindings are named as -output for functions with single output and -output-0...-output-n for functions with n-1 number of outputs. * If applications prefer custom input binding names, the defaults can be overridden through spring.cloud.stream.function.inputBindings.. Similarly if custom outpub binding names are needed, then that can be done through spring.cloud.streamfunction.outputBindings.. * Test changes * Refactoring and polishing Resolves #688 --- .../kafka/streams/GlobalKTableBinder.java | 5 + ...msApplicationSupportAutoConfiguration.java | 3 +- ...StreamsBinderSupportAutoConfiguration.java | 8 +- .../KafkaStreamsFunctionProcessor.java | 98 ++++--- .../KafkaStreamsBindableProxyFactory.java | 241 ++++++++++++++++++ ...KafkaStreamsFunctionAutoConfiguration.java | 30 +++ ...KafkaStreamsFunctionBeanPostProcessor.java | 21 +- .../KafkaStreamsFunctionProcessorInvoker.java | 4 +- .../KafkaStreamsFunctionWrapperDetector.java | 5 + ...sBinderWordCountBranchesFunctionTests.java | 26 +- ...kaStreamsBinderWordCountFunctionTests.java | 13 +- .../KafkaStreamsFunctionStateStoreTests.java | 12 +- .../StreamToGlobalKTableFunctionTests.java | 28 +- .../StreamToTableJoinFunctionTests.java | 61 ++--- 14 files changed, 411 insertions(+), 144 deletions(-) create mode 100644 spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBindableProxyFactory.java 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(); - - } }