Clean up resources on close. Allow cleanly terminating application process on context shutdown

Prevent IDE from using star imports

Update authors

Checkstyle fixes
This commit is contained in:
Tomek Szmytka
2022-05-27 22:05:01 +02:00
committed by Soby Chacko
parent 0d2697449e
commit 6640bf63e5
3 changed files with 105 additions and 64 deletions

View File

@@ -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

View File

@@ -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<BindingCreatedEvent> {
implements MeterBinder, ApplicationListener<BindingCreatedEvent>, 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);
}
}

View File

@@ -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<PartitionInfo> 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<PartitionInfo> 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<TopicPartition, OffsetAndMetadata> 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<PartitionInfo> 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<PartitionInfo> 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<PartitionInfo> 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<PartitionInfo> 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<PartitionInfo> 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<PartitionInfo> partitions1 = partitions(new Node(0, null, 0));
final List<PartitionInfo> 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<PartitionInfo> 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<PartitionInfo> partitions(Node... nodes) {
List<PartitionInfo> partitions = new ArrayList<>();
for (int i = 0; i < nodes.length; i++) {