From ba2c3a05c9af6eea1b127c0aac899db735eb0290 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 18 Aug 2021 19:47:27 -0400 Subject: [PATCH] GH-1129: Kafka Binder Metrics Improvements Avoid blocking committed() call in KafkaBinderMetrics in a loop for each topic partition. Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/1129 --- .../binder/kafka/KafkaBinderMetrics.java | 7 ++++--- .../binder/kafka/KafkaBinderMetricsTest.java | 18 +++++++++++++----- 2 files changed, 17 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 dc7ebb25b..f07e9aee1 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 @@ -1,5 +1,5 @@ /* - * Copyright 2016-2020 the original author or authors. + * Copyright 2016-2021 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. @@ -208,10 +208,11 @@ public class KafkaBinderMetrics Map endOffsets = metadataConsumer .endOffsets(topicPartitions); + final Map committedOffsets = metadataConsumer.committed(endOffsets.keySet()); + for (Map.Entry endOffset : endOffsets .entrySet()) { - OffsetAndMetadata current = metadataConsumer - .committed(endOffset.getKey()); + OffsetAndMetadata current = committedOffsets.get(endOffset.getKey()); lag += endOffset.getValue(); if (current != null) { lag -= current.offset(); diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java index 480d7795c..a1e3a432c 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2020 the original author or authors. + * Copyright 2016-2021 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. @@ -89,9 +89,12 @@ public class KafkaBinderMetricsTest { @Test public void shouldIndicateLag() { + final Map committed = new HashMap<>(); + TopicPartition topicPartition = new TopicPartition(TEST_TOPIC, 0); + committed.put(topicPartition, new OffsetAndMetadata(500)); org.mockito.BDDMockito - .given(consumer.committed(ArgumentMatchers.any(TopicPartition.class))) - .willReturn(new OffsetAndMetadata(500)); + .given(consumer.committed(ArgumentMatchers.anySet())) + .willReturn(committed); List partitions = partitions(new Node(0, null, 0)); topicsInUse.put(TEST_TOPIC, new TopicInformation("group1-metrics", partitions, false)); @@ -133,9 +136,14 @@ public class KafkaBinderMetricsTest { org.mockito.BDDMockito .given(consumer.endOffsets(ArgumentMatchers.anyCollection())) .willReturn(endOffsets); + final Map committed = new HashMap<>(); + TopicPartition topicPartition1 = new TopicPartition(TEST_TOPIC, 0); + TopicPartition topicPartition2 = new TopicPartition(TEST_TOPIC, 1); + committed.put(topicPartition1, new OffsetAndMetadata(500)); + committed.put(topicPartition2, new OffsetAndMetadata(500)); org.mockito.BDDMockito - .given(consumer.committed(ArgumentMatchers.any(TopicPartition.class))) - .willReturn(new OffsetAndMetadata(500)); + .given(consumer.committed(ArgumentMatchers.anySet())) + .willReturn(committed); List partitions = partitions(new Node(0, null, 0), new Node(0, null, 0)); topicsInUse.put(TEST_TOPIC,