From 4d1b69943cc1d8e3ec5e92d8d048ca5f816b0ea0 Mon Sep 17 00:00:00 2001 From: Sam Lendle Date: Fri, 3 Aug 2018 17:23:22 -0700 Subject: [PATCH 1/7] Test incorrect behavior in StreamsMetricsImplTest, add CumulativeCount class, use CumulativeCount in StreamsMetricsImpl#addThroughputMetrics --- .../internals/metrics/CumulativeCount.java | 38 +++++++++++++++++ .../internals/metrics/StreamsMetricsImpl.java | 2 +- .../internals/StreamsMetricsImplTest.java | 42 +++++++++++++++++++ 3 files changed, 81 insertions(+), 1 deletion(-) create mode 100644 streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/CumulativeCount.java diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/CumulativeCount.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/CumulativeCount.java new file mode 100644 index 0000000000000..54e698324dd75 --- /dev/null +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/CumulativeCount.java @@ -0,0 +1,38 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.streams.processor.internals.metrics; + +import org.apache.kafka.common.metrics.MeasurableStat; +import org.apache.kafka.common.metrics.MetricConfig; + +/** + * A non-SampledStat version of Count for measuring -total metrics in streams + */ +public class CumulativeCount implements MeasurableStat { + + private double count = 0.0; + + @Override + public void record(MetricConfig config, double value, long timeMs) { + count += 1; + } + + @Override + public double measure(MetricConfig config, long now) { + return count; + } +} 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 56166a4d372a4..170311238ecf3 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 @@ -398,7 +398,7 @@ public static void addInvocationRateAndCount(final Sensor sensor, "The total number of occurrence of " + operation + " operations.", tags ), - new Count() + new CumulativeCount() ); } 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/StreamsMetricsImplTest.java index b065e2ca7b51c..d5ef826959ee8 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/StreamsMetricsImplTest.java @@ -17,8 +17,12 @@ package org.apache.kafka.streams.processor.internals; +import org.apache.kafka.common.MetricName; +import org.apache.kafka.common.metrics.KafkaMetric; +import org.apache.kafka.common.metrics.MetricConfig; 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; @@ -96,4 +100,42 @@ public void testThroughputMetrics() { streamsMetrics.removeSensor(sensor1); assertEquals(defaultMetrics, streamsMetrics.metrics().size()); } + + @Test + public void testTotalMetricDoesntDecrease() { + final MockTime time = new MockTime(); + final MetricConfig config = new MetricConfig().eventWindow(1).samples(2); + final Metrics metrics = new Metrics(config, time); + final StreamsMetricsImpl streamsMetrics = new StreamsMetricsImpl(metrics, ""); + + final String scope = "scope"; + final String entity = "entity"; + final String operation = "op"; + + final Sensor sensor = streamsMetrics.addLatencyAndThroughputSensor( + 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" + ); + + final KafkaMetric totalMetric = metrics.metric(totalMetricName); + + for (int i = 0; i < 10; i++) { + assertEquals(i, totalMetric.measurable().measure(config, time.milliseconds()), 0.0001); + sensor.record(latency, time.milliseconds()); + } + + } } From d18d419d50b7d9850b3ab92d43f71014667ab276 Mon Sep 17 00:00:00 2001 From: Sam Lendle Date: Fri, 3 Aug 2018 21:42:21 -0700 Subject: [PATCH 2/7] Use CumulativeCount for -total metrics in StreamTask and StreamThread --- .../kafka/streams/processor/internals/StreamTask.java | 7 ++++--- .../kafka/streams/processor/internals/StreamThread.java | 3 ++- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index 79df5d158a169..7f3d31fd77607 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -41,6 +41,7 @@ import org.apache.kafka.streams.processor.Punctuator; import org.apache.kafka.streams.processor.TaskId; import org.apache.kafka.streams.processor.TimestampExtractor; +import org.apache.kafka.streams.processor.internals.metrics.CumulativeCount; import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl; import org.apache.kafka.streams.state.internals.ThreadCache; @@ -109,7 +110,7 @@ protected static final class TaskMetrics { ); parent.add( new MetricName("commit-total", group, "The total number of occurrence of commit operations.", allTagMap), - new Count() + new CumulativeCount() ); // add the operation metrics with additional tags @@ -129,7 +130,7 @@ protected static final class TaskMetrics { ); taskCommitTimeSensor.add( new MetricName("commit-total", group, "The total number of occurrence of commit operations.", tagMap), - new Count() + new CumulativeCount() ); // add the metrics for enforced processing @@ -140,7 +141,7 @@ protected static final class TaskMetrics { ); taskEnforcedProcessSensor.add( new MetricName("enforced-process-total", group, "The total number of occurrence of enforced-process operations.", tagMap), - new Count() + new CumulativeCount() ); } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index efd94eaf6373f..87b7adce53211 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -44,6 +44,7 @@ import org.apache.kafka.streams.processor.TaskId; import org.apache.kafka.streams.processor.TaskMetadata; import org.apache.kafka.streams.processor.ThreadMetadata; +import org.apache.kafka.streams.processor.internals.metrics.CumulativeCount; import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl; import org.apache.kafka.streams.state.internals.ThreadCache; import org.slf4j.Logger; @@ -532,7 +533,7 @@ static class StreamsMetricsThreadImpl extends StreamsMetricsImpl { addAvgMaxLatency(pollTimeSensor, group, tagMap(), "poll"); // can't use addInvocationRateAndCount due to non-standard description string pollTimeSensor.add(metrics.metricName("poll-rate", group, "The average per-second number of record-poll calls", tagMap()), new Rate(TimeUnit.SECONDS, new Count())); - pollTimeSensor.add(metrics.metricName("poll-total", group, "The total number of record-poll calls", tagMap()), new Count()); + pollTimeSensor.add(metrics.metricName("poll-total", group, "The total number of record-poll calls", tagMap()), new CumulativeCount()); processTimeSensor = threadLevelSensor("process-latency", Sensor.RecordingLevel.INFO); addAvgMaxLatency(processTimeSensor, group, tagMap(), "process"); From ab7378f97294ff2616fc6ffa0faaa6f2e07ce881 Mon Sep 17 00:00:00 2001 From: Sam Lendle Date: Mon, 6 Aug 2018 13:41:35 -0700 Subject: [PATCH 3/7] Rename tasksClosedSensor -> taskClosedSensor for consistency, add calls to taskClosedSensor.record --- .../processor/internals/StreamThread.java | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index 87b7adce53211..1d31ebb8f5d64 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -438,7 +438,7 @@ StreamTask createTask(final Consumer consumer, cache, time, () -> createProducer(taskId), - streamsMetrics.tasksClosedSensor); + streamsMetrics.taskClosedSensor); } private Producer createProducer(final TaskId id) { @@ -455,6 +455,7 @@ private Producer createProducer(final TaskId id) { @Override public void close() { + streamsMetrics.taskClosedSensor.record(); if (threadProducer != null) { try { threadProducer.close(); @@ -510,6 +511,11 @@ StandbyTask createTask(final Consumer consumer, return null; } } + + @Override + public void close() { + streamsMetrics.taskClosedSensor.record(); + } } static class StreamsMetricsThreadImpl extends StreamsMetricsImpl { @@ -519,7 +525,7 @@ static class StreamsMetricsThreadImpl extends StreamsMetricsImpl { private final Sensor processTimeSensor; private final Sensor punctuateTimeSensor; private final Sensor taskCreatedSensor; - private final Sensor tasksClosedSensor; + private final Sensor taskClosedSensor; StreamsMetricsThreadImpl(final Metrics metrics, final String threadName) { super(metrics, threadName); @@ -547,9 +553,9 @@ static class StreamsMetricsThreadImpl extends StreamsMetricsImpl { taskCreatedSensor.add(metrics.metricName("task-created-rate", "stream-metrics", "The average per-second number of newly created tasks", tagMap()), new Rate(TimeUnit.SECONDS, new Count())); taskCreatedSensor.add(metrics.metricName("task-created-total", "stream-metrics", "The total number of newly created tasks", tagMap()), new Total()); - tasksClosedSensor = threadLevelSensor("task-closed", Sensor.RecordingLevel.INFO); - tasksClosedSensor.add(metrics.metricName("task-closed-rate", group, "The average per-second number of closed tasks", tagMap()), new Rate(TimeUnit.SECONDS, new Count())); - tasksClosedSensor.add(metrics.metricName("task-closed-total", group, "The total number of closed tasks", tagMap()), new Total()); + taskClosedSensor = threadLevelSensor("task-closed", Sensor.RecordingLevel.INFO); + taskClosedSensor.add(metrics.metricName("task-closed-rate", group, "The average per-second number of closed tasks", tagMap()), new Rate(TimeUnit.SECONDS, new Count())); + taskClosedSensor.add(metrics.metricName("task-closed-total", group, "The total number of closed tasks", tagMap()), new Total()); } } From 69f22626c648138abf382ddf3b45891c74b635d6 Mon Sep 17 00:00:00 2001 From: Sam Lendle Date: Mon, 6 Aug 2018 15:34:58 -0700 Subject: [PATCH 4/7] Missed some checkstyle errors --- .../kafka/streams/processor/internals/StreamThread.java | 2 +- .../streams/processor/internals/metrics/CumulativeCount.java | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index 1d31ebb8f5d64..ed38018755029 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -514,7 +514,7 @@ StandbyTask createTask(final Consumer consumer, @Override public void close() { - streamsMetrics.taskClosedSensor.record(); + streamsMetrics.taskClosedSensor.record(); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/CumulativeCount.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/CumulativeCount.java index 54e698324dd75..2c12c2b6e9d41 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/CumulativeCount.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/metrics/CumulativeCount.java @@ -27,12 +27,12 @@ public class CumulativeCount implements MeasurableStat { private double count = 0.0; @Override - public void record(MetricConfig config, double value, long timeMs) { + public void record(final MetricConfig config, final double value, final long timeMs) { count += 1; } @Override - public double measure(MetricConfig config, long now) { + public double measure(final MetricConfig config, final long now) { return count; } } From 29ebacb8c8c62274fb5cbc138b57406b54ef9d04 Mon Sep 17 00:00:00 2001 From: Sam Lendle Date: Tue, 7 Aug 2018 17:00:40 -0700 Subject: [PATCH 5/7] Use small timeWindow instead of eventWindow in testTotalMetricDoesntDecrease --- .../streams/processor/internals/StreamsMetricsImplTest.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) 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/StreamsMetricsImplTest.java index d5ef826959ee8..9ea813908300f 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/StreamsMetricsImplTest.java @@ -26,6 +26,8 @@ import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl; import org.junit.Test; +import java.util.concurrent.TimeUnit; + import static org.junit.Assert.assertEquals; public class StreamsMetricsImplTest { @@ -103,8 +105,8 @@ public void testThroughputMetrics() { @Test public void testTotalMetricDoesntDecrease() { - final MockTime time = new MockTime(); - final MetricConfig config = new MetricConfig().eventWindow(1).samples(2); + final MockTime time = new MockTime(1); + final MetricConfig config = new MetricConfig().timeWindow(1, TimeUnit.MILLISECONDS); final Metrics metrics = new Metrics(config, time); final StreamsMetricsImpl streamsMetrics = new StreamsMetricsImpl(metrics, ""); From 0847bb3b6f99c3892b375c42d62ede8431777fbc Mon Sep 17 00:00:00 2001 From: Sam Lendle Date: Tue, 7 Aug 2018 17:02:32 -0700 Subject: [PATCH 6/7] remove calls to taskClosedSensor.record, it is called in the StreamTask class --- .../kafka/streams/processor/internals/StreamThread.java | 6 ------ 1 file changed, 6 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index ed38018755029..28cedbe41976b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -455,7 +455,6 @@ private Producer createProducer(final TaskId id) { @Override public void close() { - streamsMetrics.taskClosedSensor.record(); if (threadProducer != null) { try { threadProducer.close(); @@ -511,11 +510,6 @@ StandbyTask createTask(final Consumer consumer, return null; } } - - @Override - public void close() { - streamsMetrics.taskClosedSensor.record(); - } } static class StreamsMetricsThreadImpl extends StreamsMetricsImpl { From 5e98c3d1847cb900f4ab0b64ac8c99baac9b7572 Mon Sep 17 00:00:00 2001 From: Sam Lendle Date: Tue, 21 Aug 2018 12:50:15 -0700 Subject: [PATCH 7/7] Round metric to int in test --- .../streams/processor/internals/StreamsMetricsImplTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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/StreamsMetricsImplTest.java index 9ea813908300f..7ce27b4b6d668 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/StreamsMetricsImplTest.java @@ -135,7 +135,7 @@ public void testTotalMetricDoesntDecrease() { final KafkaMetric totalMetric = metrics.metric(totalMetricName); for (int i = 0; i < 10; i++) { - assertEquals(i, totalMetric.measurable().measure(config, time.milliseconds()), 0.0001); + assertEquals(i, Math.round(totalMetric.measurable().measure(config, time.milliseconds()))); sensor.record(latency, time.milliseconds()); }