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 81a695842..7eae7953d 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,7 +21,6 @@ import java.util.Map; import org.springframework.beans.factory.ObjectProvider; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.config.BeanFactoryPostProcessor; -import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; @@ -44,26 +43,25 @@ import org.springframework.context.annotation.Import; @Import({ KafkaAutoConfiguration.class, MultiBinderPropertiesConfiguration.class, KafkaStreamsBinderHealthIndicatorConfiguration.class}) -@AutoConfigureAfter(KafkaStreamsBinderSupportAutoConfiguration.class) public class GlobalKTableBinderConfiguration { @Bean - public KafkaTopicProvisioner globalKTableProvisioningProvider( - @Qualifier("binderConfigurationProperties") KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, + public KafkaTopicProvisioner provisioningProvider( + KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, KafkaProperties kafkaProperties, ObjectProvider adminClientConfigCustomizer) { return new KafkaTopicProvisioner(binderConfigurationProperties, kafkaProperties, adminClientConfigCustomizer.getIfUnique()); } - @Bean("globalktable") + @Bean public GlobalKTableBinder GlobalKTableBinder( KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, - KafkaTopicProvisioner globalKTableProvisioningProvider, + KafkaTopicProvisioner kafkaTopicProvisioner, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, @Qualifier("streamConfigGlobalProperties") Map streamConfigGlobalProperties) { GlobalKTableBinder globalKTableBinder = new GlobalKTableBinder(binderConfigurationProperties, - globalKTableProvisioningProvider, kafkaStreamsBindingInformationCatalogue); + kafkaTopicProvisioner, kafkaStreamsBindingInformationCatalogue); globalKTableBinder.setKafkaStreamsExtendedBindingProperties( kafkaStreamsExtendedBindingProperties); return globalKTableBinder; @@ -71,7 +69,7 @@ public class GlobalKTableBinderConfiguration { @Bean @ConditionalOnBean(name = "outerContext") - public static BeanFactoryPostProcessor globalKTableOuterContextBeanFactoryPostProcessor() { + public static BeanFactoryPostProcessor outerContextBeanFactoryPostProcessor() { return beanFactory -> { // It is safe to call getBean("outerContext") here, because this bean is 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 a7ee25cc1..9450bc392 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 @@ -17,9 +17,7 @@ package org.springframework.cloud.stream.binder.kafka.streams; import org.springframework.beans.factory.ObjectProvider; -import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.config.BeanFactoryPostProcessor; -import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; @@ -43,27 +41,26 @@ import org.springframework.context.annotation.Import; @Import({ KafkaAutoConfiguration.class, MultiBinderPropertiesConfiguration.class, KafkaStreamsBinderHealthIndicatorConfiguration.class}) -@AutoConfigureAfter(KafkaStreamsBinderSupportAutoConfiguration.class) public class KStreamBinderConfiguration { @Bean - public KafkaTopicProvisioner kstreamProvisioningProvider( - @Qualifier("binderConfigurationProperties") KafkaStreamsBinderConfigurationProperties kafkaStreamsBinderConfigurationProperties, + public KafkaTopicProvisioner provisioningProvider( + KafkaStreamsBinderConfigurationProperties kafkaStreamsBinderConfigurationProperties, KafkaProperties kafkaProperties, ObjectProvider adminClientConfigCustomizer) { return new KafkaTopicProvisioner(kafkaStreamsBinderConfigurationProperties, kafkaProperties, adminClientConfigCustomizer.getIfUnique()); } - @Bean("kstream") + @Bean public KStreamBinder kStreamBinder( KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, - KafkaTopicProvisioner kstreamProvisioningProvider, + KafkaTopicProvisioner kafkaTopicProvisioner, KafkaStreamsMessageConversionDelegate KafkaStreamsMessageConversionDelegate, KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue, KeyValueSerdeResolver keyValueSerdeResolver, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties) { KStreamBinder kStreamBinder = new KStreamBinder(binderConfigurationProperties, - kstreamProvisioningProvider, KafkaStreamsMessageConversionDelegate, + kafkaTopicProvisioner, KafkaStreamsMessageConversionDelegate, KafkaStreamsBindingInformationCatalogue, keyValueSerdeResolver); kStreamBinder.setKafkaStreamsExtendedBindingProperties( kafkaStreamsExtendedBindingProperties); @@ -72,7 +69,7 @@ public class KStreamBinderConfiguration { @Bean @ConditionalOnBean(name = "outerContext") - public static BeanFactoryPostProcessor kstreamOuterContextBeanFactoryPostProcessor() { + public static BeanFactoryPostProcessor outerContextBeanFactoryPostProcessor() { return beanFactory -> { // It is safe to call getBean("outerContext") here, because this bean is 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 4da804e15..825e4a2ae 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 @@ -46,28 +46,28 @@ import org.springframework.context.annotation.Import; public class KTableBinderConfiguration { @Bean - public KafkaTopicProvisioner ktableProvisioningProvider( - @Qualifier("binderConfigurationProperties") KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, + public KafkaTopicProvisioner provisioningProvider( + KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, KafkaProperties kafkaProperties, ObjectProvider adminClientConfigCustomizer) { return new KafkaTopicProvisioner(binderConfigurationProperties, kafkaProperties, adminClientConfigCustomizer.getIfUnique()); } - @Bean("ktable") + @Bean public KTableBinder kTableBinder( KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, - KafkaTopicProvisioner ktableProvisioningProvider, + KafkaTopicProvisioner kafkaTopicProvisioner, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, @Qualifier("streamConfigGlobalProperties") Map streamConfigGlobalProperties) { KTableBinder kTableBinder = new KTableBinder(binderConfigurationProperties, - ktableProvisioningProvider, kafkaStreamsBindingInformationCatalogue); + kafkaTopicProvisioner, kafkaStreamsBindingInformationCatalogue); kTableBinder.setKafkaStreamsExtendedBindingProperties(kafkaStreamsExtendedBindingProperties); return kTableBinder; } @Bean @ConditionalOnBean(name = "outerContext") - public static BeanFactoryPostProcessor ktableOuterContextBeanFactoryPostProcessor() { + public static BeanFactoryPostProcessor outerContextBeanFactoryPostProcessor() { return beanFactory -> { // It is safe to call getBean("outerContext") here, because this bean is diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/resources/spring.binders b/spring-cloud-stream-binder-kafka-streams/src/main/resources/META-INF/spring.binders similarity index 100% rename from spring-cloud-stream-binder-kafka-streams/src/main/resources/spring.binders rename to spring-cloud-stream-binder-kafka-streams/src/main/resources/META-INF/spring.binders diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/resources/META-INF/spring.factories b/spring-cloud-stream-binder-kafka-streams/src/main/resources/META-INF/spring.factories index 8f59c3d66..ca6492f05 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/resources/META-INF/spring.factories +++ b/spring-cloud-stream-binder-kafka-streams/src/main/resources/META-INF/spring.factories @@ -1,7 +1,4 @@ -org.springframework.boot.autoconfigure.EnableAutoConfiguration:\ -org.springframework.cloud.stream.binder.kafka.streams.KafkaStreamsBinderSupportAutoConfiguration,\ -org.springframework.cloud.stream.binder.kafka.streams.function.KafkaStreamsFunctionAutoConfiguration,\ -org.springframework.cloud.stream.binder.kafka.streams.endpoint.KafkaStreamsTopologyEndpointAutoConfiguration,\ -org.springframework.cloud.stream.binder.kafka.streams.KStreamBinderConfiguration,\ -org.springframework.cloud.stream.binder.kafka.streams.KTableBinderConfiguration,\ -org.springframework.cloud.stream.binder.kafka.streams.GlobalKTableBinderConfiguration +org.springframework.boot.autoconfigure.EnableAutoConfiguration=\ + org.springframework.cloud.stream.binder.kafka.streams.KafkaStreamsBinderSupportAutoConfiguration,\ + org.springframework.cloud.stream.binder.kafka.streams.function.KafkaStreamsFunctionAutoConfiguration,\ + org.springframework.cloud.stream.binder.kafka.streams.endpoint.KafkaStreamsTopologyEndpointAutoConfiguration 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 1d8a28bed..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 @@ -50,21 +50,21 @@ public class KafkaStreamsBinderBootstrapTest { "--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=kstream", - "--spring.cloud.stream.binders.kstream.type=kstream", - "--spring.cloud.stream.binders.kstream.environment" + "--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(), "--spring.cloud.stream.bindings.input-2.destination=bar", - "--spring.cloud.stream.bindings.input-2.binder=ktable", - "--spring.cloud.stream.binders.ktable.type=ktable", - "--spring.cloud.stream.binders.ktable.environment" + "--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=globalktable", - "--spring.cloud.stream.binders.globalktable.type=globalktable", - "--spring.cloud.stream.binders.globalktable.environment" + "--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()); diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java index e0dd96e48..54dcf9d58 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java @@ -34,6 +34,7 @@ import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; +import org.springframework.boot.actuate.health.CompositeHealthContributor; import org.springframework.boot.actuate.health.Health; import org.springframework.boot.actuate.health.Status; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; @@ -187,13 +188,9 @@ public class KafkaStreamsBinderHealthIndicatorTests { private static void checkHealth(ConfigurableApplicationContext context, Status expected) throws InterruptedException { -// CompositeHealthContributor healthIndicator = context -// .getBean("bindersHealthContributor", CompositeHealthContributor.class); - - KafkaStreamsBinderHealthIndicator kafkaStreamsBinderHealthIndicator = context - .getBean("kafkaStreamsBinderHealthIndicator", KafkaStreamsBinderHealthIndicator.class); - - //KafkaStreamsBinderHealthIndicator kafkaStreamsBinderHealthIndicator = (KafkaStreamsBinderHealthIndicator) healthIndicator.getContributor("kstream"); + CompositeHealthContributor healthIndicator = context + .getBean("bindersHealthContributor", CompositeHealthContributor.class); + KafkaStreamsBinderHealthIndicator kafkaStreamsBinderHealthIndicator = (KafkaStreamsBinderHealthIndicator) healthIndicator.getContributor("kstream"); Health health = kafkaStreamsBinderHealthIndicator.health(); while (waitFor(health.getStatus(), health.getDetails())) { TimeUnit.SECONDS.sleep(2);