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