From efabc148414265677a981de277798badb858ee2d Mon Sep 17 00:00:00 2001 From: John Roesler Date: Tue, 27 Nov 2018 15:56:45 -0600 Subject: [PATCH 1/4] KAFKA-7660: fix parentSensors memory leak --- .../streams/processor/internals/metrics/StreamsMetricsImpl.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java index 8ec2711e764cc..c38e0a25fba4e 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java @@ -413,7 +413,7 @@ public void removeSensor(final Sensor sensor) { Objects.requireNonNull(sensor, "Sensor is null"); metrics.removeSensor(sensor.name()); - final Sensor parent = parentSensors.get(sensor); + final Sensor parent = parentSensors.remove(sensor); if (parent != null) { metrics.removeSensor(parent.name()); } From 9a6456a86d82418649de476d16b8a9e75ce0ce89 Mon Sep 17 00:00:00 2001 From: John Roesler Date: Tue, 27 Nov 2018 16:14:22 -0600 Subject: [PATCH 2/4] unit test --- .../processor/internals/metrics/StreamsMetricsImpl.java | 5 +++++ .../internals/{ => metrics}/StreamsMetricsImplTest.java | 5 ++++- 2 files changed, 9 insertions(+), 1 deletion(-) rename streams/src/test/java/org/apache/kafka/streams/processor/internals/{ => metrics}/StreamsMetricsImplTest.java (97%) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java index c38e0a25fba4e..991b81de9b210 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java @@ -419,6 +419,11 @@ public void removeSensor(final Sensor sensor) { } } + /** visible for testing */ + Map getParentSensors() { + return Collections.unmodifiableMap(parentSensors); + } + private static String groupNameFromScope(final String scopeName) { return "stream-" + scopeName + "-metrics"; } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsMetricsImplTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java similarity index 97% rename from streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsMetricsImplTest.java rename to streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java index 7ce27b4b6d668..b16082c70121d 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsMetricsImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.streams.processor.internals; +package org.apache.kafka.streams.processor.internals.metrics; import org.apache.kafka.common.MetricName; @@ -26,6 +26,7 @@ import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl; import org.junit.Test; +import java.util.Collections; import java.util.concurrent.TimeUnit; import static org.junit.Assert.assertEquals; @@ -62,6 +63,8 @@ public void testRemoveSensor() { final Sensor sensor3 = streamsMetrics.addThroughputSensor(scope, entity, operation, Sensor.RecordingLevel.DEBUG); streamsMetrics.removeSensor(sensor3); + + assertEquals(Collections.emptyMap(), streamsMetrics.getParentSensors()); } @Test From e68167cbfbd68e3aa14dc9867e10b9938d8b2f8a Mon Sep 17 00:00:00 2001 From: John Roesler Date: Wed, 28 Nov 2018 15:17:32 -0600 Subject: [PATCH 3/4] clean up and test creation and removal of x-Level sensors --- .../processor/internals/ProcessorNode.java | 2 + .../internals/metrics/StreamsMetricsImpl.java | 32 +++---- .../metrics/StreamsMetricsImplTest.java | 84 ++++++++++++++++--- 3 files changed, 87 insertions(+), 31 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorNode.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorNode.java index 8483791e2e79b..ec4f8e5e5db14 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorNode.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorNode.java @@ -234,9 +234,11 @@ private static Sensor createTaskAndNodeLatencyAndThroughputSensors(final String final Sensor parent = metrics.taskLevelSensor(taskName, operation, Sensor.RecordingLevel.DEBUG); addAvgMaxLatency(parent, PROCESSOR_NODE_METRICS_GROUP, taskTags, operation); addInvocationRateAndCount(parent, PROCESSOR_NODE_METRICS_GROUP, taskTags, operation); + final Sensor sensor = metrics.nodeLevelSensor(taskName, processorNodeName, operation, Sensor.RecordingLevel.DEBUG, parent); addAvgMaxLatency(sensor, PROCESSOR_NODE_METRICS_GROUP, nodeTags, operation); addInvocationRateAndCount(sensor, PROCESSOR_NODE_METRICS_GROUP, nodeTags, operation); + return sensor; } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java index 991b81de9b210..e335e197792c1 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java @@ -115,11 +115,9 @@ public final Sensor taskLevelSensor(final String taskName, public final void removeAllTaskLevelSensors(final String taskName) { final String key = taskSensorPrefix(taskName); synchronized (taskLevelSensors) { - if (taskLevelSensors.containsKey(key)) { - while (!taskLevelSensors.get(key).isEmpty()) { - metrics.removeSensor(taskLevelSensors.get(key).pop()); - } - taskLevelSensors.remove(key); + final Deque sensors = taskLevelSensors.remove(key); + while (sensors != null && !sensors.isEmpty()) { + metrics.removeSensor(sensors.pop()); } } } @@ -152,10 +150,9 @@ public Sensor nodeLevelSensor(final String taskName, public final void removeAllNodeLevelSensors(final String taskName, final String processorNodeName) { final String key = nodeSensorPrefix(taskName, processorNodeName); synchronized (nodeLevelSensors) { - if (nodeLevelSensors.containsKey(key)) { - while (!nodeLevelSensors.get(key).isEmpty()) { - metrics.removeSensor(nodeLevelSensors.get(key).pop()); - } + final Deque sensors = nodeLevelSensors.remove(key); + while (sensors != null && !sensors.isEmpty()) { + metrics.removeSensor(sensors.pop()); } } } @@ -188,11 +185,9 @@ public final Sensor cacheLevelSensor(final String taskName, public final void removeAllCacheLevelSensors(final String taskName, final String cacheName) { final String key = cacheSensorPrefix(taskName, cacheName); synchronized (cacheLevelSensors) { - if (cacheLevelSensors.containsKey(key)) { - while (!cacheLevelSensors.get(key).isEmpty()) { - metrics.removeSensor(cacheLevelSensors.get(key).pop()); - } - cacheLevelSensors.remove(key); + final Deque strings = cacheLevelSensors.remove(key); + while (strings != null && !strings.isEmpty()) { + metrics.removeSensor(strings.pop()); } } } @@ -225,11 +220,9 @@ public final Sensor storeLevelSensor(final String taskName, public final void removeAllStoreLevelSensors(final String taskName, final String storeName) { final String key = storeSensorPrefix(taskName, storeName); synchronized (storeLevelSensors) { - if (storeLevelSensors.containsKey(key)) { - while (!storeLevelSensors.get(key).isEmpty()) { - metrics.removeSensor(storeLevelSensors.get(key).pop()); - } - storeLevelSensors.remove(key); + final Deque sensors = storeLevelSensors.remove(key); + while (sensors != null && !sensors.isEmpty()) { + metrics.removeSensor(sensors.pop()); } } } @@ -415,6 +408,7 @@ public void removeSensor(final Sensor sensor) { final Sensor parent = parentSensors.remove(sensor); if (parent != null) { + metrics.removeSensor(parent.name()); } } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java index b16082c70121d..d4b981e4ebafc 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java @@ -23,13 +23,21 @@ import org.apache.kafka.common.metrics.Metrics; import org.apache.kafka.common.metrics.Sensor; import org.apache.kafka.common.utils.MockTime; -import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl; import org.junit.Test; import java.util.Collections; +import java.util.Map; import java.util.concurrent.TimeUnit; +import static org.apache.kafka.common.utils.Utils.mkEntry; +import static org.apache.kafka.common.utils.Utils.mkMap; +import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.PROCESSOR_NODE_METRICS_GROUP; +import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.addAvgMaxLatency; +import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.addInvocationRateAndCount; +import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.greaterThan; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertThat; public class StreamsMetricsImplTest { @@ -67,6 +75,58 @@ public void testRemoveSensor() { assertEquals(Collections.emptyMap(), streamsMetrics.getParentSensors()); } + @Test + public void testMutiLevelSensorRemoval() { + final Metrics registry = new Metrics(); + final StreamsMetricsImpl metrics = new StreamsMetricsImpl(registry, ""); + for (final MetricName defaultMetric : registry.metrics().keySet()) { + registry.removeMetric(defaultMetric); + } + + final String taskName = "taskName"; + final String operation = "operation"; + final Map taskTags = mkMap(mkEntry("tkey", "value")); + + final String processorNodeName = "processorNodeName"; + final Map nodeTags = mkMap(mkEntry("nkey", "value")); + + final Sensor parent1 = metrics.taskLevelSensor(taskName, operation, Sensor.RecordingLevel.DEBUG); + addAvgMaxLatency(parent1, PROCESSOR_NODE_METRICS_GROUP, taskTags, operation); + addInvocationRateAndCount(parent1, PROCESSOR_NODE_METRICS_GROUP, taskTags, operation); + + final int numberOfTaskMetrics = registry.metrics().size(); + + final Sensor sensor1 = metrics.nodeLevelSensor(taskName, processorNodeName, operation, Sensor.RecordingLevel.DEBUG, parent1); + addAvgMaxLatency(sensor1, PROCESSOR_NODE_METRICS_GROUP, nodeTags, operation); + addInvocationRateAndCount(sensor1, PROCESSOR_NODE_METRICS_GROUP, nodeTags, operation); + + assertThat(registry.metrics().size(), greaterThan(numberOfTaskMetrics)); + + metrics.removeAllNodeLevelSensors(taskName, processorNodeName); + + assertThat(registry.metrics().size(), equalTo(numberOfTaskMetrics)); + + final Sensor parent2 = metrics.taskLevelSensor(taskName, operation, Sensor.RecordingLevel.DEBUG); + addAvgMaxLatency(parent2, PROCESSOR_NODE_METRICS_GROUP, taskTags, operation); + addInvocationRateAndCount(parent2, PROCESSOR_NODE_METRICS_GROUP, taskTags, operation); + + assertThat(registry.metrics().size(), equalTo(numberOfTaskMetrics)); + + final Sensor sensor2 = metrics.nodeLevelSensor(taskName, processorNodeName, operation, Sensor.RecordingLevel.DEBUG, parent2); + addAvgMaxLatency(sensor2, PROCESSOR_NODE_METRICS_GROUP, nodeTags, operation); + addInvocationRateAndCount(sensor2, PROCESSOR_NODE_METRICS_GROUP, nodeTags, operation); + + assertThat(registry.metrics().size(), greaterThan(numberOfTaskMetrics)); + + metrics.removeAllNodeLevelSensors(taskName, processorNodeName); + + assertThat(registry.metrics().size(), equalTo(numberOfTaskMetrics)); + + metrics.removeAllTaskLevelSensors(taskName); + + assertThat(registry.metrics().size(), equalTo(0)); + } + @Test public void testLatencyMetrics() { final StreamsMetricsImpl streamsMetrics = new StreamsMetricsImpl(new Metrics(), ""); @@ -118,21 +178,21 @@ public void testTotalMetricDoesntDecrease() { final String operation = "op"; final Sensor sensor = streamsMetrics.addLatencyAndThroughputSensor( - scope, - entity, - operation, - Sensor.RecordingLevel.INFO + scope, + entity, + operation, + Sensor.RecordingLevel.INFO ); final double latency = 100.0; final MetricName totalMetricName = metrics.metricName( - "op-total", - "stream-scope-metrics", - "", - "client-id", - "", - "scope-id", - "entity" + "op-total", + "stream-scope-metrics", + "", + "client-id", + "", + "scope-id", + "entity" ); final KafkaMetric totalMetric = metrics.metric(totalMetricName); From 4db2541474143efda6df82e13f2d70bb2bcd0eba Mon Sep 17 00:00:00 2001 From: John Roesler Date: Wed, 28 Nov 2018 15:23:31 -0600 Subject: [PATCH 4/4] cr comments --- .../processor/internals/metrics/StreamsMetricsImpl.java | 7 ++++--- .../internals/metrics/StreamsMetricsImplTest.java | 2 +- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java index e335e197792c1..dd6cc4a218573 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImpl.java @@ -408,13 +408,14 @@ public void removeSensor(final Sensor sensor) { final Sensor parent = parentSensors.remove(sensor); if (parent != null) { - metrics.removeSensor(parent.name()); } } - /** visible for testing */ - Map getParentSensors() { + /** + * Visible for testing + */ + Map parentSensors() { return Collections.unmodifiableMap(parentSensors); } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java index d4b981e4ebafc..cadfdb0cef12b 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/metrics/StreamsMetricsImplTest.java @@ -72,7 +72,7 @@ public void testRemoveSensor() { final Sensor sensor3 = streamsMetrics.addThroughputSensor(scope, entity, operation, Sensor.RecordingLevel.DEBUG); streamsMetrics.removeSensor(sensor3); - assertEquals(Collections.emptyMap(), streamsMetrics.getParentSensors()); + assertEquals(Collections.emptyMap(), streamsMetrics.parentSensors()); } @Test