From d44f2348e67bfb972d0fde4961ba165ee58651b1 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 2 Oct 2018 14:43:07 -0400 Subject: [PATCH] Polishing, renamed method Resolves #452 --- .../stream/binder/kafka/KafkaBinderMetrics.java | 12 +++++------- .../kafka/integration/KafkaBinderActuatorTests.java | 2 +- 2 files changed, 6 insertions(+), 8 deletions(-) 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 aed73dc02..03dbcbd95 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 @@ -112,15 +112,15 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener calculateConsumerLagOnTopic(topic, group)) + o -> computeUnconsumedMessages(topic, group)) .tag("group", group) .tag("topic", topic) - .description("Consumer lag for a particular group and topic") + .description("Unconsumed messages for a particular group and topic") .register(registry); } } - private double calculateConsumerLagOnTopic(String topic, String group) { + private long computeUnconsumedMessages(String topic, String group) { ExecutorService exec = Executors.newSingleThreadExecutor(); Future future = exec.submit(() -> { @@ -144,11 +144,9 @@ public class KafkaBinderMetrics implements MeterBinder, ApplicationListener endOffset : endOffsets.entrySet()) { OffsetAndMetadata current = metadataConsumer.committed(endOffset.getKey()); + lag += endOffset.getValue(); if (current != null) { - lag += endOffset.getValue() - current.offset(); - } - else { - lag += endOffset.getValue(); + lag -= current.offset(); } } } diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java index 237a64dbd..926b0ed3b 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderActuatorTests.java @@ -91,7 +91,7 @@ public class KafkaBinderActuatorTests { assertThat(this.meterRegistry.get("spring.cloud.stream.binder.kafka.offset") .tag("group", TEST_CONSUMER_GROUP) .tag("topic", Sink.INPUT) - .timeGauge().value()).isGreaterThan(0); + .gauge().value()).isGreaterThan(0); } @Test