From a9c05350da47a674f7dd5e9d8fddc5859ae116de Mon Sep 17 00:00:00 2001 From: Rajini Sivaram Date: Thu, 9 Aug 2018 11:56:44 +0100 Subject: [PATCH 1/3] KAFKA-7261: Fix request total metric to count requests instead of bytes --- .../kafka/common/metrics/stats/Meter.java | 4 +++- .../kafka/common/network/NetworkReceive.java | 8 +++++++ .../apache/kafka/common/network/Selector.java | 2 +- .../kafka/common/network/NioEchoServer.java | 2 +- .../common/network/SslTransportLayerTest.java | 21 +++++++++++++++++++ 5 files changed, 34 insertions(+), 3 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/metrics/stats/Meter.java b/clients/src/main/java/org/apache/kafka/common/metrics/stats/Meter.java index 09263cecae89c..a25a4f0844ac5 100644 --- a/clients/src/main/java/org/apache/kafka/common/metrics/stats/Meter.java +++ b/clients/src/main/java/org/apache/kafka/common/metrics/stats/Meter.java @@ -77,6 +77,8 @@ public List stats() { @Override public void record(MetricConfig config, double value, long timeMs) { rate.record(config, value, timeMs); - total.record(config, value, timeMs); + // Total metrics with Count stat should record 1.0 (as recorded in the count) + double totalValue = (rate.stat instanceof Count) ? 1.0 : value; + total.record(config, totalValue, timeMs); } } diff --git a/clients/src/main/java/org/apache/kafka/common/network/NetworkReceive.java b/clients/src/main/java/org/apache/kafka/common/network/NetworkReceive.java index 355233125bc9b..55354ac8d6417 100644 --- a/clients/src/main/java/org/apache/kafka/common/network/NetworkReceive.java +++ b/clients/src/main/java/org/apache/kafka/common/network/NetworkReceive.java @@ -146,4 +146,12 @@ public ByteBuffer payload() { return this.buffer; } + /** + * Returns the total size of the receive including payload and size buffer + * for use in metrics. This is consistent with {@link NetworkSend#size()} + */ + public int size() { + return payload().limit() + size.limit(); + } + } diff --git a/clients/src/main/java/org/apache/kafka/common/network/Selector.java b/clients/src/main/java/org/apache/kafka/common/network/Selector.java index 8ca7fff381a0b..7e32509933e55 100644 --- a/clients/src/main/java/org/apache/kafka/common/network/Selector.java +++ b/clients/src/main/java/org/apache/kafka/common/network/Selector.java @@ -862,7 +862,7 @@ private void addToCompletedReceives() { private void addToCompletedReceives(KafkaChannel channel, Deque stagedDeque) { NetworkReceive networkReceive = stagedDeque.poll(); this.completedReceives.add(networkReceive); - this.sensors.recordBytesReceived(channel.id(), networkReceive.payload().limit()); + this.sensors.recordBytesReceived(channel.id(), networkReceive.size()); } // only for testing diff --git a/clients/src/test/java/org/apache/kafka/common/network/NioEchoServer.java b/clients/src/test/java/org/apache/kafka/common/network/NioEchoServer.java index 53f9d95a55a67..64b7e4e679225 100644 --- a/clients/src/test/java/org/apache/kafka/common/network/NioEchoServer.java +++ b/clients/src/test/java/org/apache/kafka/common/network/NioEchoServer.java @@ -120,7 +120,7 @@ public void verifyAuthenticationMetrics(int successfulAuthentications, final int waitForMetric("failed-authentication", failedAuthentications); } - private void waitForMetric(String name, final double expectedValue) throws InterruptedException { + public void waitForMetric(String name, final double expectedValue) throws InterruptedException { final String totalName = name + "-total"; final String rateName = name + "-rate"; if (expectedValue == 0.0) { diff --git a/clients/src/test/java/org/apache/kafka/common/network/SslTransportLayerTest.java b/clients/src/test/java/org/apache/kafka/common/network/SslTransportLayerTest.java index 6aef2f7eda6f2..d70a448df228f 100644 --- a/clients/src/test/java/org/apache/kafka/common/network/SslTransportLayerTest.java +++ b/clients/src/test/java/org/apache/kafka/common/network/SslTransportLayerTest.java @@ -556,6 +556,27 @@ public void testUnsupportedCiphers() throws Exception { server.verifyAuthenticationMetrics(0, 1); } + @Test + public void testServerRequestMetrics() throws Exception { + String node = "0"; + server = createEchoServer(SecurityProtocol.SSL); + createSelector(sslClientConfigs, 16384, 16384, 16384); + InetSocketAddress addr = new InetSocketAddress("localhost", server.port()); + selector.connect(node, addr, 102400, 102400); + NetworkTestUtils.waitForChannelReady(selector, node); + int messageSize = 1024 * 1024; + String message = TestUtils.randomString(messageSize); + selector.send(new NetworkSend(node, ByteBuffer.wrap(message.getBytes()))); + while (selector.completedReceives().isEmpty()) { + selector.poll(100L); + } + int totalBytes = messageSize + 4; // including 4-byte size + server.waitForMetric("incoming-byte", totalBytes); + server.waitForMetric("outgoing-byte", totalBytes); + server.waitForMetric("request", 1); + server.waitForMetric("response", 1); + } + /** * selector.poll() should be able to fetch more data than netReadBuffer from the socket. */ From 9b382ef8d00eb4ac5a8a0a7e61b0a65f84910bfb Mon Sep 17 00:00:00 2001 From: Rajini Sivaram Date: Fri, 10 Aug 2018 11:28:42 +0100 Subject: [PATCH 2/3] Add unit test --- .../apache/kafka/common/metrics/MetricsTest.java | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/clients/src/test/java/org/apache/kafka/common/metrics/MetricsTest.java b/clients/src/test/java/org/apache/kafka/common/metrics/MetricsTest.java index 59bc84e40decf..c7ed236be06c9 100644 --- a/clients/src/test/java/org/apache/kafka/common/metrics/MetricsTest.java +++ b/clients/src/test/java/org/apache/kafka/common/metrics/MetricsTest.java @@ -465,8 +465,12 @@ public void testRateWindowing() throws Exception { Sensor s = metrics.sensor("test.sensor", cfg); MetricName rateMetricName = metrics.metricName("test.rate", "grp1"); MetricName totalMetricName = metrics.metricName("test.total", "grp1"); + MetricName countRateMetricName = metrics.metricName("test.count.rate", "grp1"); + MetricName countTotalMetricName = metrics.metricName("test.count.total", "grp1"); s.add(new Meter(TimeUnit.SECONDS, rateMetricName, totalMetricName)); + s.add(new Meter(TimeUnit.SECONDS, new Count(), countRateMetricName, countTotalMetricName)); KafkaMetric totalMetric = metrics.metrics().get(metrics.metricName("test.total", "grp1")); + KafkaMetric countTotalMetric = metrics.metrics().get(metrics.metricName("test.count.total", "grp1")); int sum = 0; int count = cfg.samples() - 1; @@ -485,10 +489,20 @@ public void testRateWindowing() throws Exception { double elapsedSecs = (cfg.timeWindowMs() * (cfg.samples() - 1) + cfg.timeWindowMs() / 2) / 1000.0; KafkaMetric rateMetric = metrics.metrics().get(metrics.metricName("test.rate", "grp1")); + KafkaMetric countRateMetric = metrics.metrics().get(metrics.metricName("test.count.rate", "grp1")); assertEquals("Rate(0...2) = 2.666", sum / elapsedSecs, rateMetric.value(), EPS); + assertEquals("Count rate(0...2) = 0.02666", count / elapsedSecs, countRateMetric.value(), EPS); assertEquals("Elapsed Time = 75 seconds", elapsedSecs, ((Rate) rateMetric.measurable()).windowSize(cfg, time.milliseconds()) / 1000, EPS); assertEquals(sum, totalMetric.value(), EPS); + assertEquals(count, countTotalMetric.value(), EPS); + + // Verify that rates are expired, but total is cumulative + time.sleep(cfg.timeWindowMs() * cfg.samples()); + assertEquals(0, rateMetric.value(), EPS); + assertEquals(0, countRateMetric.value(), EPS); + assertEquals(sum, totalMetric.value(), EPS); + assertEquals(count, countTotalMetric.value(), EPS); } public static class ConstantMeasurable implements Measurable { From 40db0de2e042738420750f71b15c1218c758bf84 Mon Sep 17 00:00:00 2001 From: Rajini Sivaram Date: Fri, 10 Aug 2018 21:37:48 +0100 Subject: [PATCH 3/3] Address review comments --- .../java/org/apache/kafka/common/metrics/stats/Meter.java | 3 +++ .../java/org/apache/kafka/common/metrics/MetricsTest.java | 8 ++++---- 2 files changed, 7 insertions(+), 4 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/metrics/stats/Meter.java b/clients/src/main/java/org/apache/kafka/common/metrics/stats/Meter.java index a25a4f0844ac5..91d4461d2b52f 100644 --- a/clients/src/main/java/org/apache/kafka/common/metrics/stats/Meter.java +++ b/clients/src/main/java/org/apache/kafka/common/metrics/stats/Meter.java @@ -61,6 +61,9 @@ public Meter(SampledStat rateStat, MetricName rateMetricName, MetricName totalMe * Construct a Meter with provided time unit and provided {@link SampledStat} stats for Rate */ public Meter(TimeUnit unit, SampledStat rateStat, MetricName rateMetricName, MetricName totalMetricName) { + if (!(rateStat instanceof SampledTotal) && !(rateStat instanceof Count)) { + throw new IllegalArgumentException("Meter is supported only for SampledTotal and Count"); + } this.total = new Total(); this.rate = new Rate(unit, rateStat); this.rateMetricName = rateMetricName; diff --git a/clients/src/test/java/org/apache/kafka/common/metrics/MetricsTest.java b/clients/src/test/java/org/apache/kafka/common/metrics/MetricsTest.java index c7ed236be06c9..5c75d03b0e642 100644 --- a/clients/src/test/java/org/apache/kafka/common/metrics/MetricsTest.java +++ b/clients/src/test/java/org/apache/kafka/common/metrics/MetricsTest.java @@ -469,8 +469,8 @@ public void testRateWindowing() throws Exception { MetricName countTotalMetricName = metrics.metricName("test.count.total", "grp1"); s.add(new Meter(TimeUnit.SECONDS, rateMetricName, totalMetricName)); s.add(new Meter(TimeUnit.SECONDS, new Count(), countRateMetricName, countTotalMetricName)); - KafkaMetric totalMetric = metrics.metrics().get(metrics.metricName("test.total", "grp1")); - KafkaMetric countTotalMetric = metrics.metrics().get(metrics.metricName("test.count.total", "grp1")); + KafkaMetric totalMetric = metrics.metrics().get(totalMetricName); + KafkaMetric countTotalMetric = metrics.metrics().get(countTotalMetricName); int sum = 0; int count = cfg.samples() - 1; @@ -488,8 +488,8 @@ public void testRateWindowing() throws Exception { // prior to any time passing double elapsedSecs = (cfg.timeWindowMs() * (cfg.samples() - 1) + cfg.timeWindowMs() / 2) / 1000.0; - KafkaMetric rateMetric = metrics.metrics().get(metrics.metricName("test.rate", "grp1")); - KafkaMetric countRateMetric = metrics.metrics().get(metrics.metricName("test.count.rate", "grp1")); + KafkaMetric rateMetric = metrics.metrics().get(rateMetricName); + KafkaMetric countRateMetric = metrics.metrics().get(countRateMetricName); assertEquals("Rate(0...2) = 2.666", sum / elapsedSecs, rateMetric.value(), EPS); assertEquals("Count rate(0...2) = 0.02666", count / elapsedSecs, countRateMetric.value(), EPS); assertEquals("Elapsed Time = 75 seconds", elapsedSecs,