diff --git a/pom.xml b/pom.xml
index 4601c93c7..7384eb77b 100644
--- a/pom.xml
+++ b/pom.xml
@@ -149,10 +149,6 @@
org.apache.maven.plugins
maven-checkstyle-plugin
-
- io.spring.javaformat
- spring-javaformat-maven-plugin
-
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 c83e49cdf..3aba94a6d 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
@@ -40,6 +40,8 @@ import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfi
import org.springframework.cloud.stream.binder.kafka.properties.KafkaExtendedBindingProperties;
import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner;
import org.springframework.cloud.stream.config.ListenerContainerCustomizer;
+import org.springframework.context.ApplicationContext;
+import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
@@ -144,17 +146,13 @@ public class KafkaBinderConfiguration {
return 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)
+ @ConditionalOnMissingBean(value = KafkaBinderMetrics.class, name = "outerContext")
+ @ConditionalOnClass(name = "io.micrometer.core.instrument.MeterRegistry")
protected class KafkaBinderMetricsConfiguration {
@Bean
+ @ConditionalOnBean(MeterRegistry.class)
@ConditionalOnMissingBean(KafkaBinderMetrics.class)
public MeterBinder kafkaBinderMetrics(
KafkaMessageChannelBinder kafkaMessageChannelBinder,
@@ -164,7 +162,25 @@ public class KafkaBinderConfiguration {
return new KafkaBinderMetrics(kafkaMessageChannelBinder,
configurationProperties, null, meterRegistry);
}
+ }
+ @Configuration
+ @ConditionalOnBean(name = "outerContext")
+ @ConditionalOnMissingBean(KafkaBinderMetrics.class)
+ @ConditionalOnClass(name = "io.micrometer.core.instrument.MeterRegistry")
+ protected class KafkaBinderMetricsConfigurationWithMultiBinder {
+
+ @Bean
+ public MeterBinder kafkaBinderMetrics(
+ KafkaMessageChannelBinder kafkaMessageChannelBinder,
+ KafkaBinderConfigurationProperties configurationProperties,
+ ConfigurableApplicationContext context) {
+
+ MeterRegistry meterRegistry = context.getBean("outerContext", ApplicationContext.class)
+ .getBean(MeterRegistry.class);
+ return new KafkaBinderMetrics(kafkaMessageChannelBinder,
+ configurationProperties, null, meterRegistry);
+ }
}
/**
@@ -178,5 +194,4 @@ public class KafkaBinderConfiguration {
private JaasLoginModuleConfiguration zookeeper;
}
-
}
diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/bootstrap/MultiBinderMeterRegistryTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/bootstrap/MultiBinderMeterRegistryTest.java
new file mode 100644
index 000000000..8d4ad35f9
--- /dev/null
+++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/bootstrap/MultiBinderMeterRegistryTest.java
@@ -0,0 +1,70 @@
+/*
+ * Copyright 2019-2019 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.stream.binder.kafka.bootstrap;
+
+import io.micrometer.core.instrument.MeterRegistry;
+import org.junit.ClassRule;
+import org.junit.Test;
+
+import org.springframework.boot.WebApplicationType;
+import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.boot.builder.SpringApplicationBuilder;
+import org.springframework.cloud.stream.annotation.EnableBinding;
+import org.springframework.cloud.stream.messaging.Sink;
+import org.springframework.context.ConfigurableApplicationContext;
+import org.springframework.kafka.test.rule.KafkaEmbedded;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * @author Soby Chacko
+ */
+public class MultiBinderMeterRegistryTest {
+
+ @ClassRule
+ public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, 10);
+
+ @Test
+ public void testMetricsWorkWithMultiBinders() {
+ ConfigurableApplicationContext applicationContext = new SpringApplicationBuilder(
+ SimpleApplication.class).web(WebApplicationType.NONE).run(
+ "--spring.cloud.stream.bindings.input.destination=foo",
+ "--spring.cloud.stream.bindings.input.binder=inbound",
+ "--spring.cloud.stream.bindings.input.group=testGroupabc",
+ "--spring.cloud.stream.binders.inbound.type=kafka",
+ "--spring.cloud.stream.binders.inbound.environment"
+ + ".spring.cloud.stream.kafka.binder.brokers" + "="
+ + embeddedKafka.getBrokersAsString());
+
+ final MeterRegistry meterRegistry = applicationContext.getBean(MeterRegistry.class);
+
+ assertThat(meterRegistry).isNotNull();
+
+ assertThat(meterRegistry.get("spring.cloud.stream.binder.kafka.offset")
+ .tag("group", "testGroupabc")
+ .tag("topic", "foo").gauge().value()).isNotNull();
+
+ applicationContext.close();
+ }
+
+ @SpringBootApplication
+ @EnableBinding(Sink.class)
+ static class SimpleApplication {
+
+ }
+
+}