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 62d177777..cac3f93fa 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 @@ -86,14 +86,21 @@ public class KafkaBinderHealthIndicator implements HealthIndicator { Set downMessages = new HashSet<>(); final Map topicsInUse = KafkaBinderHealthIndicator.this.binder.getTopicsInUse(); - for (String topic : topicsInUse.keySet()) { - KafkaMessageChannelBinder.TopicInformation topicInformation = topicsInUse.get(topic); - if (!topicInformation.isTopicPattern()) { - List partitionInfos = this.metadataConsumer.partitionsFor(topic); - for (PartitionInfo partitionInfo : partitionInfos) { - if (topicInformation.getPartitionInfos() - .contains(partitionInfo) && partitionInfo.leader().id() == -1) { - downMessages.add(partitionInfo.toString()); + if (topicsInUse.isEmpty()) { + return Health.down() + .withDetail("No topic information available", "Kafka broker is not reachable") + .build(); + } + else { + for (String topic : topicsInUse.keySet()) { + KafkaMessageChannelBinder.TopicInformation topicInformation = topicsInUse.get(topic); + if (!topicInformation.isTopicPattern()) { + List partitionInfos = this.metadataConsumer.partitionsFor(topic); + for (PartitionInfo partitionInfo : partitionInfos) { + if (topicInformation.getPartitionInfos() + .contains(partitionInfo) && partitionInfo.leader().id() == -1) { + downMessages.add(partitionInfo.toString()); + } } } } 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 6be995fe5..1230be02a 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 @@ -738,7 +738,7 @@ public class KafkaMessageChannelBinder extends private Collection getPartitionInfo(String topic, final ExtendedConsumerProperties extendedConsumerProperties, final ConsumerFactory consumerFactory, int partitionCount) { - Collection allPartitions = provisioningProvider.getPartitionsForTopic(partitionCount, + return provisioningProvider.getPartitionsForTopic(partitionCount, extendedConsumerProperties.getExtension().isAutoRebalanceEnabled(), () -> { try (Consumer consumer = consumerFactory.createConsumer()) { @@ -746,7 +746,6 @@ public class KafkaMessageChannelBinder extends return partitionsFor; } }); - return allPartitions; } @Override 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 8b3db635d..6fa837c12 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 @@ -30,8 +30,6 @@ import org.junit.Test; import org.mockito.Mock; import org.mockito.Mockito; import org.mockito.MockitoAnnotations; -import org.mockito.invocation.InvocationOnMock; -import org.mockito.stubbing.Answer; import org.springframework.boot.actuate.health.Health; import org.springframework.boot.actuate.health.Status; @@ -105,15 +103,10 @@ public class KafkaBinderHealthIndicatorTest { public void kafkaBinderDoesNotAnswer() { final List partitions = partitions(new Node(-1, null, 0)); topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group3-healthIndicator", partitions, false)); - org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willAnswer(new Answer() { - - @Override - public Object answer(InvocationOnMock invocation) throws Throwable { - final int fiveMinutes = 1000 * 60 * 5; - Thread.sleep(fiveMinutes); - return partitions; - } - + org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willAnswer(invocation -> { + final int fiveMinutes = 1000 * 60 * 5; + Thread.sleep(fiveMinutes); + return partitions; }); this.indicator.setTimeout(1); Health health = indicator.health(); @@ -135,6 +128,9 @@ public class KafkaBinderHealthIndicatorTest { @Test public void consumerCreationFailsFirstTime() { + final List partitions = partitions(new Node(0, null, 0)); + topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("foo-healthIndicator", partitions, false)); + org.mockito.BDDMockito.given(consumerFactory.createConsumer()).willThrow(KafkaException.class) .willReturn(consumer); @@ -147,6 +143,12 @@ public class KafkaBinderHealthIndicatorTest { org.mockito.Mockito.verify(this.consumerFactory, Mockito.times(2)).createConsumer(); } + @Test + public void testIfNoTopicsRegisteredByTheBinderProvidesDownStatus() { + Health health = indicator.health(); + assertThat(health.getStatus()).isEqualTo(Status.DOWN); + } + private List partitions(Node leader) { List partitions = new ArrayList<>(); partitions.add(new PartitionInfo(TEST_TOPIC, 0, leader, null, null));