From bbb455946e27d934059c467cf0ab024dfa2ac433 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 14 Jun 2019 11:48:32 -0400 Subject: [PATCH] Refactor common code in Kafka Streams binder. Both StreamListener and the functional model require many code paths that are common. Refactoring them for better use and readability. Resolves #674 --- .../AbstractKafkaStreamsBinderProcessor.java | 279 ++++++++++++++++ .../GlobalKTableBinderConfiguration.java | 16 +- ...KStreamStreamListenerParameterAdapter.java | 3 +- .../streams/KTableBinderConfiguration.java | 16 +- ...StreamsBinderSupportAutoConfiguration.java | 6 - .../streams/KafkaStreamsBinderUtils.java | 26 ++ ...fkaStreamsBindingInformationCatalogue.java | 4 +- .../KafkaStreamsFunctionProcessor.java | 249 +++------------ ...StreamListenerSetupMethodOrchestrator.java | 297 +++--------------- .../kafka/streams/KeyValueSerdeResolver.java | 65 ++-- .../kafka/streams/QueryableStoreRegistry.java | 66 ---- .../streams/StreamsBuilderFactoryManager.java | 2 +- .../KafkaStreamsFunctionProcessorInvoker.java | 2 +- 13 files changed, 429 insertions(+), 602 deletions(-) create mode 100644 spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java delete mode 100644 spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/QueryableStoreRegistry.java diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java new file mode 100644 index 000000000..873d354ff --- /dev/null +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java @@ -0,0 +1,279 @@ +/* + * 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; + +import java.util.Map; +import java.util.Properties; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.apache.kafka.common.serialization.Serde; +import org.apache.kafka.common.serialization.Serdes; +import org.apache.kafka.common.utils.Bytes; +import org.apache.kafka.streams.StreamsBuilder; +import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.Topology; +import org.apache.kafka.streams.kstream.Consumed; +import org.apache.kafka.streams.kstream.GlobalKTable; +import org.apache.kafka.streams.kstream.KStream; +import org.apache.kafka.streams.kstream.KTable; +import org.apache.kafka.streams.kstream.Materialized; +import org.apache.kafka.streams.state.KeyValueStore; + +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.config.BeanDefinition; +import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; +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.config.BindingProperties; +import org.springframework.cloud.stream.config.BindingServiceProperties; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.core.ResolvableType; +import org.springframework.kafka.config.KafkaStreamsConfiguration; +import org.springframework.kafka.config.StreamsBuilderFactoryBean; +import org.springframework.kafka.core.CleanupConfig; +import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.support.MessageBuilder; +import org.springframework.util.StringUtils; + +/** + * @author Soby Chacko + * @since 3.0.0 + */ +public abstract class AbstractKafkaStreamsBinderProcessor implements ApplicationContextAware { + + private static final Log LOG = LogFactory.getLog(AbstractKafkaStreamsBinderProcessor.class); + + private final KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue; + private final BindingServiceProperties bindingServiceProperties; + private final KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties; + private final CleanupConfig cleanupConfig; + private final KeyValueSerdeResolver keyValueSerdeResolver; + + protected ConfigurableApplicationContext applicationContext; + + public AbstractKafkaStreamsBinderProcessor(BindingServiceProperties bindingServiceProperties, + KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, + KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, + KeyValueSerdeResolver keyValueSerdeResolver, CleanupConfig cleanupConfig) { + this.bindingServiceProperties = bindingServiceProperties; + this.kafkaStreamsBindingInformationCatalogue = kafkaStreamsBindingInformationCatalogue; + this.kafkaStreamsExtendedBindingProperties = kafkaStreamsExtendedBindingProperties; + this.keyValueSerdeResolver = keyValueSerdeResolver; + this.cleanupConfig = cleanupConfig; + } + + private KTable materializedAs(StreamsBuilder streamsBuilder, + String destination, String storeName, Serde k, Serde v, + Topology.AutoOffsetReset autoOffsetReset) { + return streamsBuilder.table( + this.bindingServiceProperties.getBindingDestination(destination), + Consumed.with(k, v).withOffsetResetPolicy(autoOffsetReset), + getMaterialized(storeName, k, v)); + } + + protected GlobalKTable materializedAsGlobalKTable( + StreamsBuilder streamsBuilder, String destination, String storeName, + Serde k, Serde v, Topology.AutoOffsetReset autoOffsetReset) { + return streamsBuilder.globalTable( + this.bindingServiceProperties.getBindingDestination(destination), + Consumed.with(k, v).withOffsetResetPolicy(autoOffsetReset), + getMaterialized(storeName, k, v)); + } + + protected GlobalKTable getGlobalKTable(StreamsBuilder streamsBuilder, + Serde keySerde, Serde valueSerde, String materializedAs, + String bindingDestination, Topology.AutoOffsetReset autoOffsetReset) { + return materializedAs != null + ? materializedAsGlobalKTable(streamsBuilder, bindingDestination, + materializedAs, keySerde, valueSerde, autoOffsetReset) + : streamsBuilder.globalTable(bindingDestination, + Consumed.with(keySerde, valueSerde) + .withOffsetResetPolicy(autoOffsetReset)); + } + + protected KTable getKTable(StreamsBuilder streamsBuilder, Serde keySerde, + Serde valueSerde, String materializedAs, String bindingDestination, + Topology.AutoOffsetReset autoOffsetReset) { + return materializedAs != null + ? materializedAs(streamsBuilder, bindingDestination, materializedAs, + keySerde, valueSerde, autoOffsetReset) + : streamsBuilder.table(bindingDestination, + Consumed.with(keySerde, valueSerde) + .withOffsetResetPolicy(autoOffsetReset)); + } + + private Materialized> getMaterialized( + String storeName, Serde k, Serde v) { + return Materialized.>as(storeName) + .withKeySerde(k).withValueSerde(v); + } + + protected Topology.AutoOffsetReset getAutoOffsetReset(String inboundName, KafkaStreamsConsumerProperties extendedConsumerProperties) { + final KafkaConsumerProperties.StartOffset startOffset = extendedConsumerProperties + .getStartOffset(); + Topology.AutoOffsetReset autoOffsetReset = null; + if (startOffset != null) { + switch (startOffset) { + case earliest: + autoOffsetReset = Topology.AutoOffsetReset.EARLIEST; + break; + case latest: + autoOffsetReset = Topology.AutoOffsetReset.LATEST; + break; + default: + break; + } + } + if (extendedConsumerProperties.isResetOffsets()) { + AbstractKafkaStreamsBinderProcessor.LOG.warn("Detected resetOffsets configured on binding " + + inboundName + ". " + + "Setting resetOffsets in Kafka Streams binder does not have any effect."); + } + return autoOffsetReset; + } + + @SuppressWarnings("unchecked") + protected void handleKTableGlobalKTableInputs(Object[] arguments, int index, String input, Class parameterType, Object targetBean, + StreamsBuilderFactoryBean streamsBuilderFactoryBean, StreamsBuilder streamsBuilder, + KafkaStreamsConsumerProperties extendedConsumerProperties, + Serde keySerde, Serde valueSerde, Topology.AutoOffsetReset autoOffsetReset) { + if (parameterType.isAssignableFrom(KTable.class)) { + String materializedAs = extendedConsumerProperties.getMaterializedAs(); + String bindingDestination = this.bindingServiceProperties.getBindingDestination(input); + KTable table = getKTable(streamsBuilder, keySerde, valueSerde, materializedAs, + bindingDestination, autoOffsetReset); + KTableBoundElementFactory.KTableWrapper kTableWrapper = + (KTableBoundElementFactory.KTableWrapper) targetBean; + //wrap the proxy created during the initial target type binding with real object (KTable) + kTableWrapper.wrap((KTable) table); + this.kafkaStreamsBindingInformationCatalogue.addStreamBuilderFactory(streamsBuilderFactoryBean); + arguments[index] = table; + } + else if (parameterType.isAssignableFrom(GlobalKTable.class)) { + String materializedAs = extendedConsumerProperties.getMaterializedAs(); + String bindingDestination = this.bindingServiceProperties.getBindingDestination(input); + GlobalKTable table = getGlobalKTable(streamsBuilder, keySerde, valueSerde, materializedAs, + bindingDestination, autoOffsetReset); + GlobalKTableBoundElementFactory.GlobalKTableWrapper globalKTableWrapper = + (GlobalKTableBoundElementFactory.GlobalKTableWrapper) targetBean; + //wrap the proxy created during the initial target type binding with real object (KTable) + globalKTableWrapper.wrap((GlobalKTable) table); + this.kafkaStreamsBindingInformationCatalogue.addStreamBuilderFactory(streamsBuilderFactoryBean); + arguments[index] = table; + } + } + + @SuppressWarnings({"unchecked"}) + protected StreamsBuilderFactoryBean buildStreamsBuilderAndRetrieveConfig(String beanNamePostPrefix, + ApplicationContext applicationContext, String inboundName) { + ConfigurableListableBeanFactory beanFactory = this.applicationContext + .getBeanFactory(); + + Map streamConfigGlobalProperties = applicationContext + .getBean("streamConfigGlobalProperties", Map.class); + + KafkaStreamsConsumerProperties extendedConsumerProperties = this.kafkaStreamsExtendedBindingProperties + .getExtendedConsumerProperties(inboundName); + streamConfigGlobalProperties + .putAll(extendedConsumerProperties.getConfiguration()); + + String applicationId = extendedConsumerProperties.getApplicationId(); + // override application.id if set at the individual binding level. + if (StringUtils.hasText(applicationId)) { + streamConfigGlobalProperties.put(StreamsConfig.APPLICATION_ID_CONFIG, + applicationId); + } + + int concurrency = this.bindingServiceProperties.getConsumerProperties(inboundName) + .getConcurrency(); + // override concurrency if set at the individual binding level. + if (concurrency > 1) { + streamConfigGlobalProperties.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, + concurrency); + } + + Map kafkaStreamsDlqDispatchers = applicationContext + .getBean("kafkaStreamsDlqDispatchers", Map.class); + + KafkaStreamsConfiguration kafkaStreamsConfiguration = new KafkaStreamsConfiguration( + streamConfigGlobalProperties) { + @Override + public Properties asProperties() { + Properties properties = super.asProperties(); + properties.put(SendToDlqAndContinue.KAFKA_STREAMS_DLQ_DISPATCHERS, + kafkaStreamsDlqDispatchers); + return properties; + } + }; + + StreamsBuilderFactoryBean streamsBuilder = this.cleanupConfig == null + ? new StreamsBuilderFactoryBean(kafkaStreamsConfiguration) + : new StreamsBuilderFactoryBean(kafkaStreamsConfiguration, + this.cleanupConfig); + streamsBuilder.setAutoStartup(false); + BeanDefinition streamsBuilderBeanDefinition = BeanDefinitionBuilder + .genericBeanDefinition( + (Class) streamsBuilder.getClass(), + () -> streamsBuilder) + .getRawBeanDefinition(); + ((BeanDefinitionRegistry) beanFactory).registerBeanDefinition( + "stream-builder-" + beanNamePostPrefix, streamsBuilderBeanDefinition); + + return applicationContext.getBean( + "&stream-builder-" + beanNamePostPrefix, StreamsBuilderFactoryBean.class); + } + + @Override + public final void setApplicationContext(ApplicationContext applicationContext) + throws BeansException { + this.applicationContext = (ConfigurableApplicationContext) applicationContext; + } + + protected KStream getkStream(BindingProperties bindingProperties, KStream stream, boolean nativeDecoding) { + stream = stream.mapValues((value) -> { + Object returnValue; + String contentType = bindingProperties.getContentType(); + if (value != null && !StringUtils.isEmpty(contentType) && !nativeDecoding) { + returnValue = MessageBuilder.withPayload(value) + .setHeader(MessageHeaders.CONTENT_TYPE, contentType).build(); + } + else { + returnValue = value; + } + return returnValue; + }); + return stream; + } + + protected Serde getValueSerde(String inboundName, KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, ResolvableType resolvableType) { + if (bindingServiceProperties.getConsumerProperties(inboundName).isUseNativeDecoding()) { + BindingProperties bindingProperties = this.bindingServiceProperties + .getBindingProperties(inboundName); + return this.keyValueSerdeResolver.getInboundValueSerde( + bindingProperties.getConsumer(), kafkaStreamsConsumerProperties, resolvableType); + } + else { + return Serdes.ByteArray(); + } + } +} 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 e88b29202..d08988bdc 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 @@ -25,7 +25,6 @@ import org.springframework.boot.autoconfigure.kafka.KafkaProperties; 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; -import org.springframework.context.ApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -41,20 +40,7 @@ public class GlobalKTableBinderConfiguration { @Bean @ConditionalOnBean(name = "outerContext") public static BeanFactoryPostProcessor outerContextBeanFactoryPostProcessor() { - return (beanFactory) -> { - // It is safe to call getBean("outerContext") here, because this bean is - // registered as first - // and as independent from the parent context. - ApplicationContext outerContext = (ApplicationContext) beanFactory - .getBean("outerContext"); - beanFactory.registerSingleton( - KafkaStreamsBinderConfigurationProperties.class.getSimpleName(), - outerContext - .getBean(KafkaStreamsBinderConfigurationProperties.class)); - beanFactory.registerSingleton( - KafkaStreamsBindingInformationCatalogue.class.getSimpleName(), - outerContext.getBean(KafkaStreamsBindingInformationCatalogue.class)); - }; + return KafkaStreamsBinderUtils.outerContextBeanFactoryPostProcessor(); } @Bean diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamStreamListenerParameterAdapter.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamStreamListenerParameterAdapter.java index f36b59bff..d2645c957 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamStreamListenerParameterAdapter.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamStreamListenerParameterAdapter.java @@ -44,8 +44,7 @@ class KStreamStreamListenerParameterAdapter @Override public boolean supports(Class bindingTargetType, MethodParameter methodParameter) { - return KStream.class.isAssignableFrom(bindingTargetType) - && KStream.class.isAssignableFrom(methodParameter.getParameterType()); + return KafkaStreamsBinderUtils.supportsKStream(methodParameter, bindingTargetType); } @Override 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 19025bbbe..691626f31 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 @@ -25,7 +25,6 @@ import org.springframework.boot.autoconfigure.kafka.KafkaProperties; 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; -import org.springframework.context.ApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -41,20 +40,7 @@ public class KTableBinderConfiguration { @Bean @ConditionalOnBean(name = "outerContext") public static BeanFactoryPostProcessor outerContextBeanFactoryPostProcessor() { - return (beanFactory) -> { - // It is safe to call getBean("outerContext") here, because this bean is - // registered as first - // and as independent from the parent context. - ApplicationContext outerContext = (ApplicationContext) beanFactory - .getBean("outerContext"); - beanFactory.registerSingleton( - KafkaStreamsBinderConfigurationProperties.class.getSimpleName(), - outerContext - .getBean(KafkaStreamsBinderConfigurationProperties.class)); - beanFactory.registerSingleton( - KafkaStreamsBindingInformationCatalogue.class.getSimpleName(), - outerContext.getBean(KafkaStreamsBindingInformationCatalogue.class)); - }; + return KafkaStreamsBinderUtils.outerContextBeanFactoryPostProcessor(); } @Bean diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java index a56226833..a12dce095 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 @@ -311,12 +311,6 @@ public class KafkaStreamsBinderSupportAutoConfiguration { (Map) streamConfigGlobalProperties, properties); } - @Bean - public QueryableStoreRegistry queryableStoreTypeRegistry( - KafkaStreamsRegistry kafkaStreamsRegistry) { - return new QueryableStoreRegistry(kafkaStreamsRegistry); - } - @Bean public InteractiveQueryService interactiveQueryServices( KafkaStreamsRegistry kafkaStreamsRegistry, diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java index 748711824..d04e99f95 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java @@ -18,12 +18,16 @@ package org.springframework.cloud.stream.binder.kafka.streams; import java.util.Map; +import org.apache.kafka.streams.kstream.KStream; + +import org.springframework.beans.factory.config.BeanFactoryPostProcessor; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; 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.KafkaStreamsConsumerProperties; import org.springframework.context.ApplicationContext; +import org.springframework.core.MethodParameter; import org.springframework.util.StringUtils; /** @@ -83,4 +87,26 @@ final class KafkaStreamsBinderUtils { } } + static boolean supportsKStream(MethodParameter methodParameter, Class targetBeanClass) { + return KStream.class.isAssignableFrom(targetBeanClass) + && KStream.class.isAssignableFrom(methodParameter.getParameterType()); + } + + static BeanFactoryPostProcessor outerContextBeanFactoryPostProcessor() { + return (beanFactory) -> { + // It is safe to call getBean("outerContext") here, because this bean is + // registered as first + // and as independent from the parent context. + ApplicationContext outerContext = (ApplicationContext) beanFactory + .getBean("outerContext"); + beanFactory.registerSingleton( + KafkaStreamsBinderConfigurationProperties.class.getSimpleName(), + outerContext + .getBean(KafkaStreamsBinderConfigurationProperties.class)); + beanFactory.registerSingleton( + KafkaStreamsBindingInformationCatalogue.class.getSimpleName(), + outerContext.getBean(KafkaStreamsBindingInformationCatalogue.class)); + }; + } + } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java index 52eb740da..2918a3095 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java @@ -125,11 +125,11 @@ class KafkaStreamsBindingInformationCatalogue { return this.streamsBuilderFactoryBeans; } - public void setOutboundKStreamResolvable(ResolvableType outboundResolvable) { + void setOutboundKStreamResolvable(ResolvableType outboundResolvable) { this.outboundKStreamResolvable = outboundResolvable; } - public ResolvableType getOutboundKStreamResolvable() { + ResolvableType getOutboundKStreamResolvable() { return outboundKStreamResolvable; } } 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 e6180e7bc..2b63b7c8f 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 @@ -21,7 +21,6 @@ import java.util.HashMap; import java.util.Iterator; import java.util.LinkedHashMap; import java.util.Map; -import java.util.Properties; import java.util.Set; import java.util.TreeSet; import java.util.function.Consumer; @@ -31,43 +30,25 @@ import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.serialization.Serdes; -import org.apache.kafka.common.utils.Bytes; import org.apache.kafka.streams.StreamsBuilder; -import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.Topology; import org.apache.kafka.streams.kstream.Consumed; -import org.apache.kafka.streams.kstream.GlobalKTable; import org.apache.kafka.streams.kstream.KStream; -import org.apache.kafka.streams.kstream.KTable; -import org.apache.kafka.streams.kstream.Materialized; -import org.apache.kafka.streams.state.KeyValueStore; import org.apache.kafka.streams.state.StoreBuilder; -import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanInitializationException; -import org.springframework.beans.factory.config.BeanDefinition; -import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; -import org.springframework.beans.factory.support.BeanDefinitionBuilder; -import org.springframework.beans.factory.support.BeanDefinitionRegistry; 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.properties.KafkaConsumerProperties; 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.context.ApplicationContext; -import org.springframework.context.ApplicationContextAware; -import org.springframework.context.ConfigurableApplicationContext; import org.springframework.core.ResolvableType; -import org.springframework.kafka.config.KafkaStreamsConfiguration; import org.springframework.kafka.config.StreamsBuilderFactoryBean; import org.springframework.kafka.core.CleanupConfig; -import org.springframework.messaging.MessageHeaders; -import org.springframework.messaging.support.MessageBuilder; import org.springframework.util.Assert; import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; @@ -76,7 +57,7 @@ import org.springframework.util.StringUtils; * @author Soby Chacko * @since 2.2.0 */ -public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { +public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderProcessor { private static final Log LOG = LogFactory.getLog(KafkaStreamsFunctionProcessor.class); @@ -86,11 +67,7 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { private final KeyValueSerdeResolver keyValueSerdeResolver; private final KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue; private final KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate; - private final CleanupConfig cleanupConfig; private final FunctionCatalog functionCatalog; - private final BindableProxyFactory bindableProxyFactory; - - private ConfigurableApplicationContext applicationContext; private Set origInputs = new TreeSet<>(); private Set origOutputs = new TreeSet<>(); @@ -105,24 +82,23 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { CleanupConfig cleanupConfig, FunctionCatalog functionCatalog, BindableProxyFactory bindableProxyFactory) { + super(bindingServiceProperties, kafkaStreamsBindingInformationCatalogue, kafkaStreamsExtendedBindingProperties, + keyValueSerdeResolver, cleanupConfig); this.bindingServiceProperties = bindingServiceProperties; this.kafkaStreamsExtendedBindingProperties = kafkaStreamsExtendedBindingProperties; this.keyValueSerdeResolver = keyValueSerdeResolver; this.kafkaStreamsBindingInformationCatalogue = kafkaStreamsBindingInformationCatalogue; this.kafkaStreamsMessageConversionDelegate = kafkaStreamsMessageConversionDelegate; - this.cleanupConfig = cleanupConfig; this.functionCatalog = functionCatalog; - this.bindableProxyFactory = bindableProxyFactory; - this.origInputs.addAll(this.bindableProxyFactory.getInputs()); - this.origOutputs.addAll(this.bindableProxyFactory.getOutputs()); + this.origInputs.addAll(bindableProxyFactory.getInputs()); + this.origOutputs.addAll(bindableProxyFactory.getOutputs()); } private Map buildTypeMap(ResolvableType resolvableType) { int inputCount = 1; ResolvableType resolvableTypeGeneric = resolvableType.getGeneric(1); - while (resolvableTypeGeneric != null && resolvableTypeGeneric.getRawClass() != null && (resolvableTypeGeneric.getRawClass().equals(Function.class) || - resolvableTypeGeneric.getRawClass().equals(Consumer.class))) { + while (resolvableTypeGeneric != null && resolvableTypeGeneric.getRawClass() != null && (functionOrConsumerFound(resolvableTypeGeneric))) { inputCount++; resolvableTypeGeneric = resolvableTypeGeneric.getGeneric(1); } @@ -131,20 +107,15 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { Map resolvableTypeMap = new LinkedHashMap<>(); final Iterator iterator = inputs.iterator(); - final String next = iterator.next(); - resolvableTypeMap.put(next, resolvableType.getGeneric(0)); - origInputs.remove(next); + popuateResolvableTypeMap(resolvableType, resolvableTypeMap, iterator); ResolvableType iterableResType = resolvableType; for (int i = 1; i < inputCount; i++) { if (iterator.hasNext()) { iterableResType = iterableResType.getGeneric(1); if (iterableResType.getRawClass() != null && - (iterableResType.getRawClass().equals(Function.class) || - iterableResType.getRawClass().equals(Consumer.class))) { - final String next1 = iterator.next(); - resolvableTypeMap.put(next1, iterableResType.getGeneric(0)); - origInputs.remove(next1); + functionOrConsumerFound(iterableResType)) { + popuateResolvableTypeMap(iterableResType, resolvableTypeMap, iterator); } } } @@ -152,8 +123,19 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { return resolvableTypeMap; } + private boolean functionOrConsumerFound(ResolvableType iterableResType) { + return iterableResType.getRawClass().equals(Function.class) || + iterableResType.getRawClass().equals(Consumer.class); + } + + private void popuateResolvableTypeMap(ResolvableType resolvableType, Map resolvableTypeMap, Iterator iterator) { + final String next = iterator.next(); + resolvableTypeMap.put(next, resolvableType.getGeneric(0)); + origInputs.remove(next); + } + @SuppressWarnings({ "unchecked", "rawtypes" }) - public void orchestrateFunctionInvoking(ResolvableType resolvableType, String functionName) { + public void setupFunctionInvokerForKafkaStreams(ResolvableType resolvableType, String functionName) { final Map stringResolvableTypeMap = buildTypeMap(resolvableType); Object[] adaptedInboundArguments = adaptAndRetrieveInboundArguments(stringResolvableTypeMap, functionName); try { @@ -176,6 +158,7 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { target = ((FluxedFunction) function).getTarget(); } function = (Function) target; + Assert.isTrue(function != null, "Function bean cannot be null"); Object result = function.apply(adaptedInboundArguments[0]); int i = 1; while (result instanceof Function || result instanceof Consumer) { @@ -192,15 +175,15 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { kafkaStreamsBindingInformationCatalogue.setOutboundKStreamResolvable( outboundResolvableType != null ? outboundResolvableType : resolvableType.getGeneric(1)); final Set outputs = new TreeSet<>(origOutputs); - final Iterator iterator = outputs.iterator(); + final Iterator outboundDefinitionIterator = outputs.iterator(); if (result.getClass().isArray()) { final int length = ((Object[]) result).length; String[] methodAnnotatedOutboundNames = new String[length]; for (int j = 0; j < length; j++) { - if (iterator.hasNext()) { - final String next = iterator.next(); + if (outboundDefinitionIterator.hasNext()) { + final String next = outboundDefinitionIterator.next(); methodAnnotatedOutboundNames[j] = next; this.origOutputs.remove(next); } @@ -216,8 +199,8 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { } } else { - if (iterator.hasNext()) { - final String next = iterator.next(); + if (outboundDefinitionIterator.hasNext()) { + final String next = outboundDefinitionIterator.next(); Object targetBean = this.applicationContext.getBean(next); this.origOutputs.remove(next); @@ -230,7 +213,7 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { } } catch (Exception ex) { - throw new BeanInitializationException("Cannot setup StreamListener for foobar", ex); + throw new BeanInitializationException("Cannot setup function invoker for this Kafka Streams function.", ex); } } @@ -243,13 +226,13 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { Class parameterType = stringResolvableTypeMap.get(input).getRawClass(); if (input != null) { - Assert.isInstanceOf(String.class, input, "Annotation value must be a String"); Object targetBean = applicationContext.getBean(input); BindingProperties bindingProperties = this.bindingServiceProperties.getBindingProperties(input); //Retrieve the StreamsConfig created for this method if available. //Otherwise, create the StreamsBuilderFactory and get the underlying config. if (!this.methodStreamsBuilderFactoryBeanMap.containsKey(functionName)) { - buildStreamsBuilderAndRetrieveConfig(functionName, applicationContext, input); + StreamsBuilderFactoryBean streamsBuilderFactoryBean = buildStreamsBuilderAndRetrieveConfig("stream-builder-" + functionName, applicationContext, input); + this.methodStreamsBuilderFactoryBeanMap.put(functionName, streamsBuilderFactoryBean); } try { StreamsBuilderFactoryBean streamsBuilderFactoryBean = @@ -260,34 +243,10 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { //get state store spec Serde keySerde = this.keyValueSerdeResolver.getInboundKeySerde(extendedConsumerProperties, stringResolvableTypeMap.get(input)); - Serde valueSerde; + Serde valueSerde = bindingServiceProperties.getConsumerProperties(input).isUseNativeDecoding() ? + getValueSerde(input, extendedConsumerProperties, stringResolvableTypeMap.get(input)) : Serdes.ByteArray(); - if (bindingServiceProperties.getConsumerProperties(input).isUseNativeDecoding()) { - valueSerde = this.keyValueSerdeResolver.getInboundValueSerde( - bindingProperties.getConsumer(), extendedConsumerProperties, stringResolvableTypeMap.get(input)); - } - else { - valueSerde = Serdes.ByteArray(); - } - - final KafkaConsumerProperties.StartOffset startOffset = extendedConsumerProperties.getStartOffset(); - Topology.AutoOffsetReset autoOffsetReset = null; - if (startOffset != null) { - switch (startOffset) { - case earliest: - autoOffsetReset = Topology.AutoOffsetReset.EARLIEST; - break; - case latest: - autoOffsetReset = Topology.AutoOffsetReset.LATEST; - break; - default: - break; - } - } - if (extendedConsumerProperties.isResetOffsets()) { - LOG.warn("Detected resetOffsets configured on binding " + input + ". " - + "Setting resetOffsets in Kafka Streams binder does not have any effect."); - } + final Topology.AutoOffsetReset autoOffsetReset = getAutoOffsetReset(input, extendedConsumerProperties); if (parameterType.isAssignableFrom(KStream.class)) { KStream stream = getkStream(input, bindingProperties, @@ -315,31 +274,11 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { if (arguments[i] == null) { arguments[i] = stream; } - Assert.notNull(arguments[i], "problems.."); + Assert.notNull(arguments[i], "Problems encountered while adapting the function argument."); } - else if (parameterType.isAssignableFrom(KTable.class)) { - String materializedAs = extendedConsumerProperties.getMaterializedAs(); - String bindingDestination = this.bindingServiceProperties.getBindingDestination(input); - KTable table = getKTable(streamsBuilder, keySerde, valueSerde, materializedAs, - bindingDestination, autoOffsetReset); - KTableBoundElementFactory.KTableWrapper kTableWrapper = - (KTableBoundElementFactory.KTableWrapper) targetBean; - //wrap the proxy created during the initial target type binding with real object (KTable) - kTableWrapper.wrap((KTable) table); - this.kafkaStreamsBindingInformationCatalogue.addStreamBuilderFactory(streamsBuilderFactoryBean); - arguments[i] = table; - } - else if (parameterType.isAssignableFrom(GlobalKTable.class)) { - String materializedAs = extendedConsumerProperties.getMaterializedAs(); - String bindingDestination = this.bindingServiceProperties.getBindingDestination(input); - GlobalKTable table = getGlobalKTable(streamsBuilder, keySerde, valueSerde, materializedAs, - bindingDestination, autoOffsetReset); - GlobalKTableBoundElementFactory.GlobalKTableWrapper globalKTableWrapper = - (GlobalKTableBoundElementFactory.GlobalKTableWrapper) targetBean; - //wrap the proxy created during the initial target type binding with real object (KTable) - globalKTableWrapper.wrap((GlobalKTable) table); - this.kafkaStreamsBindingInformationCatalogue.addStreamBuilderFactory(streamsBuilderFactoryBean); - arguments[i] = table; + else { + handleKTableGlobalKTableInputs(arguments, i, input, parameterType, targetBean, streamsBuilderFactoryBean, + streamsBuilder, extendedConsumerProperties, keySerde, valueSerde, autoOffsetReset); } i++; } @@ -354,50 +293,6 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { return arguments; } - private GlobalKTable getGlobalKTable(StreamsBuilder streamsBuilder, - Serde keySerde, Serde valueSerde, String materializedAs, - String bindingDestination, Topology.AutoOffsetReset autoOffsetReset) { - return materializedAs != null ? - materializedAsGlobalKTable(streamsBuilder, bindingDestination, materializedAs, - keySerde, valueSerde, autoOffsetReset) : - streamsBuilder.globalTable(bindingDestination, - Consumed.with(keySerde, valueSerde).withOffsetResetPolicy(autoOffsetReset)); - } - - private KTable getKTable(StreamsBuilder streamsBuilder, Serde keySerde, Serde valueSerde, - String materializedAs, - String bindingDestination, Topology.AutoOffsetReset autoOffsetReset) { - return materializedAs != null ? - materializedAs(streamsBuilder, bindingDestination, materializedAs, keySerde, valueSerde, - autoOffsetReset) : - streamsBuilder.table(bindingDestination, - Consumed.with(keySerde, valueSerde).withOffsetResetPolicy(autoOffsetReset)); - } - - private KTable materializedAs(StreamsBuilder streamsBuilder, String destination, - String storeName, Serde k, Serde v, - Topology.AutoOffsetReset autoOffsetReset) { - return streamsBuilder.table(this.bindingServiceProperties.getBindingDestination(destination), - Consumed.with(k, v).withOffsetResetPolicy(autoOffsetReset), - getMaterialized(storeName, k, v)); - } - - private GlobalKTable materializedAsGlobalKTable(StreamsBuilder streamsBuilder, - String destination, String storeName, - Serde k, Serde v, - Topology.AutoOffsetReset autoOffsetReset) { - return streamsBuilder.globalTable(this.bindingServiceProperties.getBindingDestination(destination), - Consumed.with(k, v).withOffsetResetPolicy(autoOffsetReset), - getMaterialized(storeName, k, v)); - } - - private Materialized> getMaterialized(String storeName, - Serde k, Serde v) { - return Materialized.>as(storeName) - .withKeySerde(k) - .withValueSerde(v); - } - private KStream getkStream(String inboundName, BindingProperties bindingProperties, StreamsBuilder streamsBuilder, @@ -435,75 +330,7 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { "Inbound message conversion done by Spring Cloud Stream."); } - stream = stream.mapValues((value) -> { - Object returnValue; - String contentType = bindingProperties.getContentType(); - if (value != null && !StringUtils.isEmpty(contentType) && !nativeDecoding) { - returnValue = MessageBuilder.withPayload(value) - .setHeader(MessageHeaders.CONTENT_TYPE, contentType).build(); - } - else { - returnValue = value; - } - return returnValue; - }); - return stream; + return getkStream(bindingProperties, stream, nativeDecoding); } - @SuppressWarnings({"unchecked"}) - private void buildStreamsBuilderAndRetrieveConfig(String functionName, ApplicationContext applicationContext, - String inboundName) { - ConfigurableListableBeanFactory beanFactory = this.applicationContext.getBeanFactory(); - - Map streamConfigGlobalProperties = applicationContext.getBean("streamConfigGlobalProperties", - Map.class); - - KafkaStreamsConsumerProperties extendedConsumerProperties = this.kafkaStreamsExtendedBindingProperties - .getExtendedConsumerProperties(inboundName); - streamConfigGlobalProperties.putAll(extendedConsumerProperties.getConfiguration()); - - String applicationId = extendedConsumerProperties.getApplicationId(); - //override application.id if set at the individual binding level. - if (StringUtils.hasText(applicationId)) { - streamConfigGlobalProperties.put(StreamsConfig.APPLICATION_ID_CONFIG, applicationId); - } - - int concurrency = this.bindingServiceProperties.getConsumerProperties(inboundName).getConcurrency(); - // override concurrency if set at the individual binding level. - if (concurrency > 1) { - streamConfigGlobalProperties.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, concurrency); - } - - Map kafkaStreamsDlqDispatchers = applicationContext.getBean( - "kafkaStreamsDlqDispatchers", Map.class); - - KafkaStreamsConfiguration kafkaStreamsConfiguration = - new KafkaStreamsConfiguration(streamConfigGlobalProperties) { - @Override - public Properties asProperties() { - Properties properties = super.asProperties(); - properties.put(SendToDlqAndContinue.KAFKA_STREAMS_DLQ_DISPATCHERS, kafkaStreamsDlqDispatchers); - return properties; - } - }; - - StreamsBuilderFactoryBean streamsBuilder = this.cleanupConfig == null - ? new StreamsBuilderFactoryBean(kafkaStreamsConfiguration) - : new StreamsBuilderFactoryBean(kafkaStreamsConfiguration, this.cleanupConfig); - streamsBuilder.setAutoStartup(false); - BeanDefinition streamsBuilderBeanDefinition = - BeanDefinitionBuilder.genericBeanDefinition( - (Class) streamsBuilder.getClass(), () -> streamsBuilder) - .getRawBeanDefinition(); - ((BeanDefinitionRegistry) beanFactory).registerBeanDefinition("stream-builder-" + - functionName, streamsBuilderBeanDefinition); - StreamsBuilderFactoryBean streamsBuilderX = applicationContext.getBean("&stream-builder-" + - functionName, StreamsBuilderFactoryBean.class); - this.methodStreamsBuilderFactoryBeanMap.put(functionName, streamsBuilderX); - } - - @Override - public final void setApplicationContext(ApplicationContext applicationContext) throws BeansException { - this.applicationContext = (ConfigurableApplicationContext) applicationContext; - } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java index 56dc608a4..ac97119f3 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java @@ -23,13 +23,11 @@ import java.util.Collection; import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.Properties; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.serialization.Serdes; -import org.apache.kafka.common.utils.Bytes; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.Topology; @@ -37,20 +35,12 @@ import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.GlobalKTable; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.KTable; -import org.apache.kafka.streams.kstream.Materialized; -import org.apache.kafka.streams.state.KeyValueStore; import org.apache.kafka.streams.state.StoreBuilder; import org.apache.kafka.streams.state.Stores; -import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanInitializationException; -import org.springframework.beans.factory.config.BeanDefinition; -import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; -import org.springframework.beans.factory.support.BeanDefinitionBuilder; -import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.cloud.stream.annotation.Input; import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsStateStore; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsConsumerProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties; @@ -62,17 +52,12 @@ import org.springframework.cloud.stream.binding.StreamListenerSetupMethodOrchest import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.context.ApplicationContext; -import org.springframework.context.ApplicationContextAware; -import org.springframework.context.ConfigurableApplicationContext; import org.springframework.core.MethodParameter; import org.springframework.core.ResolvableType; import org.springframework.core.annotation.AnnotationUtils; -import org.springframework.kafka.config.KafkaStreamsConfiguration; import org.springframework.kafka.config.StreamsBuilderFactoryBean; import org.springframework.kafka.core.CleanupConfig; -import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.handler.annotation.SendTo; -import org.springframework.messaging.support.MessageBuilder; import org.springframework.util.Assert; import org.springframework.util.ObjectUtils; import org.springframework.util.ReflectionUtils; @@ -93,8 +78,8 @@ import org.springframework.util.StringUtils; * @author Lei Chen * @author Gary Russell */ -class KafkaStreamsStreamListenerSetupMethodOrchestrator - implements StreamListenerSetupMethodOrchestrator, ApplicationContextAware { +class KafkaStreamsStreamListenerSetupMethodOrchestrator extends AbstractKafkaStreamsBinderProcessor + implements StreamListenerSetupMethodOrchestrator { private static final Log LOG = LogFactory .getLog(KafkaStreamsStreamListenerSetupMethodOrchestrator.class); @@ -111,13 +96,9 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator private final KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue; - private final Map methodStreamsBuilderFactoryBeanMap = new HashMap<>(); - private final Map> registeredStoresPerMethod = new HashMap<>(); - private final CleanupConfig cleanupConfig; - - private ConfigurableApplicationContext applicationContext; + private final Map methodStreamsBuilderFactoryBeanMap = new HashMap<>(); KafkaStreamsStreamListenerSetupMethodOrchestrator( BindingServiceProperties bindingServiceProperties, @@ -127,13 +108,13 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator StreamListenerParameterAdapter streamListenerParameterAdapter, Collection listenerResultAdapters, CleanupConfig cleanupConfig) { + super(bindingServiceProperties, bindingInformationCatalogue, extendedBindingProperties, keyValueSerdeResolver, cleanupConfig); this.bindingServiceProperties = bindingServiceProperties; this.kafkaStreamsExtendedBindingProperties = extendedBindingProperties; this.keyValueSerdeResolver = keyValueSerdeResolver; this.kafkaStreamsBindingInformationCatalogue = bindingInformationCatalogue; this.streamListenerParameterAdapter = streamListenerParameterAdapter; this.streamListenerResultAdapters = listenerResultAdapters; - this.cleanupConfig = cleanupConfig; } @Override @@ -183,42 +164,33 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator else { Object result = method.invoke(bean, adaptedInboundArguments); - if (result.getClass().isArray()) { - Assert.isTrue( - methodAnnotatedOutboundNames.length == ((Object[]) result).length, - "Result does not match with the number of declared outbounds"); - } - else { - Assert.isTrue(methodAnnotatedOutboundNames.length == 1, - "Result does not match with the number of declared outbounds"); - } - kafkaStreamsBindingInformationCatalogue.setOutboundKStreamResolvable(ResolvableType.forMethodReturnType(method)); - if (result.getClass().isArray()) { - Object[] outboundKStreams = (Object[]) result; - int i = 0; - for (Object outboundKStream : outboundKStreams) { - Object targetBean = this.applicationContext - .getBean(methodAnnotatedOutboundNames[i++]); - for (StreamListenerResultAdapter streamListenerResultAdapter : this.streamListenerResultAdapters) { - if (streamListenerResultAdapter.supports( - outboundKStream.getClass(), targetBean.getClass())) { - streamListenerResultAdapter.adapt(outboundKStream, - targetBean); - break; - } - } + if (methodAnnotatedOutboundNames != null && methodAnnotatedOutboundNames.length > 0) { + if (result.getClass().isArray()) { + Assert.isTrue( + methodAnnotatedOutboundNames.length == ((Object[]) result).length, + "Result does not match with the number of declared outbounds"); + } + else { + Assert.isTrue(methodAnnotatedOutboundNames.length == 1, + "Result does not match with the number of declared outbounds"); } } - else { - Object targetBean = this.applicationContext - .getBean(methodAnnotatedOutboundNames[0]); - for (StreamListenerResultAdapter streamListenerResultAdapter : this.streamListenerResultAdapters) { - if (streamListenerResultAdapter.supports(result.getClass(), - targetBean.getClass())) { - streamListenerResultAdapter.adapt(result, targetBean); - break; + kafkaStreamsBindingInformationCatalogue.setOutboundKStreamResolvable(ResolvableType.forMethodReturnType(method)); + if (methodAnnotatedOutboundNames != null && methodAnnotatedOutboundNames.length > 0) { + if (result.getClass().isArray()) { + Object[] outboundKStreams = (Object[]) result; + int i = 0; + for (Object outboundKStream : outboundKStreams) { + Object targetBean = this.applicationContext + .getBean(methodAnnotatedOutboundNames[i++]); + adaptStreamListenerResult(outboundKStream, targetBean); } } + else { + Object targetBean = this.applicationContext + .getBean(methodAnnotatedOutboundNames[0]); + adaptStreamListenerResult(result, targetBean); + } } } } @@ -228,6 +200,18 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator } } + @SuppressWarnings("unchecked") + private void adaptStreamListenerResult(Object outboundKStream, Object targetBean) { + for (StreamListenerResultAdapter streamListenerResultAdapter : this.streamListenerResultAdapters) { + if (streamListenerResultAdapter.supports( + outboundKStream.getClass(), targetBean.getClass())) { + streamListenerResultAdapter.adapt(outboundKStream, + targetBean); + break; + } + } + } + @Override @SuppressWarnings({"unchecked"}) public Object[] adaptAndRetrieveInboundArguments(Method method, String inboundName, @@ -260,8 +244,10 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator // Otherwise, create the StreamsBuilderFactory and get the underlying // config. if (!this.methodStreamsBuilderFactoryBeanMap.containsKey(method)) { - buildStreamsBuilderAndRetrieveConfig(method, applicationContext, + StreamsBuilderFactoryBean streamsBuilderFactoryBean = buildStreamsBuilderAndRetrieveConfig(method.getDeclaringClass().getSimpleName() + "-" + method.getName(), + applicationContext, inboundName); + this.methodStreamsBuilderFactoryBeanMap.put(method, streamsBuilderFactoryBean); } try { StreamsBuilderFactoryBean streamsBuilderFactoryBean = this.methodStreamsBuilderFactoryBeanMap @@ -274,37 +260,10 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator Serde keySerde = this.keyValueSerdeResolver .getInboundKeySerde(extendedConsumerProperties, ResolvableType.forMethodParameter(methodParameter)); - Serde valueSerde; + Serde valueSerde = bindingServiceProperties.getConsumerProperties(inboundName).isUseNativeDecoding() ? + getValueSerde(inboundName, extendedConsumerProperties, ResolvableType.forMethodParameter(methodParameter)) : Serdes.ByteArray(); - if (bindingServiceProperties.getConsumerProperties(inboundName).isUseNativeDecoding()) { - valueSerde = this.keyValueSerdeResolver.getInboundValueSerde( - bindingProperties.getConsumer(), extendedConsumerProperties, ResolvableType.forMethodParameter(methodParameter)); - } - else { - //keySerde = Serdes.ByteArray(); - valueSerde = Serdes.ByteArray(); - } - - final KafkaConsumerProperties.StartOffset startOffset = extendedConsumerProperties - .getStartOffset(); - Topology.AutoOffsetReset autoOffsetReset = null; - if (startOffset != null) { - switch (startOffset) { - case earliest: - autoOffsetReset = Topology.AutoOffsetReset.EARLIEST; - break; - case latest: - autoOffsetReset = Topology.AutoOffsetReset.LATEST; - break; - default: - break; - } - } - if (extendedConsumerProperties.isResetOffsets()) { - LOG.warn("Detected resetOffsets configured on binding " - + inboundName + ". " - + "Setting resetOffsets in Kafka Streams binder does not have any effect."); - } + Topology.AutoOffsetReset autoOffsetReset = getAutoOffsetReset(inboundName, extendedConsumerProperties); if (parameterType.isAssignableFrom(KStream.class)) { KStream stream = getkStream(inboundName, spec, @@ -333,39 +292,9 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator + method + "from " + stream.getClass() + " to " + parameterType); } - else if (parameterType.isAssignableFrom(KTable.class)) { - String materializedAs = extendedConsumerProperties - .getMaterializedAs(); - String bindingDestination = this.bindingServiceProperties - .getBindingDestination(inboundName); - KTable table = getKTable(streamsBuilder, keySerde, - valueSerde, materializedAs, bindingDestination, - autoOffsetReset); - KTableBoundElementFactory.KTableWrapper kTableWrapper = (KTableBoundElementFactory.KTableWrapper) targetBean; - // wrap the proxy created during the initial target type binding - // with real object (KTable) - kTableWrapper.wrap((KTable) table); - this.kafkaStreamsBindingInformationCatalogue - .addStreamBuilderFactory(streamsBuilderFactoryBean); - arguments[parameterIndex] = table; - } - else if (parameterType.isAssignableFrom(GlobalKTable.class)) { - String materializedAs = extendedConsumerProperties - .getMaterializedAs(); - String bindingDestination = this.bindingServiceProperties - .getBindingDestination(inboundName); - GlobalKTable table = getGlobalKTable(streamsBuilder, - keySerde, valueSerde, materializedAs, bindingDestination, - autoOffsetReset); - // @checkstyle:off - GlobalKTableBoundElementFactory.GlobalKTableWrapper globalKTableWrapper = (GlobalKTableBoundElementFactory.GlobalKTableWrapper) targetBean; - // @checkstyle:on - // wrap the proxy created during the initial target type binding - // with real object (KTable) - globalKTableWrapper.wrap((GlobalKTable) table); - this.kafkaStreamsBindingInformationCatalogue - .addStreamBuilderFactory(streamsBuilderFactoryBean); - arguments[parameterIndex] = table; + else { + handleKTableGlobalKTableInputs(arguments, parameterIndex, inboundName, parameterType, targetBean, streamsBuilderFactoryBean, + streamsBuilder, extendedConsumerProperties, keySerde, valueSerde, autoOffsetReset); } } catch (Exception ex) { @@ -380,52 +309,6 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator return arguments; } - private GlobalKTable getGlobalKTable(StreamsBuilder streamsBuilder, - Serde keySerde, Serde valueSerde, String materializedAs, - String bindingDestination, Topology.AutoOffsetReset autoOffsetReset) { - return materializedAs != null - ? materializedAsGlobalKTable(streamsBuilder, bindingDestination, - materializedAs, keySerde, valueSerde, autoOffsetReset) - : streamsBuilder.globalTable(bindingDestination, - Consumed.with(keySerde, valueSerde) - .withOffsetResetPolicy(autoOffsetReset)); - } - - private KTable getKTable(StreamsBuilder streamsBuilder, Serde keySerde, - Serde valueSerde, String materializedAs, String bindingDestination, - Topology.AutoOffsetReset autoOffsetReset) { - return materializedAs != null - ? materializedAs(streamsBuilder, bindingDestination, materializedAs, - keySerde, valueSerde, autoOffsetReset) - : streamsBuilder.table(bindingDestination, - Consumed.with(keySerde, valueSerde) - .withOffsetResetPolicy(autoOffsetReset)); - } - - private KTable materializedAs(StreamsBuilder streamsBuilder, - String destination, String storeName, Serde k, Serde v, - Topology.AutoOffsetReset autoOffsetReset) { - return streamsBuilder.table( - this.bindingServiceProperties.getBindingDestination(destination), - Consumed.with(k, v).withOffsetResetPolicy(autoOffsetReset), - getMaterialized(storeName, k, v)); - } - - private GlobalKTable materializedAsGlobalKTable( - StreamsBuilder streamsBuilder, String destination, String storeName, - Serde k, Serde v, Topology.AutoOffsetReset autoOffsetReset) { - return streamsBuilder.globalTable( - this.bindingServiceProperties.getBindingDestination(destination), - Consumed.with(k, v).withOffsetResetPolicy(autoOffsetReset), - getMaterialized(storeName, k, v)); - } - - private Materialized> getMaterialized( - String storeName, Serde k, Serde v) { - return Materialized.>as(storeName) - .withKeySerde(k).withValueSerde(v); - } - private StoreBuilder buildStateStore(KafkaStreamsStateStoreProperties spec) { try { @@ -498,86 +381,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator + ". Inbound message conversion done by Spring Cloud Stream."); } - stream = stream.mapValues((value) -> { - Object returnValue; - String contentType = bindingProperties.getContentType(); - if (value != null && !StringUtils.isEmpty(contentType) && !nativeDecoding) { - returnValue = MessageBuilder.withPayload(value) - .setHeader(MessageHeaders.CONTENT_TYPE, contentType).build(); - } - else { - returnValue = value; - } - return returnValue; - }); - return stream; - } - - @SuppressWarnings({"unchecked"}) - private void buildStreamsBuilderAndRetrieveConfig(Method method, - ApplicationContext applicationContext, String inboundName) { - ConfigurableListableBeanFactory beanFactory = this.applicationContext - .getBeanFactory(); - - Map streamConfigGlobalProperties = applicationContext - .getBean("streamConfigGlobalProperties", Map.class); - - KafkaStreamsConsumerProperties extendedConsumerProperties = this.kafkaStreamsExtendedBindingProperties - .getExtendedConsumerProperties(inboundName); - streamConfigGlobalProperties - .putAll(extendedConsumerProperties.getConfiguration()); - - String applicationId = extendedConsumerProperties.getApplicationId(); - // override application.id if set at the individual binding level. - if (StringUtils.hasText(applicationId)) { - streamConfigGlobalProperties.put(StreamsConfig.APPLICATION_ID_CONFIG, - applicationId); - } - - int concurrency = this.bindingServiceProperties.getConsumerProperties(inboundName) - .getConcurrency(); - // override concurrency if set at the individual binding level. - if (concurrency > 1) { - streamConfigGlobalProperties.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, - concurrency); - } - - Map kafkaStreamsDlqDispatchers = applicationContext - .getBean("kafkaStreamsDlqDispatchers", Map.class); - - KafkaStreamsConfiguration kafkaStreamsConfiguration = new KafkaStreamsConfiguration( - streamConfigGlobalProperties) { - @Override - public Properties asProperties() { - Properties properties = super.asProperties(); - properties.put(SendToDlqAndContinue.KAFKA_STREAMS_DLQ_DISPATCHERS, - kafkaStreamsDlqDispatchers); - return properties; - } - }; - - StreamsBuilderFactoryBean streamsBuilder = this.cleanupConfig == null - ? new StreamsBuilderFactoryBean(kafkaStreamsConfiguration) - : new StreamsBuilderFactoryBean(kafkaStreamsConfiguration, - this.cleanupConfig); - streamsBuilder.setAutoStartup(false); - BeanDefinition streamsBuilderBeanDefinition = BeanDefinitionBuilder - .genericBeanDefinition( - (Class) streamsBuilder.getClass(), - () -> streamsBuilder) - .getRawBeanDefinition(); - final String beanNamePostFix = method.getDeclaringClass().getSimpleName() + "-" + method.getName(); - ((BeanDefinitionRegistry) beanFactory).registerBeanDefinition( - "stream-builder-" + beanNamePostFix, streamsBuilderBeanDefinition); - StreamsBuilderFactoryBean streamsBuilderX = applicationContext.getBean( - "&stream-builder-" + beanNamePostFix, StreamsBuilderFactoryBean.class); - this.methodStreamsBuilderFactoryBeanMap.put(method, streamsBuilderX); - } - - @Override - public final void setApplicationContext(ApplicationContext applicationContext) - throws BeansException { - this.applicationContext = (ConfigurableApplicationContext) applicationContext; + return getkStream(bindingProperties, stream, nativeDecoding); } private void validateStreamListenerMethod(StreamListener streamListener, @@ -628,8 +432,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator && this.applicationContext.containsBean(targetBeanName)) { Class targetBeanClass = this.applicationContext.getType(targetBeanName); if (targetBeanClass != null) { - boolean supports = KStream.class.isAssignableFrom(targetBeanClass) - && KStream.class.isAssignableFrom(methodParameter.getParameterType()); + boolean supports = KafkaStreamsBinderUtils.supportsKStream(methodParameter, targetBeanClass); if (!supports) { supports = KTable.class.isAssignableFrom(targetBeanClass) && KTable.class.isAssignableFrom(methodParameter.getParameterType()); diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KeyValueSerdeResolver.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KeyValueSerdeResolver.java index 08c2648c9..cada2d380 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KeyValueSerdeResolver.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KeyValueSerdeResolver.java @@ -227,12 +227,7 @@ public class KeyValueSerdeResolver { keySerde = Utils.newInstance(keySerdeString, Serde.class); } else { - keySerde = this.binderConfigurationProperties.getConfiguration() - .containsKey("default.key.serde") - ? Utils.newInstance(this.binderConfigurationProperties - .getConfiguration().get("default.key.serde"), - Serde.class) - : Serdes.ByteArray(); + keySerde = getFallbackSerde("default.key.serde"); } keySerde.configure(this.streamConfigGlobalProperties, true); @@ -253,15 +248,10 @@ public class KeyValueSerdeResolver { if (resolvableType != null && (isResolvalbeKafkaStreamsType(resolvableType) || isResolvableKStreamArrayType(resolvableType))) { ResolvableType generic = resolvableType.isArray() ? resolvableType.getComponentType().getGeneric(0) : resolvableType.getGeneric(0); - keySerde = getSerde(keySerde, generic); + keySerde = getSerde(generic); } if (keySerde == null) { - keySerde = this.binderConfigurationProperties.getConfiguration() - .containsKey("default.key.serde") - ? Utils.newInstance(this.binderConfigurationProperties - .getConfiguration().get("default.key.serde"), - Serde.class) - : Serdes.ByteArray(); + keySerde = getFallbackSerde("default.key.serde"); } } keySerde.configure(this.streamConfigGlobalProperties, true); @@ -282,34 +272,38 @@ public class KeyValueSerdeResolver { GlobalKTable.class.isAssignableFrom(resolvableType.getRawClass())); } - private Serde getSerde(Serde keySerde, ResolvableType generic) { + private Serde getSerde(ResolvableType generic) { + Serde serde = null; if (generic.getRawClass() != null) { if (Integer.class.isAssignableFrom(generic.getRawClass())) { - keySerde = Serdes.Integer(); + serde = Serdes.Integer(); } else if (Long.class.isAssignableFrom(generic.getRawClass())) { - keySerde = Serdes.Long(); + serde = Serdes.Long(); } else if (Short.class.isAssignableFrom(generic.getRawClass())) { - keySerde = Serdes.Short(); + serde = Serdes.Short(); } else if (Double.class.isAssignableFrom(generic.getRawClass())) { - keySerde = Serdes.Double(); + serde = Serdes.Double(); } else if (Float.class.isAssignableFrom(generic.getRawClass())) { - keySerde = Serdes.Float(); + serde = Serdes.Float(); } else if (byte[].class.isAssignableFrom(generic.getRawClass())) { - keySerde = Serdes.ByteArray(); + serde = Serdes.ByteArray(); } else if (String.class.isAssignableFrom(generic.getRawClass())) { - keySerde = Serdes.String(); + serde = Serdes.String(); } else { - keySerde = new JsonSerde(generic.getRawClass()); + // If the type is Object, then skip assigning the JsonSerde and let the fallback mechanism takes precedence. + if (!generic.getRawClass().isAssignableFrom((Object.class))) { + serde = new JsonSerde(generic.getRawClass()); + } } } - return keySerde; + return serde; } @@ -320,16 +314,20 @@ public class KeyValueSerdeResolver { valueSerde = Utils.newInstance(valueSerdeString, Serde.class); } else { - valueSerde = this.binderConfigurationProperties.getConfiguration() - .containsKey("default.value.serde") - ? Utils.newInstance(this.binderConfigurationProperties - .getConfiguration().get("default.value.serde"), - Serde.class) - : Serdes.ByteArray(); + valueSerde = getFallbackSerde("default.value.serde"); } return valueSerde; } + private Serde getFallbackSerde(String s) throws ClassNotFoundException { + return this.binderConfigurationProperties.getConfiguration() + .containsKey(s) + ? Utils.newInstance(this.binderConfigurationProperties + .getConfiguration().get(s), + Serde.class) + : Serdes.ByteArray(); + } + @SuppressWarnings("unchecked") private Serde getValueSerde(String valueSerdeString, ResolvableType resolvableType) throws ClassNotFoundException { @@ -342,17 +340,12 @@ public class KeyValueSerdeResolver { if (resolvableType != null && ((isResolvalbeKafkaStreamsType(resolvableType)) || (isResolvableKStreamArrayType(resolvableType)))) { ResolvableType generic = resolvableType.isArray() ? resolvableType.getComponentType().getGeneric(1) : resolvableType.getGeneric(1); - valueSerde = getSerde(valueSerde, generic); + valueSerde = getSerde(generic); } if (valueSerde == null) { - valueSerde = this.binderConfigurationProperties.getConfiguration() - .containsKey("default.value.serde") - ? Utils.newInstance(this.binderConfigurationProperties - .getConfiguration().get("default.value.serde"), - Serde.class) - : Serdes.ByteArray(); + valueSerde = getFallbackSerde("default.value.serde"); } } return valueSerde; diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/QueryableStoreRegistry.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/QueryableStoreRegistry.java deleted file mode 100644 index 3b2a44113..000000000 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/QueryableStoreRegistry.java +++ /dev/null @@ -1,66 +0,0 @@ -/* - * Copyright 2018-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; - -import org.apache.kafka.streams.KafkaStreams; -import org.apache.kafka.streams.errors.InvalidStateStoreException; -import org.apache.kafka.streams.state.QueryableStoreType; - -/** - * Registry that contains {@link QueryableStoreType}s those created from the user - * applications. - * - * @author Soby Chacko - * @author Renwei Han - * @since 2.0.0 - * @deprecated in favor of {@link InteractiveQueryService} - */ -public class QueryableStoreRegistry { - - private final KafkaStreamsRegistry kafkaStreamsRegistry; - - public QueryableStoreRegistry(KafkaStreamsRegistry kafkaStreamsRegistry) { - this.kafkaStreamsRegistry = kafkaStreamsRegistry; - } - - /** - * Retrieve and return a queryable store by name created in the application. - * @param storeName name of the queryable store - * @param storeType type of the queryable store - * @param generic queryable store - * @return queryable store. - * @deprecated in favor of - * {@link InteractiveQueryService#getQueryableStore(String, QueryableStoreType)} - */ - public T getQueryableStoreType(String storeName, - QueryableStoreType storeType) { - - for (KafkaStreams kafkaStream : this.kafkaStreamsRegistry.getKafkaStreams()) { - try { - T store = kafkaStream.store(storeName, storeType); - if (store != null) { - return store; - } - } - catch (InvalidStateStoreException ignored) { - // pass through - } - } - return null; - } - -} diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java index 6e4ffa5d1..7b775b6b0 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java @@ -25,7 +25,7 @@ import org.springframework.kafka.config.StreamsBuilderFactoryBean; /** * Iterate through all {@link StreamsBuilderFactoryBean} in the application context and * start them. As each one completes starting, register the associated KafkaStreams object - * into {@link QueryableStoreRegistry}. + * into {@link InteractiveQueryService}. * * This {@link SmartLifecycle} class ensures that the bean created from it is started very * late through the bootstrap process by setting the phase value closer to 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 fa491192b..86b762e59 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 @@ -42,6 +42,6 @@ class KafkaStreamsFunctionProcessorInvoker { @PostConstruct void invoke() { resolvableTypeMap.forEach((key, value) -> - this.kafkaStreamsFunctionProcessor.orchestrateFunctionInvoking(value, key)); + this.kafkaStreamsFunctionProcessor.setupFunctionInvokerForKafkaStreams(value, key)); } }