diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java index 1de105d74..ecdfd5177 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderHealthIndicator.java @@ -16,6 +16,15 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.time.Duration; +import java.util.HashMap; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; +import java.util.stream.Collectors; + import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.clients.admin.AdminClient; @@ -24,6 +33,8 @@ import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.processor.TaskMetadata; import org.apache.kafka.streams.processor.ThreadMetadata; + +import org.springframework.beans.factory.DisposableBean; import org.springframework.boot.actuate.health.AbstractHealthIndicator; import org.springframework.boot.actuate.health.Health; import org.springframework.boot.actuate.health.Status; @@ -32,19 +43,13 @@ import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProv import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; import org.springframework.kafka.config.StreamsBuilderFactoryBean; -import java.util.HashMap; -import java.util.Map; -import java.util.Set; -import java.util.concurrent.TimeUnit; -import java.util.stream.Collectors; - /** * Health indicator for Kafka Streams. * * @author Arnaud Jardiné * @author Soby Chacko */ -public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator { +public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator implements DisposableBean { private final Log logger = LogFactory.getLog(getClass()); @@ -57,6 +62,10 @@ public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator { private static final ThreadLocal healthStatusThreadLocal = new ThreadLocal<>(); + private AdminClient adminClient; + + private final Lock lock = new ReentrantLock(); + KafkaStreamsBinderHealthIndicator(KafkaStreamsRegistry kafkaStreamsRegistry, KafkaStreamsBinderConfigurationProperties kafkaStreamsBinderConfigurationProperties, KafkaProperties kafkaProperties, @@ -73,7 +82,11 @@ public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator { @Override protected void doHealthCheck(Health.Builder builder) throws Exception { - try (AdminClient adminClient = AdminClient.create(this.adminClientProperties)) { + try { + this.lock.lock(); + if (this.adminClient == null) { + this.adminClient = AdminClient.create(this.adminClientProperties); + } final Status status = healthStatusThreadLocal.get(); //If one of the kafka streams binders (kstream, ktable, globalktable) was down before on the same request, //retrieve that from the thead local storage where it was saved before. This is done in order to avoid @@ -82,14 +95,16 @@ public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator { if (status == Status.DOWN) { builder.withDetail("No topic information available", "Kafka broker is not reachable"); builder.status(Status.DOWN); - } else { - final ListTopicsResult listTopicsResult = adminClient.listTopics(); + } + else { + final ListTopicsResult listTopicsResult = this.adminClient.listTopics(); listTopicsResult.listings().get(this.configurationProperties.getHealthTimeout(), TimeUnit.SECONDS); if (this.kafkaStreamsBindingInformationCatalogue.getStreamsBuilderFactoryBeans().isEmpty()) { builder.withDetail("No Kafka Streams bindings have been established", "Kafka Streams binder did not detect any processors"); builder.status(Status.UNKNOWN); - } else { + } + else { boolean up = true; for (KafkaStreams kStream : kafkaStreamsRegistry.getKafkaStreams()) { up &= kStream.state().isRunning(); @@ -98,13 +113,17 @@ public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator { builder.status(up ? Status.UP : Status.DOWN); } } - } catch (Exception e) { + } + catch (Exception e) { builder.withDetail("No topic information available", "Kafka broker is not reachable"); builder.status(Status.DOWN); builder.withException(e); //Store binder down status into a thread local storage. healthStatusThreadLocal.set(Status.DOWN); } + finally { + this.lock.unlock(); + } } private Map buildDetails(KafkaStreams kafkaStreams) { @@ -146,4 +165,10 @@ public class KafkaStreamsBinderHealthIndicator extends AbstractHealthIndicator { return details; } + @Override + public void destroy() throws Exception { + if (adminClient != null) { + adminClient.close(Duration.ofSeconds(0)); + } + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java index 37c8b5c93..54dcf9d58 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java @@ -16,6 +16,11 @@ package org.springframework.cloud.stream.binder.kafka.streams.integration; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.producer.ProducerRecord; @@ -26,6 +31,7 @@ import org.junit.Assert; import org.junit.BeforeClass; import org.junit.ClassRule; import org.junit.Test; + import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.actuate.health.CompositeHealthContributor; @@ -50,11 +56,6 @@ import org.springframework.messaging.handler.annotation.SendTo; import org.springframework.util.concurrent.ListenableFuture; import org.springframework.util.concurrent.ListenableFutureCallback; -import java.util.List; -import java.util.Map; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.TimeUnit; - import static org.assertj.core.api.Assertions.assertThat; /**