From a7299df63f495af3a798637551c2179c947af9cf Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 10 Mar 2021 19:11:15 -0500 Subject: [PATCH] Cleanup Kafka Streams metrics support StreamsListener (for Micrometer) is now directly available in Spring Kafka starting from 2.5.3 (as KafkaStreamsMicrometerListener). Removing the temporory interface added in the binder. Addressing PR review comments. Modifying tests to verify. Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/1040 --- ...StreamsBinderSupportAutoConfiguration.java | 64 ++----------------- .../streams/StreamsBuilderFactoryManager.java | 13 ++-- .../binder/kafka/streams/StreamsListener.java | 46 ------------- ...kaStreamsBinderWordCountFunctionTests.java | 8 ++- 4 files changed, 18 insertions(+), 113 deletions(-) delete mode 100644 spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsListener.java diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java index 9551bd964..9d292848a 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java @@ -17,7 +17,6 @@ package org.springframework.cloud.stream.binder.kafka.streams; import java.lang.reflect.Constructor; -import java.util.ArrayList; import java.util.Collection; import java.util.HashMap; import java.util.List; @@ -25,12 +24,8 @@ import java.util.Map; import java.util.Properties; import java.util.stream.Collectors; -import io.micrometer.core.instrument.ImmutableTag; import io.micrometer.core.instrument.MeterRegistry; -import io.micrometer.core.instrument.Tag; -import io.micrometer.core.instrument.binder.kafka.KafkaStreamsMetrics; import org.apache.kafka.common.serialization.Serdes; -import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.errors.LogAndContinueExceptionHandler; import org.apache.kafka.streams.errors.LogAndFailExceptionHandler; @@ -76,6 +71,7 @@ import org.springframework.integration.support.utils.IntegrationUtils; import org.springframework.kafka.config.KafkaStreamsConfiguration; import org.springframework.kafka.config.StreamsBuilderFactoryBeanCustomizer; import org.springframework.kafka.core.CleanupConfig; +import org.springframework.kafka.streams.KafkaStreamsMicrometerListener; import org.springframework.kafka.streams.RecoveringDeserializationExceptionHandler; import org.springframework.lang.Nullable; import org.springframework.messaging.converter.CompositeMessageConverter; @@ -402,7 +398,7 @@ public class KafkaStreamsBinderSupportAutoConfiguration { KafkaStreamsBindingInformationCatalogue catalogue, KafkaStreamsRegistry kafkaStreamsRegistry, @Nullable KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics, - @Nullable StreamsListener listener) { + @Nullable KafkaStreamsMicrometerListener listener) { return new StreamsBuilderFactoryManager(catalogue, kafkaStreamsRegistry, kafkaStreamsBinderMetrics, listener); } @@ -448,33 +444,10 @@ public class KafkaStreamsBinderSupportAutoConfiguration { @Bean @ConditionalOnMissingBean(name = "binderStreamsListener") - public StreamsListener binderStreamsListener(MeterRegistry meterRegistry) { - return new StreamsListener() { - - private final Map metrics = new HashMap<>(); - - @Override - public synchronized void streamsAdded(String id, KafkaStreams kafkaStreams) { - if (!this.metrics.containsKey(id)) { - List streamsTags = new ArrayList<>(); - streamsTags.add(new ImmutableTag("spring.id", id)); - this.metrics.put(id, new KafkaStreamsMetrics(kafkaStreams, streamsTags)); - this.metrics.get(id).bindTo(meterRegistry); - } - } - - @Override - public synchronized void streamsRemoved(String id, KafkaStreams streams) { - KafkaStreamsMetrics removed = this.metrics.remove(id); - if (removed != null) { - removed.close(); - } - } - - }; + public KafkaStreamsMicrometerListener binderStreamsListener(MeterRegistry meterRegistry) { + return new KafkaStreamsMicrometerListener(meterRegistry); } } - } @Configuration @@ -498,34 +471,9 @@ public class KafkaStreamsBinderSupportAutoConfiguration { @Bean @ConditionalOnMissingBean(name = "binderStreamsListener") - public StreamsListener binderStreamsListener(ConfigurableApplicationContext context) { - MeterRegistry meterRegistry = context.getBean("outerContext", ApplicationContext.class) - .getBean(MeterRegistry.class); - return new StreamsListener() { - - private final Map metrics = new HashMap<>(); - - @Override - public synchronized void streamsAdded(String id, KafkaStreams kafkaStreams) { - if (!this.metrics.containsKey(id)) { - List streamsTags = new ArrayList<>(); - streamsTags.add(new ImmutableTag("spring.id", id)); - this.metrics.put(id, new KafkaStreamsMetrics(kafkaStreams, streamsTags)); - this.metrics.get(id).bindTo(meterRegistry); - } - } - - @Override - public synchronized void streamsRemoved(String id, KafkaStreams streams) { - KafkaStreamsMetrics removed = this.metrics.remove(id); - if (removed != null) { - removed.close(); - } - } - - }; + public KafkaStreamsMicrometerListener binderStreamsListener(MeterRegistry meterRegistry) { + return new KafkaStreamsMicrometerListener(meterRegistry); } } - } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java index bbea2c38e..7a15eae7d 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsBuilderFactoryManager.java @@ -23,6 +23,7 @@ import org.springframework.context.SmartLifecycle; import org.springframework.kafka.KafkaException; import org.springframework.kafka.config.StreamsBuilderFactoryBean; import org.springframework.kafka.core.ProducerFactory; +import org.springframework.kafka.streams.KafkaStreamsMicrometerListener; /** * Iterate through all {@link StreamsBuilderFactoryBean} in the application context and @@ -45,14 +46,14 @@ class StreamsBuilderFactoryManager implements SmartLifecycle { private final KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics; - private final StreamsListener listener; + private final KafkaStreamsMicrometerListener listener; private volatile boolean running; StreamsBuilderFactoryManager(KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, KafkaStreamsRegistry kafkaStreamsRegistry, KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics, - StreamsListener listener) { + KafkaStreamsMicrometerListener listener) { this.kafkaStreamsBindingInformationCatalogue = kafkaStreamsBindingInformationCatalogue; this.kafkaStreamsRegistry = kafkaStreamsRegistry; this.kafkaStreamsBinderMetrics = kafkaStreamsBinderMetrics; @@ -80,11 +81,9 @@ class StreamsBuilderFactoryManager implements SmartLifecycle { .getStreamsBuilderFactoryBeans(); int n = 0; for (StreamsBuilderFactoryBean streamsBuilderFactoryBean : streamsBuilderFactoryBeans) { + streamsBuilderFactoryBean.addListener(this.listener); streamsBuilderFactoryBean.start(); this.kafkaStreamsRegistry.registerKafkaStreams(streamsBuilderFactoryBean); - if (this.listener != null) { - this.listener.streamsAdded("streams." + n++, streamsBuilderFactoryBean.getKafkaStreams()); - } } if (this.kafkaStreamsBinderMetrics != null) { this.kafkaStreamsBinderMetrics.addMetrics(streamsBuilderFactoryBeans); @@ -105,10 +104,8 @@ class StreamsBuilderFactoryManager implements SmartLifecycle { .getStreamsBuilderFactoryBeans(); int n = 0; for (StreamsBuilderFactoryBean streamsBuilderFactoryBean : streamsBuilderFactoryBeans) { + streamsBuilderFactoryBean.removeListener(this.listener); streamsBuilderFactoryBean.stop(); - if (this.listener != null) { - this.listener.streamsRemoved("streams." + n++, streamsBuilderFactoryBean.getKafkaStreams()); - } } for (ProducerFactory dlqProducerFactory : this.kafkaStreamsBindingInformationCatalogue.getDlqProducerFactories()) { ((DisposableBean) dlqProducerFactory).destroy(); diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsListener.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsListener.java deleted file mode 100644 index 807c74a4c..000000000 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/StreamsListener.java +++ /dev/null @@ -1,46 +0,0 @@ -/* - * Copyright 2020-2020 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 - * - * https://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.streams; - -import org.apache.kafka.streams.KafkaStreams; - -/** - * Temporary workaround until SK 2.5.3 is available. - * - * @author Gary Russell - * @since 3.0.6 - * - */ -interface StreamsListener { - - /** - * A new {@link KafkaStreams} was created. - * @param id the streams id (factory bean name). - * @param streams the streams; - */ - default void streamsAdded(String id, KafkaStreams streams) { - } - - /** - * An existing {@link KafkaStreams} was removed. - * @param id the streams id (factory bean name). - * @param streams the streams; - */ - default void streamsRemoved(String id, KafkaStreams streams) { - } - -} diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java index b5204b96f..3de45fdc8 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java @@ -109,9 +109,15 @@ public class KafkaStreamsBinderWordCountFunctionTests { "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) { receiveAndValidate("words", "counts"); + final MeterRegistry meterRegistry = context.getBean(MeterRegistry.class); Thread.sleep(100); - assertThat(meterRegistry.getMeters().size() > 1).isTrue(); + + assertThat(meterRegistry.getMeters().stream().anyMatch(m -> m.getId().getName().equals("kafka.stream.thread.poll.records.max"))).isTrue(); + assertThat(meterRegistry.getMeters().stream().anyMatch(m -> m.getId().getName().equals("kafka.consumer.network.io.total"))).isTrue(); + assertThat(meterRegistry.getMeters().stream().anyMatch(m -> m.getId().getName().equals("kafka.producer.record.send.total"))).isTrue(); + assertThat(meterRegistry.getMeters().stream().anyMatch(m -> m.getId().getName().equals("kafka.admin.client.network.io.total"))).isTrue(); + Assert.isTrue(LATCH.await(5, TimeUnit.SECONDS), "Failed to call customizers"); //Testing topology endpoint final KafkaStreamsRegistry kafkaStreamsRegistry = context.getBean(KafkaStreamsRegistry.class);