From 4e250b34cf45f4192601a9acbfab31a5edcd326c Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 4 Sep 2019 18:10:00 -0400 Subject: [PATCH] Multi binder issues for Kafka Streams table types When KTable or GlobalKTable binders are used in a multi bindder environment, it has difficulty finding certain beans/properties. Addressing this issue. Adding tests. Resolves #681 --- .../GlobalKTableBinderConfiguration.java | 31 +++++-- .../streams/KTableBinderConfiguration.java | 32 ++++++-- .../streams/KafkaStreamsBinderUtils.java | 18 ----- .../KafkaStreamsBinderBootstrapTest.java | 80 ++++++++++++++----- 4 files changed, 113 insertions(+), 48 deletions(-) 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 687971ac4..e8b65664a 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 @@ -21,13 +21,16 @@ 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.KafkaAutoConfiguration; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.cloud.stream.annotation.BindingProvider; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties; +import org.springframework.context.ApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Import; /** * Configuration for GlobalKTable binder. @@ -37,14 +40,10 @@ import org.springframework.context.annotation.Configuration; */ @Configuration @BindingProvider +@Import({ KafkaAutoConfiguration.class, + KafkaStreamsBinderHealthIndicatorConfiguration.class }) public class GlobalKTableBinderConfiguration { - @Bean - @ConditionalOnBean(name = "outerContext") - public static BeanFactoryPostProcessor outerContextBeanFactoryPostProcessor() { - return KafkaStreamsBinderUtils.outerContextBeanFactoryPostProcessor(); - } - @Bean public KafkaTopicProvisioner provisioningProvider( KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, @@ -65,4 +64,24 @@ public class GlobalKTableBinderConfiguration { return globalKTableBinder; } + @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( + KafkaStreamsExtendedBindingProperties.class.getSimpleName(), + outerContext.getBean(KafkaStreamsExtendedBindingProperties.class)); + }; + } + } 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 978e116d8..12e7ef427 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 @@ -21,13 +21,16 @@ 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.KafkaAutoConfiguration; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.cloud.stream.annotation.BindingProvider; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties; +import org.springframework.context.ApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Import; /** * Configuration for KTable binder. @@ -37,14 +40,10 @@ import org.springframework.context.annotation.Configuration; @SuppressWarnings("ALL") @Configuration @BindingProvider +@Import({ KafkaAutoConfiguration.class, + KafkaStreamsBinderHealthIndicatorConfiguration.class }) public class KTableBinderConfiguration { - @Bean - @ConditionalOnBean(name = "outerContext") - public static BeanFactoryPostProcessor outerContextBeanFactoryPostProcessor() { - return KafkaStreamsBinderUtils.outerContextBeanFactoryPostProcessor(); - } - @Bean public KafkaTopicProvisioner provisioningProvider( KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, @@ -64,4 +63,25 @@ public class KTableBinderConfiguration { return kTableBinder; } + @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( + KafkaStreamsExtendedBindingProperties.class.getSimpleName(), + outerContext.getBean(KafkaStreamsExtendedBindingProperties.class)); + }; + } + + } 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 59ad90035..4f7fcb883 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 @@ -20,14 +20,12 @@ 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.context.support.GenericApplicationContext; import org.springframework.core.MethodParameter; import org.springframework.util.StringUtils; @@ -92,20 +90,4 @@ final class KafkaStreamsBinderUtils { 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 first and is independent from the parent context. - GenericApplicationContext outerContext = (GenericApplicationContext) beanFactory - .getBean("outerContext"); - - outerContext.registerBean(KafkaStreamsBinderConfigurationProperties.class, - () -> outerContext.getBean(KafkaStreamsBinderConfigurationProperties.class)); - outerContext.registerBean(KafkaStreamsBindingInformationCatalogue.class, - () -> outerContext.getBean(KafkaStreamsBindingInformationCatalogue.class)); - - }; - } - } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/bootstrap/KafkaStreamsBinderBootstrapTest.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/bootstrap/KafkaStreamsBinderBootstrapTest.java index 0c1cb267a..5ed9c0b64 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/bootstrap/KafkaStreamsBinderBootstrapTest.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/bootstrap/KafkaStreamsBinderBootstrapTest.java @@ -16,7 +16,9 @@ package org.springframework.cloud.stream.binder.kafka.streams.bootstrap; +import org.apache.kafka.streams.kstream.GlobalKTable; import org.apache.kafka.streams.kstream.KStream; +import org.apache.kafka.streams.kstream.KTable; import org.junit.ClassRule; import org.junit.Test; @@ -38,17 +40,33 @@ public class KafkaStreamsBinderBootstrapTest { public static EmbeddedKafkaRule embeddedKafka = new EmbeddedKafkaRule(1, true, 10); @Test - public void testKafkaStreamsBinderWithCustomEnvironmentCanStart() { + public void testKStreamBinderWithCustomEnvironmentCanStart() { ConfigurableApplicationContext applicationContext = new SpringApplicationBuilder( - SimpleApplication.class).web(WebApplicationType.NONE).run( - "--spring.cloud.stream.kafka.streams.default.consumer.application-id" - + "=testKafkaStreamsBinderWithCustomEnvironmentCanStart", - "--spring.cloud.stream.bindings.input.destination=foo", - "--spring.cloud.stream.bindings.input.binder=kBind1", - "--spring.cloud.stream.binders.kBind1.type=kstream", - "--spring.cloud.stream.binders.kBind1.environment" + SimpleKafkaStreamsApplication.class).web(WebApplicationType.NONE).run( + "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.application-id" + + "=testKStreamBinderWithCustomEnvironmentCanStart", + "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.application-id" + + "=testKStreamBinderWithCustomEnvironmentCanStart-foo", + "--spring.cloud.stream.kafka.streams.bindings.input-3.consumer.application-id" + + "=testKStreamBinderWithCustomEnvironmentCanStart-foobar", + "--spring.cloud.stream.bindings.input-1.destination=foo", + "--spring.cloud.stream.bindings.input-1.binder=kstreamBinder", + "--spring.cloud.stream.binders.kstreamBinder.type=kstream", + "--spring.cloud.stream.binders.kstreamBinder.environment" + ".spring.cloud.stream.kafka.streams.binder.brokers" - + "=" + embeddedKafka.getEmbeddedKafka().getBrokersAsString()); + + "=" + embeddedKafka.getEmbeddedKafka().getBrokersAsString(), + "--spring.cloud.stream.bindings.input-2.destination=bar", + "--spring.cloud.stream.bindings.input-2.binder=ktableBinder", + "--spring.cloud.stream.binders.ktableBinder.type=ktable", + "--spring.cloud.stream.binders.ktableBinder.environment" + + ".spring.cloud.stream.kafka.streams.binder.brokers" + + "=" + embeddedKafka.getEmbeddedKafka().getBrokersAsString(), + "--spring.cloud.stream.bindings.input-3.destination=foobar", + "--spring.cloud.stream.bindings.input-3.binder=globalktableBinder", + "--spring.cloud.stream.binders.globalktableBinder.type=globalktable", + "--spring.cloud.stream.binders.globalktableBinder.environment" + + ".spring.cloud.stream.kafka.streams.binder.brokers" + + "=" + embeddedKafka.getEmbeddedKafka().getBrokersAsString()); applicationContext.close(); } @@ -56,10 +74,13 @@ public class KafkaStreamsBinderBootstrapTest { @Test public void testKafkaStreamsBinderWithStandardConfigurationCanStart() { ConfigurableApplicationContext applicationContext = new SpringApplicationBuilder( - SimpleApplication.class).web(WebApplicationType.NONE).run( - "--spring.cloud.stream.kafka.streams.default.consumer.application-id" - + "=testKafkaStreamsBinderWithStandardConfigurationCanStart", - "--spring.cloud.stream.bindings.input.destination=foo", + SimpleKafkaStreamsApplication.class).web(WebApplicationType.NONE).run( + "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.application-id" + + "=testKafkaStreamsBinderWithStandardConfigurationCanStart", + "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.application-id" + + "=testKafkaStreamsBinderWithStandardConfigurationCanStart-foo", + "--spring.cloud.stream.kafka.streams.bindings.input-3.consumer.application-id" + + "=testKafkaStreamsBinderWithStandardConfigurationCanStart-foobar", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getEmbeddedKafka().getBrokersAsString()); @@ -67,21 +88,44 @@ public class KafkaStreamsBinderBootstrapTest { } @SpringBootApplication - @EnableBinding(StreamSourceProcessor.class) - static class SimpleApplication { + @EnableBinding({SimpleKStreamBinding.class, SimpleKTableBinding.class, SimpleGlobalKTableBinding.class}) + static class SimpleKafkaStreamsApplication { @StreamListener - public void handle(@Input("input") KStream stream) { + public void handle(@Input("input-1") KStream stream) { + + } + + @StreamListener + public void handleX(@Input("input-2") KTable stream) { + + } + + @StreamListener + public void handleY(@Input("input-3") GlobalKTable stream) { } } - interface StreamSourceProcessor { + interface SimpleKStreamBinding { - @Input("input") + @Input("input-1") KStream inputStream(); } + interface SimpleKTableBinding { + + @Input("input-2") + KTable inputStream(); + + } + + interface SimpleGlobalKTableBinding { + + @Input("input-3") + GlobalKTable inputStream(); + + } }