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 3ac0b87e0..1923b8518 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 @@ -16,6 +16,8 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.util.Map; + import org.apache.kafka.streams.kstream.GlobalKTable; import org.springframework.cloud.stream.binder.AbstractBinder; @@ -49,15 +51,15 @@ public class GlobalKTableBinder extends private final KafkaTopicProvisioner kafkaTopicProvisioner; - private final KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue; + private final Map kafkaStreamsDlqDispatchers; private KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties = new KafkaStreamsExtendedBindingProperties(); public GlobalKTableBinder(KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, KafkaTopicProvisioner kafkaTopicProvisioner, - KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue) { + Map kafkaStreamsDlqDispatchers) { this.binderConfigurationProperties = binderConfigurationProperties; this.kafkaTopicProvisioner = kafkaTopicProvisioner; - this.kafkaStreamsBindingInformationCatalogue = kafkaStreamsBindingInformationCatalogue; + this.kafkaStreamsDlqDispatchers = kafkaStreamsDlqDispatchers; } @Override @@ -67,11 +69,9 @@ public class GlobalKTableBinder extends if (!StringUtils.hasText(group)) { group = binderConfigurationProperties.getApplicationId(); } - KafkaStreamsBinderUtils.prepareConsumerBinding(name, group, inputTarget, - getApplicationContext(), + KafkaStreamsBinderUtils.prepareConsumerBinding(name, group, getApplicationContext(), kafkaTopicProvisioner, - kafkaStreamsBindingInformationCatalogue, - binderConfigurationProperties, properties); + binderConfigurationProperties, properties, kafkaStreamsDlqDispatchers); return new DefaultBinding<>(name, group, inputTarget, null); } 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 1f8451044..917899840 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 @@ -16,6 +16,9 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.util.Map; + +import org.springframework.beans.factory.annotation.Qualifier; 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; @@ -41,8 +44,7 @@ public class GlobalKTableBinderConfiguration { @Bean public GlobalKTableBinder GlobalKTableBinder(KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, KafkaTopicProvisioner kafkaTopicProvisioner, - KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue) { - return new GlobalKTableBinder(binderConfigurationProperties, kafkaTopicProvisioner, - KafkaStreamsBindingInformationCatalogue); + @Qualifier("kafkaStreamsDlqDispatchers") Map kafkaStreamsDlqDispatchers) { + return new GlobalKTableBinder(binderConfigurationProperties, kafkaTopicProvisioner, kafkaStreamsDlqDispatchers); } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java index e8a292623..1b1fe8e36 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java @@ -16,6 +16,8 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.util.Map; + import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.common.serialization.Serde; @@ -64,16 +66,20 @@ class KStreamBinder extends private final KeyValueSerdeResolver keyValueSerdeResolver; + private final Map kafkaStreamsDlqDispatchers; + KStreamBinder(KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, KafkaTopicProvisioner kafkaTopicProvisioner, KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate, KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue, - KeyValueSerdeResolver keyValueSerdeResolver) { + KeyValueSerdeResolver keyValueSerdeResolver, + Map kafkaStreamsDlqDispatchers) { this.binderConfigurationProperties = binderConfigurationProperties; this.kafkaTopicProvisioner = kafkaTopicProvisioner; this.kafkaStreamsMessageConversionDelegate = kafkaStreamsMessageConversionDelegate; this.kafkaStreamsBindingInformationCatalogue = KafkaStreamsBindingInformationCatalogue; this.keyValueSerdeResolver = keyValueSerdeResolver; + this.kafkaStreamsDlqDispatchers = kafkaStreamsDlqDispatchers; } @Override @@ -85,11 +91,9 @@ class KStreamBinder extends if (!StringUtils.hasText(group)) { group = binderConfigurationProperties.getApplicationId(); } - KafkaStreamsBinderUtils.prepareConsumerBinding(name, group, inputTarget, - getApplicationContext(), + KafkaStreamsBinderUtils.prepareConsumerBinding(name, group, getApplicationContext(), kafkaTopicProvisioner, - kafkaStreamsBindingInformationCatalogue, - binderConfigurationProperties, properties); + binderConfigurationProperties, properties, kafkaStreamsDlqDispatchers); return new DefaultBinding<>(name, group, inputTarget, null); } 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 2aa2fc828..2ac8709a6 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 @@ -16,6 +16,9 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.util.Map; + +import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.config.MethodInvokingFactoryBean; import org.springframework.beans.factory.support.AbstractBeanDefinition; import org.springframework.beans.factory.support.BeanDefinitionBuilder; @@ -86,14 +89,15 @@ public class KStreamBinderConfiguration { @Bean public KStreamBinder kStreamBinder(KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, - KafkaTopicProvisioner kafkaTopicProvisioner, - KafkaStreamsMessageConversionDelegate KafkaStreamsMessageConversionDelegate, - KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue, - KeyValueSerdeResolver keyValueSerdeResolver, - KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties) { + KafkaTopicProvisioner kafkaTopicProvisioner, + KafkaStreamsMessageConversionDelegate KafkaStreamsMessageConversionDelegate, + KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue, + KeyValueSerdeResolver keyValueSerdeResolver, + KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, + @Qualifier("kafkaStreamsDlqDispatchers") Map kafkaStreamsDlqDispatchers) { KStreamBinder kStreamBinder = new KStreamBinder(binderConfigurationProperties, kafkaTopicProvisioner, KafkaStreamsMessageConversionDelegate, KafkaStreamsBindingInformationCatalogue, - keyValueSerdeResolver); + keyValueSerdeResolver, kafkaStreamsDlqDispatchers); kStreamBinder.setKafkaStreamsExtendedBindingProperties(kafkaStreamsExtendedBindingProperties); return kStreamBinder; } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java index 23e9d86a9..7d006e76f 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java @@ -16,6 +16,8 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.util.Map; + import org.apache.kafka.streams.kstream.KTable; import org.springframework.cloud.stream.binder.AbstractBinder; @@ -48,15 +50,15 @@ class KTableBinder extends private final KafkaTopicProvisioner kafkaTopicProvisioner; - private final KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue; + private Map kafkaStreamsDlqDispatchers; private KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties = new KafkaStreamsExtendedBindingProperties(); KTableBinder(KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, KafkaTopicProvisioner kafkaTopicProvisioner, - KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue) { + Map kafkaStreamsDlqDispatchers) { this.binderConfigurationProperties = binderConfigurationProperties; this.kafkaTopicProvisioner = kafkaTopicProvisioner; - this.kafkaStreamsBindingInformationCatalogue = kafkaStreamsBindingInformationCatalogue; + this.kafkaStreamsDlqDispatchers = kafkaStreamsDlqDispatchers; } @Override @@ -66,11 +68,9 @@ class KTableBinder extends if (!StringUtils.hasText(group)) { group = binderConfigurationProperties.getApplicationId(); } - KafkaStreamsBinderUtils.prepareConsumerBinding(name, group, inputTarget, - getApplicationContext(), + KafkaStreamsBinderUtils.prepareConsumerBinding(name, group, getApplicationContext(), kafkaTopicProvisioner, - kafkaStreamsBindingInformationCatalogue, - binderConfigurationProperties, properties); + binderConfigurationProperties, properties, kafkaStreamsDlqDispatchers); return new DefaultBinding<>(name, group, inputTarget, null); } 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 e4c4598b0..6425f3c78 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 @@ -16,6 +16,9 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.util.Map; + +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; @@ -56,9 +59,8 @@ public class KTableBinderConfiguration { @Bean public KTableBinder kTableBinder(KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, KafkaTopicProvisioner kafkaTopicProvisioner, - KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue) { - KTableBinder kStreamBinder = new KTableBinder(binderConfigurationProperties, kafkaTopicProvisioner, - KafkaStreamsBindingInformationCatalogue); + @Qualifier("kafkaStreamsDlqDispatchers") Map kafkaStreamsDlqDispatchers) { + KTableBinder kStreamBinder = new KTableBinder(binderConfigurationProperties, kafkaTopicProvisioner, kafkaStreamsDlqDispatchers); return kStreamBinder; } } 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 41b232471..ea7d415fa 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 @@ -17,6 +17,7 @@ package org.springframework.cloud.stream.binder.kafka.streams; import java.util.Collection; +import java.util.HashMap; import java.util.Map; import java.util.Properties; import java.util.stream.Collectors; @@ -109,13 +110,13 @@ public class KafkaStreamsBinderSupportAutoConfiguration { if (binderConfigurationProperties.getSerdeError() == KafkaStreamsBinderConfigurationProperties.SerdeError.logAndContinue) { properties.put(StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, - LogAndContinueExceptionHandler.class); + LogAndContinueExceptionHandler.class.getName()); } else if (binderConfigurationProperties.getSerdeError() == KafkaStreamsBinderConfigurationProperties.SerdeError.logAndFail) { properties.put(StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, - LogAndFailExceptionHandler.class); + LogAndFailExceptionHandler.class.getName()); } else if (binderConfigurationProperties.getSerdeError() == KafkaStreamsBinderConfigurationProperties.SerdeError.sendToDlq) { properties.put(StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, - SendToDlqAndContinue.class); + SendToDlqAndContinue.class.getName()); } if (!ObjectUtils.isEmpty(binderConfigurationProperties.getConfiguration())) { @@ -144,11 +145,10 @@ public class KafkaStreamsBinderSupportAutoConfiguration { KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, KStreamStreamListenerParameterAdapter kafkaStreamListenerParameterAdapter, Collection streamListenerResultAdapters, - KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, ObjectProvider cleanupConfig) { return new KafkaStreamsStreamListenerSetupMethodOrchestrator(bindingServiceProperties, kafkaStreamsExtendedBindingProperties, keyValueSerdeResolver, kafkaStreamsBindingInformationCatalogue, - kafkaStreamListenerParameterAdapter, streamListenerResultAdapters, binderConfigurationProperties, + kafkaStreamListenerParameterAdapter, streamListenerResultAdapters, cleanupConfig.getIfUnique()); } @@ -217,4 +217,9 @@ public class KafkaStreamsBinderSupportAutoConfiguration { return new StreamsBuilderFactoryManager(kafkaStreamsBindingInformationCatalogue, kafkaStreamsRegistry); } + @Bean("kafkaStreamsDlqDispatchers") + public Map dlqDispatchers() { + return new HashMap<>(); + } + } 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 baf51c1a6..bc13d8d43 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 @@ -16,8 +16,7 @@ package org.springframework.cloud.stream.binder.kafka.streams; -import org.apache.kafka.streams.StreamsConfig; -import org.apache.kafka.streams.errors.DeserializationExceptionHandler; +import java.util.Map; import org.springframework.beans.factory.config.MethodInvokingFactoryBean; import org.springframework.beans.factory.support.AbstractBeanDefinition; @@ -38,12 +37,11 @@ import org.springframework.util.StringUtils; */ class KafkaStreamsBinderUtils { - static void prepareConsumerBinding(String name, String group, Object inputTarget, - ApplicationContext context, + static void prepareConsumerBinding(String name, String group, ApplicationContext context, KafkaTopicProvisioner kafkaTopicProvisioner, - KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, - ExtendedConsumerProperties properties) { + ExtendedConsumerProperties properties, + Map kafkaStreamsDlqDispatchers) { ExtendedConsumerProperties extendedConsumerProperties = new ExtendedConsumerProperties<>( properties.getExtension()); if (binderConfigurationProperties.getSerdeError() == KafkaStreamsBinderConfigurationProperties.SerdeError.sendToDlq) { @@ -56,8 +54,6 @@ class KafkaStreamsBinderUtils { } if (extendedConsumerProperties.getExtension().isEnableDlq()) { - StreamsConfig streamsConfig = kafkaStreamsBindingInformationCatalogue.getStreamsConfig(inputTarget); - KafkaStreamsDlqDispatch kafkaStreamsDlqDispatch = !StringUtils.isEmpty(extendedConsumerProperties.getExtension().getDlqName()) ? new KafkaStreamsDlqDispatch(extendedConsumerProperties.getExtension().getDlqName(), binderConfigurationProperties, extendedConsumerProperties.getExtension()) : null; @@ -67,13 +63,11 @@ class KafkaStreamsBinderUtils { kafkaStreamsDlqDispatch = new KafkaStreamsDlqDispatch(dlqName, binderConfigurationProperties, extendedConsumerProperties.getExtension()); } + SendToDlqAndContinue sendToDlqAndContinue = context.getBean(SendToDlqAndContinue.class); sendToDlqAndContinue.addKStreamDlqDispatch(inputTopic, kafkaStreamsDlqDispatch); - DeserializationExceptionHandler deserializationExceptionHandler = streamsConfig.defaultDeserializationExceptionHandler(); - if (deserializationExceptionHandler instanceof SendToDlqAndContinue) { - ((SendToDlqAndContinue) deserializationExceptionHandler).addKStreamDlqDispatch(inputTopic, kafkaStreamsDlqDispatch); - } + kafkaStreamsDlqDispatchers.put(inputTopic, kafkaStreamsDlqDispatch); } } } 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 cf0c0eba4..d48d70bbf 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 @@ -43,8 +43,6 @@ class KafkaStreamsBindingInformationCatalogue { private final Map, KafkaStreamsConsumerProperties> consumerProperties = new ConcurrentHashMap<>(); - private final Map streamsConfigs = new ConcurrentHashMap<>(); - private final Set streamsBuilderFactoryBeans = new HashSet<>(); /** @@ -94,16 +92,6 @@ class KafkaStreamsBindingInformationCatalogue { return bindingProperties.getContentType(); } - /** - * Retrieve and return the registered {@link StreamsBuilderFactoryBean} for the given KStream - * - * @param bindingTarget KStream binding target - * @return corresponding {@link StreamsBuilderFactoryBean} - */ - StreamsConfig getStreamsConfig(Object bindingTarget) { - return streamsConfigs.get(bindingTarget); - } - /** * Register a cache for bounded KStream -> {@link BindingProperties} * @@ -133,10 +121,6 @@ class KafkaStreamsBindingInformationCatalogue { this.streamsBuilderFactoryBeans.add(streamsBuilderFactoryBean); } - void addStreamsConfigs(Object bindingTarget, StreamsConfig streamsConfig) { - this.streamsConfigs.put(bindingTarget, streamsConfig); - } - Set getStreamsBuilderFactoryBeans() { return streamsBuilderFactoryBeans; } 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 19f288eab..8ecde8f98 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 @@ -21,6 +21,7 @@ import java.util.Arrays; import java.util.Collection; import java.util.HashMap; import java.util.Map; +import java.util.Properties; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -28,7 +29,6 @@ import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.utils.Bytes; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; -import org.apache.kafka.streams.errors.DeserializationExceptionHandler; import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.GlobalKTable; import org.apache.kafka.streams.kstream.KStream; @@ -53,7 +53,6 @@ import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.cloud.stream.binder.BinderSpecificPropertiesProvider; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsStateStore; -import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; 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.binder.kafka.streams.properties.KafkaStreamsStateStoreProperties; @@ -70,6 +69,7 @@ import org.springframework.context.ConfigurableApplicationContext; import org.springframework.core.MethodParameter; import org.springframework.core.annotation.AnnotationUtils; import org.springframework.integration.support.utils.IntegrationUtils; +import org.springframework.kafka.config.KafkaStreamsConfiguration; import org.springframework.kafka.config.StreamsBuilderFactoryBean; import org.springframework.kafka.core.CleanupConfig; import org.springframework.messaging.MessageHeaders; @@ -112,8 +112,6 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene private final Map methodStreamsBuilderFactoryBeanMap = new HashMap<>(); - private final KafkaStreamsBinderConfigurationProperties binderConfigurationProperties; - private final CleanupConfig cleanupConfig; private ConfigurableApplicationContext applicationContext; @@ -124,7 +122,6 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, StreamListenerParameterAdapter streamListenerParameterAdapter, Collection streamListenerResultAdapters, - KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, CleanupConfig cleanupConfig) { this.bindingServiceProperties = bindingServiceProperties; this.kafkaStreamsExtendedBindingProperties = kafkaStreamsExtendedBindingProperties; @@ -132,7 +129,6 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene this.kafkaStreamsBindingInformationCatalogue = kafkaStreamsBindingInformationCatalogue; this.streamListenerParameterAdapter = streamListenerParameterAdapter; this.streamListenerResultAdapters = streamListenerResultAdapters; - this.binderConfigurationProperties = binderConfigurationProperties; this.cleanupConfig = cleanupConfig; } @@ -239,11 +235,10 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene Object targetBean = applicationContext.getBean((String) targetReferenceValue); BindingProperties bindingProperties = bindingServiceProperties.getBindingProperties(inboundName); enableNativeDecodingForKTableAlways(parameterType, bindingProperties); - StreamsConfig streamsConfig = null; //Retrieve the StreamsConfig created for this method if available. //Otherwise, create the StreamsBuilderFactory and get the underlying config. if (!methodStreamsBuilderFactoryBeanMap.containsKey(method)) { - streamsConfig = buildStreamsBuilderAndRetrieveConfig(method, applicationContext, inboundName); + buildStreamsBuilderAndRetrieveConfig(method, applicationContext, inboundName); } try { StreamsBuilderFactoryBean streamsBuilderFactoryBean = methodStreamsBuilderFactoryBeanMap.get(method); @@ -259,9 +254,6 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene //wrap the proxy created during the initial target type binding with real object (KStream) kStreamWrapper.wrap((KStream) stream); kafkaStreamsBindingInformationCatalogue.addStreamBuilderFactory(streamsBuilderFactoryBean); - if (streamsConfig != null){ - kafkaStreamsBindingInformationCatalogue.addStreamsConfigs(kStreamWrapper, streamsConfig); - } for (StreamListenerParameterAdapter streamListenerParameterAdapter : streamListenerParameterAdapters) { if (streamListenerParameterAdapter.supports(stream.getClass(), methodParameter)) { arguments[parameterIndex] = streamListenerParameterAdapter.adapt(kStreamWrapper, methodParameter); @@ -285,9 +277,6 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene //wrap the proxy created during the initial target type binding with real object (KTable) kTableWrapper.wrap((KTable) table); kafkaStreamsBindingInformationCatalogue.addStreamBuilderFactory(streamsBuilderFactoryBean); - if (streamsConfig != null){ - kafkaStreamsBindingInformationCatalogue.addStreamsConfigs(kTableWrapper, streamsConfig); - } arguments[parameterIndex] = table; } else if (parameterType.isAssignableFrom(GlobalKTable.class)) { @@ -301,9 +290,6 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene //wrap the proxy created during the initial target type binding with real object (KTable) globalKTableWrapper.wrap((GlobalKTable) table); kafkaStreamsBindingInformationCatalogue.addStreamBuilderFactory(streamsBuilderFactoryBean); - if (streamsConfig != null){ - kafkaStreamsBindingInformationCatalogue.addStreamsConfigs(globalKTableWrapper, streamsConfig); - } arguments[parameterIndex] = table; } } @@ -417,7 +403,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene } @SuppressWarnings({"unchecked"}) - private StreamsConfig buildStreamsBuilderAndRetrieveConfig(Method method, ApplicationContext applicationContext, + private void buildStreamsBuilderAndRetrieveConfig(Method method, ApplicationContext applicationContext, String inboundName) { ConfigurableListableBeanFactory beanFactory = this.applicationContext.getBeanFactory(); @@ -435,41 +421,27 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator implements StreamListene streamConfigGlobalProperties.put(StreamsConfig.APPLICATION_ID_CONFIG, applicationId); } - //Custom StreamsConfig implementation that overrides to guarantee that the deserialization handler is cached. - StreamsConfig streamsConfig = new StreamsConfig(streamConfigGlobalProperties) { - DeserializationExceptionHandler deserializationExceptionHandler; + Map kafkaStreamsDlqDispatchers = applicationContext.getBean("kafkaStreamsDlqDispatchers", Map.class); + + KafkaStreamsConfiguration kafkaStreamsConfiguration = new KafkaStreamsConfiguration(streamConfigGlobalProperties) { @Override - @SuppressWarnings("unchecked") - public T getConfiguredInstance(String key, Class clazz) { - if (key.equals(StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG)){ - if (deserializationExceptionHandler != null){ - return (T)deserializationExceptionHandler; - } - else { - T t = super.getConfiguredInstance(key, clazz); - deserializationExceptionHandler = (DeserializationExceptionHandler)t; - return t; - } - } - return super.getConfiguredInstance(key, clazz); + public Properties asProperties() { + Properties properties = super.asProperties(); + properties.put(SendToDlqAndContinue.KAFKA_STREAMS_DLQ_DISPATCHERS, kafkaStreamsDlqDispatchers); + return properties; } }; + StreamsBuilderFactoryBean streamsBuilder = this.cleanupConfig == null - ? new StreamsBuilderFactoryBean(streamsConfig) - : new StreamsBuilderFactoryBean(streamsConfig, this.cleanupConfig); + ? 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-" + method.getName(), streamsBuilderBeanDefinition); StreamsBuilderFactoryBean streamsBuilderX = applicationContext.getBean("&stream-builder-" + method.getName(), StreamsBuilderFactoryBean.class); - BeanDefinition streamsConfigBeanDefinition = - BeanDefinitionBuilder.genericBeanDefinition((Class) streamsConfig.getClass(), () -> streamsConfig) - .getRawBeanDefinition(); - ((BeanDefinitionRegistry) beanFactory).registerBeanDefinition("streamsConfig-" + method.getName(), streamsConfigBeanDefinition); - methodStreamsBuilderFactoryBeanMap.put(method, streamsBuilderX); - return streamsConfig; } // This method is mostly copied from core. We should refactor the original method in core so that it is publicly available diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/SendToDlqAndContinue.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/SendToDlqAndContinue.java index 8d2225e12..6d2f6e6de 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/SendToDlqAndContinue.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/SendToDlqAndContinue.java @@ -41,6 +41,8 @@ import org.springframework.util.ReflectionUtils; */ public class SendToDlqAndContinue implements DeserializationExceptionHandler{ + public static final String KAFKA_STREAMS_DLQ_DISPATCHERS = "spring.cloud.stream.kafka.streams.dlq.dispatchers"; + /** * DLQ dispatcher per topic in the application context. The key here is not the actual DLQ topic * but the incoming topic that caused the error. @@ -56,14 +58,14 @@ public class SendToDlqAndContinue implements DeserializationExceptionHandler{ * @param partition for the topic where this record should be sent */ public void sendToDlq(String topic, byte[] key, byte[] value, int partition){ - KafkaStreamsDlqDispatch kafkaStreamsDlqDispatch = dlqDispatchers.get(topic); + KafkaStreamsDlqDispatch kafkaStreamsDlqDispatch = this.dlqDispatchers.get(topic); kafkaStreamsDlqDispatch.sendToDlq(key,value, partition); } @Override @SuppressWarnings("unchecked") public DeserializationHandlerResponse handle(ProcessorContext context, ConsumerRecord record, Exception exception) { - KafkaStreamsDlqDispatch kafkaStreamsDlqDispatch = dlqDispatchers.get(record.topic()); + KafkaStreamsDlqDispatch kafkaStreamsDlqDispatch = this.dlqDispatchers.get(record.topic()); kafkaStreamsDlqDispatch.sendToDlq(record.key(), record.value(), record.partition()); context.commit(); @@ -98,11 +100,13 @@ public class SendToDlqAndContinue implements DeserializationExceptionHandler{ } @Override + @SuppressWarnings("unchecked") public void configure(Map configs) { - + this.dlqDispatchers = (Map) configs.get(KAFKA_STREAMS_DLQ_DISPATCHERS); } void addKStreamDlqDispatch(String topic, KafkaStreamsDlqDispatch kafkaStreamsDlqDispatch){ - dlqDispatchers.put(topic, kafkaStreamsDlqDispatch); + this.dlqDispatchers.put(topic, kafkaStreamsDlqDispatch); } + }