gh-1059 : Added health indicator for kafka messages listener containers created via kafka binder
renamed variables, return unknown on empty listener list added tests, fixed PR comments checkstyle fix
This commit is contained in:
committed by
Soby Chacko
parent
4fb5037fd7
commit
cb42e80dac
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2016-2018 the original author or authors.
|
||||
* Copyright 2016-2021 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.
|
||||
@@ -17,6 +17,8 @@
|
||||
package org.springframework.cloud.stream.binder.kafka;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -34,7 +36,10 @@ import org.apache.kafka.common.PartitionInfo;
|
||||
import org.springframework.beans.factory.DisposableBean;
|
||||
import org.springframework.boot.actuate.health.Health;
|
||||
import org.springframework.boot.actuate.health.HealthIndicator;
|
||||
import org.springframework.boot.actuate.health.Status;
|
||||
import org.springframework.boot.actuate.health.StatusAggregator;
|
||||
import org.springframework.kafka.core.ConsumerFactory;
|
||||
import org.springframework.kafka.listener.AbstractMessageListenerContainer;
|
||||
import org.springframework.scheduling.concurrent.CustomizableThreadFactory;
|
||||
|
||||
/**
|
||||
@@ -48,6 +53,7 @@ import org.springframework.scheduling.concurrent.CustomizableThreadFactory;
|
||||
* @author Soby Chacko
|
||||
* @author Vladislav Fefelov
|
||||
* @author Chukwubuikem Ume-Ugwa
|
||||
* @author Taras Danylchuk
|
||||
*/
|
||||
public class KafkaBinderHealthIndicator implements HealthIndicator, DisposableBean {
|
||||
|
||||
@@ -86,7 +92,22 @@ public class KafkaBinderHealthIndicator implements HealthIndicator, DisposableBe
|
||||
|
||||
@Override
|
||||
public Health health() {
|
||||
Future<Health> future = executor.submit(this::buildHealthStatus);
|
||||
Health topicsHealth = safelyBuildTopicsHealth();
|
||||
Health listenerContainersHealth = buildListenerContainersHealth();
|
||||
return merge(topicsHealth, listenerContainersHealth);
|
||||
}
|
||||
|
||||
private Health merge(Health topicsHealth, Health listenerContainersHealth) {
|
||||
Status aggregatedStatus = StatusAggregator.getDefault()
|
||||
.getAggregateStatus(topicsHealth.getStatus(), listenerContainersHealth.getStatus());
|
||||
Map<String, Object> aggregatedDetails = new HashMap<>();
|
||||
aggregatedDetails.putAll(topicsHealth.getDetails());
|
||||
aggregatedDetails.putAll(listenerContainersHealth.getDetails());
|
||||
return Health.status(aggregatedStatus).withDetails(aggregatedDetails).build();
|
||||
}
|
||||
|
||||
private Health safelyBuildTopicsHealth() {
|
||||
Future<Health> future = executor.submit(this::buildTopicsHealth);
|
||||
try {
|
||||
return future.get(this.timeout, TimeUnit.SECONDS);
|
||||
}
|
||||
@@ -112,10 +133,11 @@ public class KafkaBinderHealthIndicator implements HealthIndicator, DisposableBe
|
||||
}
|
||||
}
|
||||
|
||||
private Health buildHealthStatus() {
|
||||
private Health buildTopicsHealth() {
|
||||
try {
|
||||
initMetadataConsumer();
|
||||
Set<String> downMessages = new HashSet<>();
|
||||
Set<String> checkedTopics = new HashSet<>();
|
||||
final Map<String, KafkaMessageChannelBinder.TopicInformation> topicsInUse = KafkaBinderHealthIndicator.this.binder
|
||||
.getTopicsInUse();
|
||||
if (topicsInUse.isEmpty()) {
|
||||
@@ -148,11 +170,12 @@ public class KafkaBinderHealthIndicator implements HealthIndicator, DisposableBe
|
||||
downMessages.add(partitionInfo.toString());
|
||||
}
|
||||
}
|
||||
checkedTopics.add(topic);
|
||||
}
|
||||
}
|
||||
}
|
||||
if (downMessages.isEmpty()) {
|
||||
return Health.up().build();
|
||||
return Health.up().withDetail("topicsInUse", checkedTopics).build();
|
||||
}
|
||||
else {
|
||||
return Health.down()
|
||||
@@ -166,6 +189,33 @@ public class KafkaBinderHealthIndicator implements HealthIndicator, DisposableBe
|
||||
}
|
||||
}
|
||||
|
||||
private Health buildListenerContainersHealth() {
|
||||
List<AbstractMessageListenerContainer<?, ?>> listenerContainers = binder.getKafkaMessageListenerContainers();
|
||||
if (listenerContainers.isEmpty()) {
|
||||
return Health.unknown().build();
|
||||
}
|
||||
|
||||
Status status = Status.UP;
|
||||
List<Map<String, Object>> containersDetails = new ArrayList<>();
|
||||
|
||||
for (AbstractMessageListenerContainer<?, ?> container : listenerContainers) {
|
||||
Map<String, Object> containerDetails = new HashMap<>();
|
||||
boolean isRunning = container.isRunning();
|
||||
if (!isRunning) {
|
||||
status = Status.DOWN;
|
||||
}
|
||||
containerDetails.put("isRunning", isRunning);
|
||||
containerDetails.put("isPaused", container.isContainerPaused());
|
||||
containerDetails.put("listenerId", container.getListenerId());
|
||||
containerDetails.put("groupId", container.getGroupId());
|
||||
|
||||
containersDetails.add(containerDetails);
|
||||
}
|
||||
return Health.status(status)
|
||||
.withDetail("listenerContainers", containersDetails)
|
||||
.build();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void destroy() throws Exception {
|
||||
executor.shutdown();
|
||||
|
||||
@@ -153,6 +153,7 @@ import org.springframework.util.concurrent.ListenableFutureCallback;
|
||||
* @author Henryk Konsek
|
||||
* @author Doug Saus
|
||||
* @author Lukasz Kaminski
|
||||
* @author Taras Danylchuk
|
||||
*/
|
||||
public class KafkaMessageChannelBinder extends
|
||||
// @checkstyle:off
|
||||
@@ -233,6 +234,8 @@ public class KafkaMessageChannelBinder extends
|
||||
|
||||
private ConsumerConfigCustomizer consumerConfigCustomizer;
|
||||
|
||||
private final List<AbstractMessageListenerContainer<?, ?>> kafkaMessageListenerContainers = new ArrayList<>();
|
||||
|
||||
public KafkaMessageChannelBinder(
|
||||
KafkaBinderConfigurationProperties configurationProperties,
|
||||
KafkaTopicProvisioner provisioningProvider) {
|
||||
@@ -681,6 +684,8 @@ public class KafkaMessageChannelBinder extends
|
||||
}
|
||||
|
||||
};
|
||||
|
||||
this.kafkaMessageListenerContainers.add(messageListenerContainer);
|
||||
messageListenerContainer.setConcurrency(concurrency);
|
||||
// these won't be needed if the container is made a bean
|
||||
AbstractApplicationContext applicationContext = getApplicationContext();
|
||||
@@ -1479,6 +1484,10 @@ public class KafkaMessageChannelBinder extends
|
||||
this.producerConfigCustomizer = producerConfigCustomizer;
|
||||
}
|
||||
|
||||
List<AbstractMessageListenerContainer<?, ?>> getKafkaMessageListenerContainers() {
|
||||
return Collections.unmodifiableList(kafkaMessageListenerContainers);
|
||||
}
|
||||
|
||||
private final class ProducerConfigurationMessageHandler
|
||||
extends KafkaProducerMessageHandler<byte[], byte[]> {
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2017-2018 the original author or authors.
|
||||
* Copyright 2017-2021 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.
|
||||
@@ -17,6 +17,8 @@
|
||||
package org.springframework.cloud.stream.binder.kafka;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -34,7 +36,9 @@ import org.mockito.MockitoAnnotations;
|
||||
import org.springframework.boot.actuate.health.Health;
|
||||
import org.springframework.boot.actuate.health.Status;
|
||||
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
|
||||
import org.springframework.kafka.listener.AbstractMessageListenerContainer;
|
||||
|
||||
import static java.util.Collections.singleton;
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
/**
|
||||
@@ -43,6 +47,7 @@ import static org.assertj.core.api.Assertions.assertThat;
|
||||
* @author Laur Aliste
|
||||
* @author Soby Chacko
|
||||
* @author Chukwubuikem Ume-Ugwa
|
||||
* @author Taras Danylchuk
|
||||
*/
|
||||
public class KafkaBinderHealthIndicatorTest {
|
||||
|
||||
@@ -58,6 +63,12 @@ public class KafkaBinderHealthIndicatorTest {
|
||||
@Mock
|
||||
private KafkaConsumer consumer;
|
||||
|
||||
@Mock
|
||||
AbstractMessageListenerContainer<?, ?> listenerContainerA;
|
||||
|
||||
@Mock
|
||||
AbstractMessageListenerContainer<?, ?> listenerContainerB;
|
||||
|
||||
@Mock
|
||||
private KafkaMessageChannelBinder binder;
|
||||
|
||||
@@ -73,6 +84,21 @@ public class KafkaBinderHealthIndicatorTest {
|
||||
this.indicator.setTimeout(10);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void kafkaBinderIsUpWithNoConsumers() {
|
||||
final List<PartitionInfo> partitions = partitions(new Node(0, null, 0));
|
||||
topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation(
|
||||
"group1-healthIndicator", partitions, false));
|
||||
org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC))
|
||||
.willReturn(partitions);
|
||||
org.mockito.BDDMockito.given(binder.getKafkaMessageListenerContainers())
|
||||
.willReturn(Collections.emptyList());
|
||||
|
||||
Health health = indicator.health();
|
||||
assertThat(health.getStatus()).isEqualTo(Status.UP);
|
||||
assertThat(health.getDetails()).containsEntry("topicsInUse", singleton(TEST_TOPIC));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void kafkaBinderIsUp() {
|
||||
final List<PartitionInfo> partitions = partitions(new Node(0, null, 0));
|
||||
@@ -80,8 +106,42 @@ public class KafkaBinderHealthIndicatorTest {
|
||||
"group1-healthIndicator", partitions, false));
|
||||
org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC))
|
||||
.willReturn(partitions);
|
||||
org.mockito.BDDMockito.given(binder.getKafkaMessageListenerContainers())
|
||||
.willReturn(Arrays.asList(listenerContainerA, listenerContainerB));
|
||||
mockContainer(listenerContainerA, true);
|
||||
mockContainer(listenerContainerB, true);
|
||||
|
||||
Health health = indicator.health();
|
||||
assertThat(health.getStatus()).isEqualTo(Status.UP);
|
||||
assertThat(health.getDetails()).containsEntry("topicsInUse", singleton(TEST_TOPIC));
|
||||
assertThat(health.getDetails()).hasEntrySatisfying("listenerContainers", value ->
|
||||
assertThat((ArrayList<?>) value).hasSize(2));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void kafkaBinderIsDownWhenOneOfConsumersIsNotRunning() {
|
||||
final List<PartitionInfo> partitions = partitions(new Node(0, null, 0));
|
||||
topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation(
|
||||
"group1-healthIndicator", partitions, false));
|
||||
org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC))
|
||||
.willReturn(partitions);
|
||||
org.mockito.BDDMockito.given(binder.getKafkaMessageListenerContainers())
|
||||
.willReturn(Arrays.asList(listenerContainerA, listenerContainerB));
|
||||
mockContainer(listenerContainerA, false);
|
||||
mockContainer(listenerContainerB, true);
|
||||
|
||||
Health health = indicator.health();
|
||||
assertThat(health.getStatus()).isEqualTo(Status.DOWN);
|
||||
assertThat(health.getDetails()).containsEntry("topicsInUse", singleton(TEST_TOPIC));
|
||||
assertThat(health.getDetails()).hasEntrySatisfying("listenerContainers", value ->
|
||||
assertThat((ArrayList<?>) value).hasSize(2));
|
||||
}
|
||||
|
||||
private void mockContainer(AbstractMessageListenerContainer<?, ?> container, boolean isRunning) {
|
||||
org.mockito.BDDMockito.given(container.isRunning()).willReturn(isRunning);
|
||||
org.mockito.BDDMockito.given(container.isContainerPaused()).willReturn(true);
|
||||
org.mockito.BDDMockito.given(container.getListenerId()).willReturn("someListenerId");
|
||||
org.mockito.BDDMockito.given(container.getGroupId()).willReturn("someGroupId");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user