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 7d16f4645..71655bc21 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,11 +100,6 @@ 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/GlobalKTableBinderConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinderConfiguration.java index b1ebedd7a..2db0a21fe 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinderConfiguration.java @@ -22,6 +22,7 @@ import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.config.BeanFactoryPostProcessor; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; +import org.springframework.cloud.stream.annotation.BindingProvider; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; @@ -36,6 +37,7 @@ import org.springframework.context.annotation.Configuration; * @since 2.1.0 */ @Configuration +@BindingProvider public class GlobalKTableBinderConfiguration { @Bean diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinderConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinderConfiguration.java index b30321d7d..33a8759bb 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinderConfiguration.java @@ -23,6 +23,7 @@ import org.springframework.beans.factory.config.BeanFactoryPostProcessor; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; +import org.springframework.cloud.stream.annotation.BindingProvider; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties; @@ -41,6 +42,7 @@ import org.springframework.context.annotation.Import; @Configuration @Import({ KafkaAutoConfiguration.class, KafkaStreamsBinderHealthIndicatorConfiguration.class }) +@BindingProvider public class KStreamBinderConfiguration { @Bean diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinderConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinderConfiguration.java index 286aaab52..09423ae63 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinderConfiguration.java @@ -22,6 +22,7 @@ import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.config.BeanFactoryPostProcessor; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; +import org.springframework.cloud.stream.annotation.BindingProvider; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; @@ -36,6 +37,7 @@ import org.springframework.context.annotation.Configuration; */ @SuppressWarnings("ALL") @Configuration +@BindingProvider public class KTableBinderConfiguration { @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 cbd6630a4..d30ab98bb 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 @@ -265,7 +265,7 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro } else { for (int i = 0; i < outputs; i++) { - outputBindingNames.add(functionName + "-" + "output" + "-" + i); + outputBindingNames.add(String.format("%s_%s_%d", functionName, KafkaStreamsBindableProxyFactory.DEFAULT_OUTPUT_SUFFIX, i)); } } return outputBindingNames; 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 index 461881d42..0085a9624 100644 --- 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 @@ -68,7 +68,12 @@ import org.springframework.util.CollectionUtils; */ public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFactory implements InitializingBean, BeanFactoryAware { - private static final String DEFAULT_INPUT_SUFFIX = "input"; + /** + * Default output binding name. Output binding may occur later on in the function invoker (outside of this class), + * thus making this field part of the API. + */ + public static final String DEFAULT_OUTPUT_SUFFIX = "out"; + private static final String DEFAULT_INPUT_SUFFIX = "in"; private static Log log = LogFactory.getLog(BindableProxyFactory.class); @@ -133,7 +138,7 @@ public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFacto } else { - outputBinding = this.functionName + "-" + "output"; + outputBinding = String.format("%s_%s", this.functionName, DEFAULT_OUTPUT_SUFFIX); } Assert.isTrue(outputBinding != null, "output binding is not inferred."); KafkaStreamsBindableProxyFactory.this.outputHolders.put(outputBinding, @@ -144,7 +149,6 @@ public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFacto rootBeanDefinition1.setInstanceSupplier(() -> outputHolders.get(outputBinding1).getBoundTarget()); BeanDefinitionRegistry registry = (BeanDefinitionRegistry) beanFactory; registry.registerBeanDefinition(outputBinding1, rootBeanDefinition1); - } } @@ -170,13 +174,13 @@ public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFacto int numberOfInputs = this.type.getRawClass() != null && this.type.getRawClass().isAssignableFrom(BiFunction.class) ? 2 : getNumberOfInputs(); if (numberOfInputs == 1) { - inputs.add(this.functionName + "-" + DEFAULT_INPUT_SUFFIX); + inputs.add(String.format("%s_%s", this.functionName, DEFAULT_INPUT_SUFFIX)); return inputs; } else { int i = 0; while (i < numberOfInputs) { - inputs.add(this.functionName + "-" + DEFAULT_INPUT_SUFFIX + "-" + i++); + inputs.add(String.format("%s_%s_%d", this.functionName, DEFAULT_INPUT_SUFFIX, i++)); } return inputs; } 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 a453598da..7a1ffa0fc 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,9 +16,7 @@ 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; @@ -56,21 +54,18 @@ public class KafkaStreamsFunctionAutoConfiguration { @Bean @Conditional(FunctionDetectorCondition.class) - public BeanFactoryPostProcessor implicitFunctionBinderhello(KafkaStreamsFunctionBeanPostProcessor kafkaStreamsFunctionBeanPostProcessor) { - return new BeanFactoryPostProcessor() { - @Override - public void postProcessBeanFactory(ConfigurableListableBeanFactory beanFactory) throws BeansException { - BeanDefinitionRegistry registry = (BeanDefinitionRegistry) beanFactory; + public BeanFactoryPostProcessor implicitFunctionKafkaStreamsBinder(KafkaStreamsFunctionBeanPostProcessor kafkaStreamsFunctionBeanPostProcessor) { + return beanFactory -> { + 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); - } + 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 fd8fc1a74..3b097a9a7 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 @@ -24,19 +24,12 @@ 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; @@ -57,7 +50,6 @@ public class KafkaStreamsFunctionBeanPostProcessor implements InitializingBean, @Override public void afterPropertiesSet() { - String[] functionNames = this.beanFactory.getBeanNamesForType(Function.class); String[] biFunctionNames = this.beanFactory.getBeanNamesForType(BiFunction.class); String[] consumerNames = this.beanFactory.getBeanNamesForType(Consumer.class); @@ -65,18 +57,6 @@ public class KafkaStreamsFunctionBeanPostProcessor implements InitializingBean, Stream.concat( Stream.concat(Stream.of(functionNames), Stream.of(consumerNames)), Stream.of(biFunctionNames)) .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/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 033ef4491..88bb43bc8 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 @@ -58,7 +58,7 @@ public class KafkaStreamsFunctionStateStoreTests { try (ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.jmx.enabled=false", - "--spring.cloud.stream.bindings.process-input.destination=words", + "--spring.cloud.stream.bindings.process_in.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" + 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 bf7328092..4b71f8a74 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 @@ -104,15 +104,15 @@ public class StreamToTableJoinFunctionTests { private void runTest(SpringApplication app, Consumer consumer) { try (ConfigurableApplicationContext ignored = app.run("--server.port=0", "--spring.jmx.enabled=false", - "--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.bindings.process_in_0.destination=user-clicks-1", + "--spring.cloud.stream.bindings.process_in_1.destination=user-regions-1", + "--spring.cloud.stream.bindings.process_out.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.process-input-0.consumer.applicationId" + + "--spring.cloud.stream.kafka.streams.bindings.process_in_0.consumer.applicationId" + "=StreamToTableJoinFunctionTests-abc", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) { @@ -237,7 +237,7 @@ public class StreamToTableJoinFunctionTests { "--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.process-output.destination=output-topic-2", + "--spring.cloud.stream.bindings.process_out.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", @@ -318,7 +318,7 @@ 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: + //culprit is BinderFactoryAutoConfiguration line 300 with the following code: //if (StringUtils.hasText(name)) { // ((StandardEnvironment) environment).getSystemProperties() // .putIfAbsent("spring.cloud.stream.function.definition", name); diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java index 65e1efba3..7dc0306b1 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java @@ -55,6 +55,8 @@ import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaSt import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsConsumerProperties; import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.kafka.core.CleanupConfig; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; @@ -371,6 +373,12 @@ public class StreamToTableJoinIntegrationTests { .toStream(); } + //This forces the state stores to be cleaned up before running the test. + @Bean + public CleanupConfig cleanupConfig() { + return new CleanupConfig(true, false); + } + } @EnableBinding(KafkaStreamsProcessorY.class)