diff --git a/spring-cloud-stream-binder-kafka/.settings/org.eclipse.jdt.ui.prefs b/spring-cloud-stream-binder-kafka/.settings/org.eclipse.jdt.ui.prefs index f9aac64a9..9b721a32a 100644 --- a/spring-cloud-stream-binder-kafka/.settings/org.eclipse.jdt.ui.prefs +++ b/spring-cloud-stream-binder-kafka/.settings/org.eclipse.jdt.ui.prefs @@ -1,5 +1,5 @@ eclipse.preferences.version=1 org.eclipse.jdt.ui.ignorelowercasenames=true -org.eclipse.jdt.ui.importorder=java;javax;com;org;org.springframework;ch.qos;\#; +org.eclipse.jdt.ui.importorder=java;javax;com;io.micrometer;org;org.springframework;ch.qos;\#; org.eclipse.jdt.ui.ondemandthreshold=99 org.eclipse.jdt.ui.staticondemandthreshold=99 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 331956fac..7677ff6c6 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 @@ -1,5 +1,5 @@ /* - * Copyright 2016-2017 the original author or authors. + * Copyright 2016-2018 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. @@ -78,25 +78,31 @@ public class KafkaBinderHealthIndicator implements HealthIndicator { public Health call() { try { if (metadataConsumer == null) { - 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()); + synchronized(KafkaBinderHealthIndicator.this) { + if (metadataConsumer == null) { + metadataConsumer = consumerFactory.createConsumer(); } } } - if (downMessages.isEmpty()) { - return Health.up().build(); - } - else { - return Health.down() - .withDetail("Following partitions in use have no leaders: ", downMessages.toString()) - .build(); + synchronized (metadataConsumer) { + 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) { diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java index 6df5f5011..5cf3db77a 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java @@ -20,11 +20,17 @@ import java.util.HashMap; import java.util.LinkedList; import java.util.List; import java.util.Map; +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 io.micrometer.core.instrument.MeterRegistry; import io.micrometer.core.instrument.TimeGauge; import io.micrometer.core.instrument.binder.MeterBinder; + import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.clients.consumer.Consumer; @@ -51,9 +57,12 @@ import org.springframework.util.ObjectUtils; * @author Oleg Zhurakousky * @author Jon Schneider * @author Thomas Cheyney + * @author Gary Russell */ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener { + private static final int DEFAULT_TIMEOUT = 60; + private final static Log LOG = LogFactory.getLog(KafkaBinderMetrics.class); static final String METRIC_NAME = "spring.cloud.stream.binder.kafka.offset"; @@ -68,6 +77,8 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener metadataConsumer; + private int timeout = DEFAULT_TIMEOUT; + public KafkaBinderMetrics(KafkaMessageChannelBinder binder, KafkaBinderConfigurationProperties binderConfigurationProperties, ConsumerFactory defaultConsumerFactory, @Nullable MeterRegistry meterRegistry) { @@ -84,6 +95,10 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener topicInfo : this.binder.getTopicsInUse() @@ -106,33 +121,56 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener future = exec.submit(() -> { + + long lag = 0; + try { + if (metadataConsumer == null) { + synchronized(KafkaBinderMetrics.this) { + if (metadataConsumer == null) { + metadataConsumer = createConsumerFactory(group).createConsumer(); + } + } + } + synchronized (metadataConsumer) { + List partitionInfos = metadataConsumer.partitionsFor(topic); + List topicPartitions = new LinkedList<>(); + for (PartitionInfo partitionInfo : partitionInfos) { + topicPartitions.add(new TopicPartition(partitionInfo.topic(), partitionInfo.partition())); + } + + Map endOffsets = metadataConsumer.endOffsets(topicPartitions); + + for (Map.Entry endOffset : endOffsets.entrySet()) { + OffsetAndMetadata current = metadataConsumer.committed(endOffset.getKey()); + if (current != null) { + lag += endOffset.getValue() - current.offset(); + } + else { + lag += endOffset.getValue(); + } + } + } + } + catch (Exception e) { + LOG.debug("Cannot generate metric for topic: " + topic, e); + } + return lag; + }); try { - if (metadataConsumer == null) { - metadataConsumer = createConsumerFactory(group).createConsumer(); - } - List partitionInfos = metadataConsumer.partitionsFor(topic); - List topicPartitions = new LinkedList<>(); - for (PartitionInfo partitionInfo : partitionInfos) { - topicPartitions.add(new TopicPartition(partitionInfo.topic(), partitionInfo.partition())); - } - - Map endOffsets = metadataConsumer.endOffsets(topicPartitions); - - for (Map.Entry endOffset : endOffsets.entrySet()) { - OffsetAndMetadata current = metadataConsumer.committed(endOffset.getKey()); - if (current != null) { - lag += endOffset.getValue() - current.offset(); - } - else { - lag += endOffset.getValue(); - } - } + return future.get(this.timeout, TimeUnit.SECONDS); } - catch (Exception e) { - LOG.debug("Cannot generate metric for topic: " + topic, e); + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return 0L; + } + catch (ExecutionException | TimeoutException e) { + return 0L; + } + finally { + exec.shutdownNow(); } - return lag; } private ConsumerFactory createConsumerFactory(String group) {