diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java index 2a9807d09..c1f3800d8 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java @@ -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 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 aggregatedDetails = new HashMap<>(); + aggregatedDetails.putAll(topicsHealth.getDetails()); + aggregatedDetails.putAll(listenerContainersHealth.getDetails()); + return Health.status(aggregatedStatus).withDetails(aggregatedDetails).build(); + } + + private Health safelyBuildTopicsHealth() { + Future 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 downMessages = new HashSet<>(); + Set checkedTopics = new HashSet<>(); final Map 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> listenerContainers = binder.getKafkaMessageListenerContainers(); + if (listenerContainers.isEmpty()) { + return Health.unknown().build(); + } + + Status status = Status.UP; + List> containersDetails = new ArrayList<>(); + + for (AbstractMessageListenerContainer container : listenerContainers) { + Map 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(); diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index ad32cb47e..730bf412a 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -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> 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> getKafkaMessageListenerContainers() { + return Collections.unmodifiableList(kafkaMessageListenerContainers); + } + private final class ProducerConfigurationMessageHandler extends KafkaProducerMessageHandler { diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicatorTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicatorTest.java index 7ad32169d..69aa0614b 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicatorTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicatorTest.java @@ -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 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 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 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