diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java index 4134a9dff..bdb124924 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java @@ -20,10 +20,10 @@ import java.io.IOException; import io.micrometer.core.instrument.MeterRegistry; import io.micrometer.core.instrument.binder.MeterBinder; -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.context.PropertyPlaceholderAutoConfiguration; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; @@ -41,7 +41,6 @@ import org.springframework.context.annotation.Import; import org.springframework.kafka.security.jaas.KafkaJaasLoginModuleInitializer; import org.springframework.kafka.support.LoggingProducerListener; import org.springframework.kafka.support.ProducerListener; -import org.springframework.lang.Nullable; /** * @author David Turanski @@ -52,15 +51,14 @@ import org.springframework.lang.Nullable; * @author Henryk Konsek * @author Gary Russell * @author Oleg Zhurakousky + * @author Artem Bilan */ @Configuration @ConditionalOnMissingBean(Binder.class) -@Import({ PropertyPlaceholderAutoConfiguration.class, KafkaBinderHealthIndicatorConfiguration.class}) +@Import({ PropertyPlaceholderAutoConfiguration.class, KafkaBinderHealthIndicatorConfiguration.class }) @EnableConfigurationProperties({ KafkaExtendedBindingProperties.class }) public class KafkaBinderConfiguration { - protected static final Log logger = LogFactory.getLog(KafkaBinderConfiguration.class); - @Autowired private KafkaExtendedBindingProperties kafkaExtendedBindingProperties; @@ -82,7 +80,8 @@ public class KafkaBinderConfiguration { @Bean KafkaMessageChannelBinder kafkaMessageChannelBinder(KafkaBinderConfigurationProperties configurationProperties, - KafkaTopicProvisioner provisioningProvider) { + KafkaTopicProvisioner provisioningProvider) { + KafkaMessageChannelBinder kafkaMessageChannelBinder = new KafkaMessageChannelBinder( configurationProperties, provisioningProvider); kafkaMessageChannelBinder.setProducerListener(producerListener); @@ -96,22 +95,37 @@ public class KafkaBinderConfiguration { return new LoggingProducerListener(); } - @Bean - public MeterBinder kafkaBinderMetrics(KafkaMessageChannelBinder kafkaMessageChannelBinder, - KafkaBinderConfigurationProperties configurationProperties, - @Nullable MeterRegistry meterRegistry) { - return new KafkaBinderMetrics(kafkaMessageChannelBinder, configurationProperties, null, meterRegistry); - } - @Bean public KafkaJaasLoginModuleInitializer jaasInitializer() throws IOException { return new KafkaJaasLoginModuleInitializer(); } + /** + * A conditional configuration for the {@link KafkaBinderMetrics} bean when the + * {@link MeterRegistry} class is in classpath, as well as a {@link MeterRegistry} bean is + * present in the application context. + */ + @Configuration + @ConditionalOnClass(MeterRegistry.class) + @ConditionalOnBean(MeterRegistry.class) + protected class KafkaBinderMetricsConfiguration { + + @Bean + @ConditionalOnMissingBean(KafkaBinderMetrics.class) + public MeterBinder kafkaBinderMetrics(KafkaMessageChannelBinder kafkaMessageChannelBinder, + KafkaBinderConfigurationProperties configurationProperties, + MeterRegistry meterRegistry) { + + return new KafkaBinderMetrics(kafkaMessageChannelBinder, configurationProperties, null, meterRegistry); + } + + } + public static class JaasConfigurationProperties { private JaasLoginModuleConfiguration kafka; private JaasLoginModuleConfiguration zookeeper; } + } diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java index 7a30df22b..c54eadbac 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java @@ -17,6 +17,7 @@ package org.springframework.cloud.stream.binder.kafka.integration; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.binder.MeterBinder; import org.junit.AfterClass; import org.junit.BeforeClass; @@ -26,7 +27,9 @@ import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.test.context.FilteredClassLoader; import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.context.runner.ApplicationContextRunner; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.cloud.stream.messaging.Sink; @@ -44,8 +47,7 @@ import static org.assertj.core.api.Assertions.assertThat; * @since 2.0 */ @RunWith(SpringRunner.class) -@SpringBootTest( - webEnvironment = SpringBootTest.WebEnvironment.NONE, +@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE, properties = "spring.cloud.stream.bindings.input.group=" + KafkaBinderActuatorTests.TEST_CONSUMER_GROUP) public class KafkaBinderActuatorTests { @@ -83,6 +85,17 @@ public class KafkaBinderActuatorTests { .timeGauge().value()).isGreaterThan(0); } + @Test + public void testKafkaBinderMetricsWhenNoMicrometer() { + new ApplicationContextRunner() + .withUserConfiguration(KafkaMetricsTestConfig.class) + .withClassLoader(new FilteredClassLoader("io.micrometer.core")) + .run(context -> { + assertThat(context.getBeanNamesForType(MeterRegistry.class)).isEmpty(); + assertThat(context.getBeanNamesForType(MeterBinder.class)).isEmpty(); + }); + } + @EnableBinding(Sink.class) @EnableAutoConfiguration public static class KafkaMetricsTestConfig {