From 69227166c7a356f975163a024ac9c33bb5efbbc4 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 29 Sep 2017 19:47:45 +0100 Subject: [PATCH] GH-215: Add timeout to health indicator Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/215 * Shutdown the executor. * Polishing - PR Comments * Re-interrupt thread. * More Polishing --- .../KafkaBinderConfigurationProperties.java | 16 +++- .../src/main/asciidoc/overview.adoc | 7 +- .../kafka/KafkaBinderHealthIndicator.java | 79 +++++++++++++++---- .../config/KafkaBinderConfiguration.java | 6 +- .../kafka/KafkaBinderHealthIndicatorTest.java | 35 ++++++-- 5 files changed, 121 insertions(+), 22 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java index 235274c5a..b5442d4ad 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java @@ -1,5 +1,5 @@ /* - * Copyright 2015-2016 the original author or authors. + * Copyright 2015-2017 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. @@ -35,6 +35,7 @@ import org.springframework.util.StringUtils; * @author Ilayaperumal Gopinathan * @author Marius Bogoevici * @author Soby Chacko + * @author Gary Russell */ @ConfigurationProperties(prefix = "spring.cloud.stream.kafka.binder") public class KafkaBinderConfigurationProperties { @@ -88,6 +89,11 @@ public class KafkaBinderConfigurationProperties { private int queueSize = 8192; + /** + * Time to wait to get partition information in seconds; default 60. + */ + private int healthTimeout = 60; + private JaasLoginModuleConfiguration jaas; public String getZkConnectionString() { @@ -228,6 +234,14 @@ public class KafkaBinderConfigurationProperties { this.minPartitionCount = minPartitionCount; } + public int getHealthTimeout() { + return this.healthTimeout; + } + + public void setHealthTimeout(int healthTimeout) { + this.healthTimeout = healthTimeout; + } + public int getQueueSize() { return this.queueSize; } diff --git a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc index d84af136b..d2c2fb590 100644 --- a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc +++ b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc @@ -73,6 +73,11 @@ spring.cloud.stream.kafka.binder.headers:: The list of custom headers that will be transported by the binder. + Default: empty. +spring.cloud.stream.kafka.binder.healthTimeout:: + The time to wait to get partition information in seconds; default 60. + Health will report as down if this timer expires. ++ +Default: 10. spring.cloud.stream.kafka.binder.offsetUpdateTimeWindow:: The frequency, in milliseconds, with which offsets are saved. Ignored if `0`. @@ -572,4 +577,4 @@ Kafka binder module exposes the following metrics: `spring.cloud.stream.binder.kafka.someGroup.someTopic.lag` - this metric indicates how many messages have not been yet consumed from given binder's topic by given consumer group. For example if the value of the metric `spring.cloud.stream.binder.kafka.myGroup.myTopic.lag` is `1000`, then consumer group `myGroup` has `1000` messages to waiting to be consumed from topic `myTopic`. -This metric is particularly useful to provide auto-scaling feedback to PaaS platform of your choice. \ No newline at end of file +This metric is particularly useful to provide auto-scaling feedback to PaaS platform of your choice. 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 99e0a9e44..32324eae8 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 @@ -19,6 +19,13 @@ package org.springframework.cloud.stream.binder.kafka; import java.util.HashSet; import java.util.List; import java.util.Set; +import java.util.concurrent.Callable; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.common.PartitionInfo; @@ -33,13 +40,18 @@ import org.springframework.kafka.core.ConsumerFactory; * @author Ilayaperumal Gopinathan * @author Marius Bogoevici * @author Henryk Konsek + * @author Gary Russell */ public class KafkaBinderHealthIndicator implements HealthIndicator { + private static final int DEFAULT_TIMEOUT = 60; + private final KafkaMessageChannelBinder binder; private final ConsumerFactory consumerFactory; + private int timeout = DEFAULT_TIMEOUT; + public KafkaBinderHealthIndicator(KafkaMessageChannelBinder binder, ConsumerFactory consumerFactory) { this.binder = binder; @@ -47,28 +59,67 @@ public class KafkaBinderHealthIndicator implements HealthIndicator { } + /** + * Set the timeout in seconds to retrieve health information. + * @param timeout the timeout - default 60. + */ + public void setTimeout(int timeout) { + this.timeout = timeout; + } + @Override public Health health() { - try (Consumer metadataConsumer = consumerFactory.createConsumer()) { - Set downMessages = new HashSet<>(); - for (String topic : this.binder.getTopicsInUse().keySet()) { - List partitionInfos = metadataConsumer.partitionsFor(topic); - for (PartitionInfo partitionInfo : partitionInfos) { - if (this.binder.getTopicsInUse().get(topic).getPartitionInfos().contains(partitionInfo) - && partitionInfo.leader() - .id() == -1) { - downMessages.add(partitionInfo.toString()); + ExecutorService exec = Executors.newSingleThreadExecutor(); + Future future = exec.submit(new Callable() { + + @Override + public Health call() { + try (Consumer metadataConsumer = consumerFactory.createConsumer()) { + Set downMessages = new HashSet<>(); + for (String topic : KafkaBinderHealthIndicator.this.binder.getTopicsInUse().keySet()) { + List partitionInfos = metadataConsumer.partitionsFor(topic); + for (PartitionInfo partitionInfo : partitionInfos) { + if (KafkaBinderHealthIndicator.this.binder.getTopicsInUse().get(topic).getPartitionInfos() + .contains(partitionInfo) && partitionInfo.leader().id() == -1) { + downMessages.add(partitionInfo.toString()); + } + } + } + if (downMessages.isEmpty()) { + return Health.up().build(); + } + else { + return Health.down() + .withDetail("Following partitions in use have no leaders: ", downMessages.toString()) + .build(); } } + catch (Exception e) { + return Health.down(e).build(); + } } - if (downMessages.isEmpty()) { - return Health.up().build(); - } - return Health.down().withDetail("Following partitions in use have no leaders: ", downMessages.toString()) + + }); + try { + return future.get(this.timeout, TimeUnit.SECONDS); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return Health.down() + .withDetail("Interrupted while waiting for partition information in", this.timeout + " seconds") .build(); } - catch (Exception e) { + catch (ExecutionException e) { return Health.down(e).build(); } + catch (TimeoutException e) { + return Health.down() + .withDetail("Failed to retrieve partition information in", this.timeout + " seconds") + .build(); + } + finally { + exec.shutdownNow(); + } } + } diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java index 5f0f55dd0..88cf0ec4f 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java @@ -68,6 +68,7 @@ import org.springframework.util.ObjectUtils; * @author Mark Fisher * @author Ilayaperumal Gopinathan * @author Henryk Konsek + * @author Gary Russell */ @Configuration @ConditionalOnMissingBean(Binder.class) @@ -128,7 +129,10 @@ public class KafkaBinderConfiguration { props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, this.configurationProperties.getKafkaConnectionString()); } ConsumerFactory consumerFactory = new DefaultKafkaConsumerFactory<>(props); - return new KafkaBinderHealthIndicator(kafkaMessageChannelBinder, consumerFactory); + KafkaBinderHealthIndicator indicator = new KafkaBinderHealthIndicator(kafkaMessageChannelBinder, + consumerFactory); + indicator.setTimeout(this.configurationProperties.getHealthTimeout()); + return indicator; } @Bean 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 6e7250341..72b948af3 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 @@ -15,6 +15,10 @@ */ package org.springframework.cloud.stream.binder.kafka; +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.BDDMockito.given; +import static org.mockito.Mockito.verify; + import java.util.ArrayList; import java.util.HashMap; import java.util.List; @@ -27,16 +31,16 @@ import org.junit.Before; import org.junit.Test; import org.mockito.Mock; 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; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; -import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.BDDMockito.given; - /** * @author Barry Commins + * @author Gary Russell */ public class KafkaBinderHealthIndicatorTest { @@ -53,14 +57,15 @@ public class KafkaBinderHealthIndicatorTest { @Mock private KafkaMessageChannelBinder binder; - private Map topicsInUse = new HashMap<>(); + private final Map topicsInUse = new HashMap<>(); @Before public void setup() { MockitoAnnotations.initMocks(this); given(consumerFactory.createConsumer()).willReturn(consumer); given(binder.getTopicsInUse()).willReturn(topicsInUse); - indicator = new KafkaBinderHealthIndicator(binder, consumerFactory); + this.indicator = new KafkaBinderHealthIndicator(binder, consumerFactory); + this.indicator.setTimeout(10); } @Test @@ -70,6 +75,7 @@ public class KafkaBinderHealthIndicatorTest { given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); Health health = indicator.health(); assertThat(health.getStatus()).isEqualTo(Status.UP); + verify(this.consumer).close(); } @Test @@ -81,6 +87,25 @@ public class KafkaBinderHealthIndicatorTest { assertThat(health.getStatus()).isEqualTo(Status.DOWN); } + @Test(timeout = 5000) + public void kafkaBinderDoesNotAnswer() { + final List partitions = partitions(new Node(-1, null, 0)); + topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group", partitions)); + 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; + } + + }); + this.indicator.setTimeout(1); + 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));