From f5739a9c7f7a8566474e2d1c393cdbe5d8237596 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Thu, 20 Sep 2018 09:38:42 -0400 Subject: [PATCH] Missing beans registration in Kafka Streams binder There are duplicate code in various binder configurations where we register missing beans. Consolidate them into a common class that implements ImportBeanDefinitionRegistrar and then import this class in the binder configurations. Resolves #445 --- .../kafka/streams/GlobalKTableBinder.java | 2 +- .../GlobalKTableBinderConfiguration.java | 36 +---------- .../binder/kafka/streams/KStreamBinder.java | 2 +- .../streams/KStreamBinderConfiguration.java | 62 +++++++++++++------ .../binder/kafka/streams/KTableBinder.java | 2 +- .../streams/KTableBinderConfiguration.java | 2 + ...tils.java => KafkaStreamsBinderUtils.java} | 37 ++++++++++- 7 files changed, 84 insertions(+), 59 deletions(-) rename spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/{KafkaStreamsConsumerBindingUtils.java => KafkaStreamsBinderUtils.java} (68%) 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 c535e694d..3ac0b87e0 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 @@ -67,7 +67,7 @@ public class GlobalKTableBinder extends if (!StringUtils.hasText(group)) { group = binderConfigurationProperties.getApplicationId(); } - KafkaStreamsConsumerBindingUtils.prepareConsumerBinding(name, group, inputTarget, + KafkaStreamsBinderUtils.prepareConsumerBinding(name, group, inputTarget, getApplicationContext(), kafkaTopicProvisioner, kafkaStreamsBindingInformationCatalogue, 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 5ec5eef73..1f8451044 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,10 +16,6 @@ package org.springframework.cloud.stream.binder.kafka.streams; -import org.springframework.beans.factory.config.MethodInvokingFactoryBean; -import org.springframework.beans.factory.support.AbstractBeanDefinition; -import org.springframework.beans.factory.support.BeanDefinitionBuilder; -import org.springframework.beans.factory.support.BeanDefinitionRegistry; 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; @@ -27,45 +23,15 @@ import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStr import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; -import org.springframework.context.annotation.ImportBeanDefinitionRegistrar; -import org.springframework.core.type.AnnotationMetadata; /** * @author Soby Chacko * @since 2.1.0 */ @Configuration -@Import(GlobalKTableBinderConfiguration.Registrar.class) +@Import(KafkaStreamsBinderUtils.KafkaStreamsMissingBeansRegistrar.class) public class GlobalKTableBinderConfiguration { - static class Registrar implements ImportBeanDefinitionRegistrar { - - private static final String BEAN_NAME = "outerContext"; - - @Override - public void registerBeanDefinitions(AnnotationMetadata importingClassMetadata, - BeanDefinitionRegistry registry) { - if (registry.containsBeanDefinition(BEAN_NAME)) { - - AbstractBeanDefinition configBean = BeanDefinitionBuilder.genericBeanDefinition(MethodInvokingFactoryBean.class) - .addPropertyReference("targetObject", BEAN_NAME) - .addPropertyValue("targetMethod", "getBean") - .addPropertyValue("arguments", KafkaStreamsBinderConfigurationProperties.class) - .getBeanDefinition(); - - registry.registerBeanDefinition(KafkaStreamsBinderConfigurationProperties.class.getSimpleName(), configBean); - - AbstractBeanDefinition catalogueBean = BeanDefinitionBuilder.genericBeanDefinition(MethodInvokingFactoryBean.class) - .addPropertyReference("targetObject", BEAN_NAME) - .addPropertyValue("targetMethod", "getBean") - .addPropertyValue("arguments", KafkaStreamsBindingInformationCatalogue.class) - .getBeanDefinition(); - - registry.registerBeanDefinition(KafkaStreamsBindingInformationCatalogue.class.getSimpleName(), catalogueBean); - } - } - } - @Bean public KafkaTopicProvisioner provisioningProvider(KafkaBinderConfigurationProperties binderConfigurationProperties, KafkaProperties kafkaProperties) { 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 60bb5619a..e8a292623 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 @@ -85,7 +85,7 @@ class KStreamBinder extends if (!StringUtils.hasText(group)) { group = binderConfigurationProperties.getApplicationId(); } - KafkaStreamsConsumerBindingUtils.prepareConsumerBinding(name, group, inputTarget, + KafkaStreamsBinderUtils.prepareConsumerBinding(name, group, inputTarget, getApplicationContext(), kafkaTopicProvisioner, kafkaStreamsBindingInformationCatalogue, 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 dbf76ba81..2aa2fc828 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,18 +16,20 @@ package org.springframework.cloud.stream.binder.kafka.streams; -import org.springframework.beans.factory.config.BeanFactoryPostProcessor; -import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; +import org.springframework.beans.factory.config.MethodInvokingFactoryBean; +import org.springframework.beans.factory.support.AbstractBeanDefinition; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration; 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.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties; -import org.springframework.context.ApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; +import org.springframework.core.type.AnnotationMetadata; /** * @author Marius Bogoevici @@ -35,25 +37,45 @@ import org.springframework.context.annotation.Import; * @author Soby Chacko */ @Configuration -@Import({KafkaAutoConfiguration.class}) +@Import({KafkaAutoConfiguration.class, KStreamBinderConfiguration.KStreamMissingBeansRegistrar.class}) public class KStreamBinderConfiguration { - @Bean - @ConditionalOnBean(name = "outerContext") - public BeanFactoryPostProcessor outerContextBeanFactoryPostProcessor() { - return beanFactory -> { - ApplicationContext outerContext = (ApplicationContext) beanFactory.getBean("outerContext"); - beanFactory.registerSingleton(KafkaStreamsBinderConfigurationProperties.class.getSimpleName(), outerContext - .getBean(KafkaStreamsBinderConfigurationProperties.class)); - beanFactory.registerSingleton(KafkaStreamsMessageConversionDelegate.class.getSimpleName(), outerContext - .getBean(KafkaStreamsMessageConversionDelegate.class)); - beanFactory.registerSingleton(KafkaStreamsBindingInformationCatalogue.class.getSimpleName(), outerContext - .getBean(KafkaStreamsBindingInformationCatalogue.class)); - beanFactory.registerSingleton(KeyValueSerdeResolver.class.getSimpleName(), outerContext - .getBean(KeyValueSerdeResolver.class)); - beanFactory.registerSingleton(KafkaStreamsExtendedBindingProperties.class.getSimpleName(), outerContext - .getBean(KafkaStreamsExtendedBindingProperties.class)); - }; + static class KStreamMissingBeansRegistrar extends KafkaStreamsBinderUtils.KafkaStreamsMissingBeansRegistrar { + + private static final String BEAN_NAME = "outerContext"; + + @Override + public void registerBeanDefinitions(AnnotationMetadata importingClassMetadata, + BeanDefinitionRegistry registry) { + super.registerBeanDefinitions(importingClassMetadata, registry); + + if (registry.containsBeanDefinition(BEAN_NAME)) { + + AbstractBeanDefinition converstionDelegateBean = BeanDefinitionBuilder.genericBeanDefinition(MethodInvokingFactoryBean.class) + .addPropertyReference("targetObject", BEAN_NAME) + .addPropertyValue("targetMethod", "getBean") + .addPropertyValue("arguments", KafkaStreamsMessageConversionDelegate.class) + .getBeanDefinition(); + + registry.registerBeanDefinition(KafkaStreamsMessageConversionDelegate.class.getSimpleName(), converstionDelegateBean); + + AbstractBeanDefinition keyValueSerdeResolverBean = BeanDefinitionBuilder.genericBeanDefinition(MethodInvokingFactoryBean.class) + .addPropertyReference("targetObject", BEAN_NAME) + .addPropertyValue("targetMethod", "getBean") + .addPropertyValue("arguments", KeyValueSerdeResolver.class) + .getBeanDefinition(); + + registry.registerBeanDefinition(KeyValueSerdeResolver.class.getSimpleName(), keyValueSerdeResolverBean); + + AbstractBeanDefinition kafkaStreamsExtendedBindingPropertiesBean = BeanDefinitionBuilder.genericBeanDefinition(MethodInvokingFactoryBean.class) + .addPropertyReference("targetObject", BEAN_NAME) + .addPropertyValue("targetMethod", "getBean") + .addPropertyValue("arguments", KafkaStreamsExtendedBindingProperties.class) + .getBeanDefinition(); + + registry.registerBeanDefinition(KafkaStreamsExtendedBindingProperties.class.getSimpleName(), kafkaStreamsExtendedBindingPropertiesBean); + } + } } @Bean 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 835436a3d..23e9d86a9 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 @@ -66,7 +66,7 @@ class KTableBinder extends if (!StringUtils.hasText(group)) { group = binderConfigurationProperties.getApplicationId(); } - KafkaStreamsConsumerBindingUtils.prepareConsumerBinding(name, group, inputTarget, + KafkaStreamsBinderUtils.prepareConsumerBinding(name, group, inputTarget, getApplicationContext(), kafkaTopicProvisioner, kafkaStreamsBindingInformationCatalogue, 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 ca8b7d96a..e4c4598b0 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,12 +25,14 @@ import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStr import org.springframework.context.ApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Import; /** * @author Soby Chacko */ @SuppressWarnings("ALL") @Configuration +@Import(KafkaStreamsBinderUtils.KafkaStreamsMissingBeansRegistrar.class) public class KTableBinderConfiguration { @Bean diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsConsumerBindingUtils.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java similarity index 68% rename from spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsConsumerBindingUtils.java rename to spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java index 1e5f991d8..baf51c1a6 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsConsumerBindingUtils.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java @@ -19,18 +19,24 @@ package org.springframework.cloud.stream.binder.kafka.streams; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.errors.DeserializationExceptionHandler; +import org.springframework.beans.factory.config.MethodInvokingFactoryBean; +import org.springframework.beans.factory.support.AbstractBeanDefinition; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; 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.context.annotation.ImportBeanDefinitionRegistrar; +import org.springframework.core.type.AnnotationMetadata; import org.springframework.util.StringUtils; /** * @author Soby Chacko */ -class KafkaStreamsConsumerBindingUtils { +class KafkaStreamsBinderUtils { static void prepareConsumerBinding(String name, String group, Object inputTarget, ApplicationContext context, @@ -71,4 +77,33 @@ class KafkaStreamsConsumerBindingUtils { } } } + + static class KafkaStreamsMissingBeansRegistrar implements ImportBeanDefinitionRegistrar { + + private static final String BEAN_NAME = "outerContext"; + + @Override + public void registerBeanDefinitions(AnnotationMetadata importingClassMetadata, + BeanDefinitionRegistry registry) { + if (registry.containsBeanDefinition(BEAN_NAME)) { + + AbstractBeanDefinition configBean = BeanDefinitionBuilder.genericBeanDefinition(MethodInvokingFactoryBean.class) + .addPropertyReference("targetObject", BEAN_NAME) + .addPropertyValue("targetMethod", "getBean") + .addPropertyValue("arguments", KafkaStreamsBinderConfigurationProperties.class) + .getBeanDefinition(); + + registry.registerBeanDefinition(KafkaStreamsBinderConfigurationProperties.class.getSimpleName(), configBean); + + AbstractBeanDefinition catalogueBean = BeanDefinitionBuilder.genericBeanDefinition(MethodInvokingFactoryBean.class) + .addPropertyReference("targetObject", BEAN_NAME) + .addPropertyValue("targetMethod", "getBean") + .addPropertyValue("arguments", KafkaStreamsBindingInformationCatalogue.class) + .getBeanDefinition(); + + registry.registerBeanDefinition(KafkaStreamsBindingInformationCatalogue.class.getSimpleName(), catalogueBean); + } + } + } + }