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