diff --git a/.editorconfig b/.editorconfig index 0679d88a9..657e65a22 100644 --- a/.editorconfig +++ b/.editorconfig @@ -12,3 +12,7 @@ insert_final_newline = true [*.yml] indent_style = space indent_size = 2 + +[*.java] +ij_java_names_count_to_use_import_on_demand = 999 +ij_java_class_count_to_use_import_on_demand = 999 diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java index aa8fe95bb..c3cad944f 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java @@ -20,6 +20,7 @@ import java.util.HashMap; import java.util.LinkedList; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; @@ -62,9 +63,10 @@ import org.springframework.util.ObjectUtils; * @author Thomas Cheyney * @author Gary Russell * @author Lars Bilger + * @author Tomek Szmytka */ public class KafkaBinderMetrics - implements MeterBinder, ApplicationListener { + implements MeterBinder, ApplicationListener, AutoCloseable { private static final int DEFAULT_TIMEOUT = 5; @@ -258,4 +260,8 @@ public class KafkaBinderMetrics } } + @Override + public void close() throws Exception { + Optional.ofNullable(scheduler).ifPresent(ExecutorService::shutdown); + } } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java index af82a0b58..941580b44 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java @@ -51,6 +51,7 @@ import static org.mockito.Mockito.mock; * @author Thomas Cheyney * @author Soby Chacko * @author Lars Bilger + * @author Tomek Szmytka */ public class KafkaBinderMetricsTest { @@ -79,14 +80,15 @@ public class KafkaBinderMetricsTest { MockitoAnnotations.openMocks(this); org.mockito.BDDMockito.given(consumerFactory .createConsumer(ArgumentMatchers.any(), ArgumentMatchers.any())) - .willReturn(consumer); + .willReturn(consumer); org.mockito.BDDMockito.given(binder.getTopicsInUse()).willReturn(topicsInUse); metrics = new KafkaBinderMetrics(binder, kafkaBinderConfigurationProperties, - consumerFactory, null); + consumerFactory, null + ); org.mockito.BDDMockito - .given(consumer.endOffsets(ArgumentMatchers.anyCollection())) - .willReturn(java.util.Collections - .singletonMap(new TopicPartition(TEST_TOPIC, 0), 1000L)); + .given(consumer.endOffsets(ArgumentMatchers.anyCollection())) + .willReturn(java.util.Collections + .singletonMap(new TopicPartition(TEST_TOPIC, 0), 1000L)); } @Test @@ -95,35 +97,39 @@ public class KafkaBinderMetricsTest { TopicPartition topicPartition = new TopicPartition(TEST_TOPIC, 0); committed.put(topicPartition, new OffsetAndMetadata(500)); org.mockito.BDDMockito - .given(consumer.committed(ArgumentMatchers.anySet())) - .willReturn(committed); + .given(consumer.committed(ArgumentMatchers.anySet())) + .willReturn(committed); List partitions = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, - new TopicInformation("group1-metrics", partitions, false)); + topicsInUse.put( + TEST_TOPIC, + new TopicInformation("group1-metrics", partitions, false) + ); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)) - .willReturn(partitions); + .willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); assertThat(meterRegistry.get(KafkaBinderMetrics.OFFSET_LAG_METRIC_NAME) - .tag("group", "group1-metrics").tag("topic", TEST_TOPIC).gauge().value()) - .isEqualTo(500.0); + .tag("group", "group1-metrics").tag("topic", TEST_TOPIC).gauge().value()) + .isEqualTo(500.0); } @Test public void shouldNotContainAnyMetricsWhenUsingNoopGauge() { // Adding NoopGauge for the offset metric. meterRegistry.config().meterFilter( - MeterFilter.denyNameStartsWith("spring.cloud.stream.binder.kafka.offset")); + MeterFilter.denyNameStartsWith("spring.cloud.stream.binder.kafka.offset")); // Because we have NoopGauge for the offset metric in the meter registry, none of these expectations matter. org.mockito.BDDMockito - .given(consumer.committed(ArgumentMatchers.any(TopicPartition.class))) - .willReturn(new OffsetAndMetadata(500)); + .given(consumer.committed(ArgumentMatchers.any(TopicPartition.class))) + .willReturn(new OffsetAndMetadata(500)); List partitions = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, - new TopicInformation("group1-metrics", partitions, false)); + topicsInUse.put( + TEST_TOPIC, + new TopicInformation("group1-metrics", partitions, false) + ); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)) - .willReturn(partitions); + .willReturn(partitions); metrics.bindTo(meterRegistry); // Because of the NoopGauge, the meterRegistry should contain no metric. @@ -136,41 +142,47 @@ public class KafkaBinderMetricsTest { endOffsets.put(new TopicPartition(TEST_TOPIC, 0), 1000L); endOffsets.put(new TopicPartition(TEST_TOPIC, 1), 1000L); org.mockito.BDDMockito - .given(consumer.endOffsets(ArgumentMatchers.anyCollection())) - .willReturn(endOffsets); + .given(consumer.endOffsets(ArgumentMatchers.anyCollection())) + .willReturn(endOffsets); final Map committed = new HashMap<>(); TopicPartition topicPartition1 = new TopicPartition(TEST_TOPIC, 0); TopicPartition topicPartition2 = new TopicPartition(TEST_TOPIC, 1); committed.put(topicPartition1, new OffsetAndMetadata(500)); committed.put(topicPartition2, new OffsetAndMetadata(500)); org.mockito.BDDMockito - .given(consumer.committed(ArgumentMatchers.anySet())) - .willReturn(committed); - List partitions = partitions(new Node(0, null, 0), - new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, - new TopicInformation("group2-metrics", partitions, false)); + .given(consumer.committed(ArgumentMatchers.anySet())) + .willReturn(committed); + List partitions = partitions( + new Node(0, null, 0), + new Node(0, null, 0) + ); + topicsInUse.put( + TEST_TOPIC, + new TopicInformation("group2-metrics", partitions, false) + ); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)) - .willReturn(partitions); + .willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); assertThat(meterRegistry.get(KafkaBinderMetrics.OFFSET_LAG_METRIC_NAME) - .tag("group", "group2-metrics").tag("topic", TEST_TOPIC).gauge().value()) - .isEqualTo(1000.0); + .tag("group", "group2-metrics").tag("topic", TEST_TOPIC).gauge().value()) + .isEqualTo(1000.0); } @Test public void shouldIndicateFullLagForNotCommittedGroups() { List partitions = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, - new TopicInformation("group3-metrics", partitions, false)); + topicsInUse.put( + TEST_TOPIC, + new TopicInformation("group3-metrics", partitions, false) + ); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)) - .willReturn(partitions); + .willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); assertThat(meterRegistry.get(KafkaBinderMetrics.OFFSET_LAG_METRIC_NAME) - .tag("group", "group3-metrics").tag("topic", TEST_TOPIC).gauge().value()) - .isEqualTo(1000.0); + .tag("group", "group3-metrics").tag("topic", TEST_TOPIC).gauge().value()) + .isEqualTo(1000.0); } @Test @@ -184,76 +196,86 @@ public class KafkaBinderMetricsTest { @Test public void createsConsumerOnceWhenInvokedMultipleTimes() { final List partitions = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, - new TopicInformation("group4-metrics", partitions, false)); + topicsInUse.put( + TEST_TOPIC, + new TopicInformation("group4-metrics", partitions, false) + ); metrics.bindTo(meterRegistry); Gauge gauge = meterRegistry.get(KafkaBinderMetrics.OFFSET_LAG_METRIC_NAME) - .tag("group", "group4-metrics").tag("topic", TEST_TOPIC).gauge(); + .tag("group", "group4-metrics").tag("topic", TEST_TOPIC).gauge(); gauge.value(); assertThat(gauge.value()).isEqualTo(1000.0); org.mockito.Mockito.verify(this.consumerFactory) - .createConsumer(ArgumentMatchers.any(), ArgumentMatchers.any()); + .createConsumer(ArgumentMatchers.any(), ArgumentMatchers.any()); } @Test public void consumerCreationFailsFirstTime() { org.mockito.BDDMockito - .given(consumerFactory.createConsumer(ArgumentMatchers.any(), - ArgumentMatchers.any())) - .willThrow(KafkaException.class).willReturn(consumer); + .given(consumerFactory.createConsumer( + ArgumentMatchers.any(), + ArgumentMatchers.any() + )) + .willThrow(KafkaException.class).willReturn(consumer); final List partitions = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, - new TopicInformation("group5-metrics", partitions, false)); + topicsInUse.put( + TEST_TOPIC, + new TopicInformation("group5-metrics", partitions, false) + ); metrics.bindTo(meterRegistry); Gauge gauge = meterRegistry.get(KafkaBinderMetrics.OFFSET_LAG_METRIC_NAME) - .tag("group", "group5-metrics").tag("topic", TEST_TOPIC).gauge(); + .tag("group", "group5-metrics").tag("topic", TEST_TOPIC).gauge(); assertThat(gauge.value()).isEqualTo(0); assertThat(gauge.value()).isEqualTo(1000.0); org.mockito.Mockito.verify(this.consumerFactory, Mockito.times(2)) - .createConsumer(ArgumentMatchers.any(), ArgumentMatchers.any()); + .createConsumer(ArgumentMatchers.any(), ArgumentMatchers.any()); } @Test public void createOneConsumerPerGroup() { final List partitions1 = partitions(new Node(0, null, 0)); final List partitions2 = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, - new TopicInformation("group1-metrics", partitions1, false)); - topicsInUse.put("test2", - new TopicInformation("group2-metrics", partitions2, false)); + topicsInUse.put( + TEST_TOPIC, + new TopicInformation("group1-metrics", partitions1, false) + ); + topicsInUse.put( + "test2", + new TopicInformation("group2-metrics", partitions2, false) + ); metrics.bindTo(meterRegistry); KafkaConsumer consumer2 = mock(KafkaConsumer.class); org.mockito.BDDMockito - .given(consumerFactory.createConsumer( - ArgumentMatchers.eq("group2-metrics"), ArgumentMatchers.any())) - .willReturn(consumer2); + .given(consumerFactory.createConsumer( + ArgumentMatchers.eq("group2-metrics"), ArgumentMatchers.any())) + .willReturn(consumer2); org.mockito.BDDMockito - .given(consumer2.endOffsets(ArgumentMatchers.anyCollection())) - .willReturn(java.util.Collections - .singletonMap(new TopicPartition("test2", 0), 50L)); + .given(consumer2.endOffsets(ArgumentMatchers.anyCollection())) + .willReturn(java.util.Collections + .singletonMap(new TopicPartition("test2", 0), 50L)); Gauge gauge1 = meterRegistry.get(KafkaBinderMetrics.OFFSET_LAG_METRIC_NAME) - .tag("group", "group1-metrics").tag("topic", TEST_TOPIC).gauge(); + .tag("group", "group1-metrics").tag("topic", TEST_TOPIC).gauge(); Gauge gauge2 = meterRegistry.get(KafkaBinderMetrics.OFFSET_LAG_METRIC_NAME) - .tag("group", "group2-metrics").tag("topic", "test2").gauge(); + .tag("group", "group2-metrics").tag("topic", "test2").gauge(); gauge1.value(); gauge2.value(); assertThat(gauge1.value()).isEqualTo(1000.0); assertThat(gauge2.value()).isEqualTo(50.0); org.mockito.Mockito.verify(this.consumerFactory).createConsumer( - ArgumentMatchers.eq("group1-metrics"), ArgumentMatchers.any()); + ArgumentMatchers.eq("group1-metrics"), ArgumentMatchers.any()); org.mockito.Mockito.verify(this.consumerFactory).createConsumer( - ArgumentMatchers.eq("group2-metrics"), ArgumentMatchers.any()); + ArgumentMatchers.eq("group2-metrics"), ArgumentMatchers.any()); } @Test @@ -268,17 +290,26 @@ public class KafkaBinderMetricsTest { .given(consumer.beginningOffsets(ArgumentMatchers.anySet())) .willReturn(beginnings); List partitions = partitions(new Node(0, null, 0)); - topicsInUse.put(TEST_TOPIC, - new TopicInformation("group1-metrics", partitions, false)); + topicsInUse.put( + TEST_TOPIC, + new TopicInformation("group1-metrics", partitions, false) + ); org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)) - .willReturn(partitions); + .willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); assertThat(meterRegistry.get(KafkaBinderMetrics.OFFSET_LAG_METRIC_NAME) - .tag("group", "group1-metrics").tag("topic", TEST_TOPIC).gauge().value()) + .tag("group", "group1-metrics").tag("topic", TEST_TOPIC).gauge().value()) .isEqualTo(500.0); } + @Test + public void shouldShutdownSchedulerOnClose() throws Exception { + metrics.bindTo(meterRegistry); + metrics.close(); + assertThat(metrics.scheduler.isShutdown()).isTrue(); + } + private List partitions(Node... nodes) { List partitions = new ArrayList<>(); for (int i = 0; i < nodes.length; i++) {