From 2efd29fb273926b44258276978dd6c2097d6bb47 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 4 Oct 2021 17:13:01 -0400 Subject: [PATCH] GH-1138: HealthIndicator Improvements Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/1138 Don't report DOWN if a container is stopped normally. This is a valid state when containers are not auto-startup or are stopped while the app remains running. Containers are stopped abnormally when - a listener throws an `Error` - a `CommonContainerStoppingErrorHandler` (or similar) is configured to stop the container after an error. --- spring-cloud-stream-binder-kafka/pom.xml | 1 + .../kafka/KafkaBinderHealthIndicator.java | 4 ++- .../kafka/KafkaBinderHealthIndicatorTest.java | 32 ++++++++++++++++--- 3 files changed, 31 insertions(+), 6 deletions(-) diff --git a/spring-cloud-stream-binder-kafka/pom.xml b/spring-cloud-stream-binder-kafka/pom.xml index df027c2af..6d1548678 100644 --- a/spring-cloud-stream-binder-kafka/pom.xml +++ b/spring-cloud-stream-binder-kafka/pom.xml @@ -44,6 +44,7 @@ org.springframework.kafka spring-kafka + 2.8.0-SNAPSHOT org.springframework.boot 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 c1f3800d8..594ca55cc 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 @@ -201,10 +201,12 @@ public class KafkaBinderHealthIndicator implements HealthIndicator, DisposableBe for (AbstractMessageListenerContainer container : listenerContainers) { Map containerDetails = new HashMap<>(); boolean isRunning = container.isRunning(); - if (!isRunning) { + boolean isOk = container.isInExpectedState(); + if (!isOk) { status = Status.DOWN; } containerDetails.put("isRunning", isRunning); + containerDetails.put("isStoppedAbnormally", !isRunning && !isOk); containerDetails.put("isPaused", container.isContainerPaused()); containerDetails.put("listenerId", container.getListenerId()); containerDetails.put("groupId", container.getGroupId()); 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 69aa0614b..64d84e842 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 @@ -108,8 +108,8 @@ public class KafkaBinderHealthIndicatorTest { .willReturn(partitions); org.mockito.BDDMockito.given(binder.getKafkaMessageListenerContainers()) .willReturn(Arrays.asList(listenerContainerA, listenerContainerB)); - mockContainer(listenerContainerA, true); - mockContainer(listenerContainerB, true); + mockContainer(listenerContainerA, true, true); + mockContainer(listenerContainerB, true, true); Health health = indicator.health(); assertThat(health.getStatus()).isEqualTo(Status.UP); @@ -127,8 +127,27 @@ public class KafkaBinderHealthIndicatorTest { .willReturn(partitions); org.mockito.BDDMockito.given(binder.getKafkaMessageListenerContainers()) .willReturn(Arrays.asList(listenerContainerA, listenerContainerB)); - mockContainer(listenerContainerA, false); - mockContainer(listenerContainerB, true); + mockContainer(listenerContainerA, false, true); + mockContainer(listenerContainerB, true, 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 kafkaBinderIsDownWhenOneOfContainersWasStoppedAbnormally() { + 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, false); + mockContainer(listenerContainerB, true, true); Health health = indicator.health(); assertThat(health.getStatus()).isEqualTo(Status.DOWN); @@ -137,11 +156,14 @@ public class KafkaBinderHealthIndicatorTest { assertThat((ArrayList) value).hasSize(2)); } - private void mockContainer(AbstractMessageListenerContainer container, boolean isRunning) { + private void mockContainer(AbstractMessageListenerContainer container, boolean isRunning, + boolean normalState) { + 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"); + org.mockito.BDDMockito.given(container.isInExpectedState()).willReturn(normalState); } @Test