diff --git a/docs/_docs/monitoring-metrics/new-metrics.adoc b/docs/_docs/monitoring-metrics/new-metrics.adoc index 7b749d7d212b3..23ce5e1a9e96e 100644 --- a/docs/_docs/monitoring-metrics/new-metrics.adoc +++ b/docs/_docs/monitoring-metrics/new-metrics.adoc @@ -520,4 +520,6 @@ Register name: `sql.queries.user` |success| long | The number of succesfully executed SQL queries. |failed| long | The number of failed SQL queries (including canceled). |canceled| long | The number of canceled SQL queries. +|resultSetSizeHistogram| histogram | Histogram of fetched result set sizes for SQL queries. +|maxResultSetSize| max value | Maximum fetched result set size for SQL queries. |=== diff --git a/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/exec/ExecutionServiceImpl.java b/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/exec/ExecutionServiceImpl.java index 5ab4e94c77330..3824416e78be7 100644 --- a/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/exec/ExecutionServiceImpl.java +++ b/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/exec/ExecutionServiceImpl.java @@ -1270,6 +1270,8 @@ private ListFieldsQueryCursor mapAndExecutePlan( ); } + ctx.query().runningQueryManager().onFullyFetched(resultSetChecker.fetchedSize()); + resultSetChecker.checkOnClose(); }; diff --git a/modules/calcite/src/test/java/org/apache/ignite/internal/processors/query/calcite/integration/ResultSetSizeMetricsTest.java b/modules/calcite/src/test/java/org/apache/ignite/internal/processors/query/calcite/integration/ResultSetSizeMetricsTest.java new file mode 100644 index 0000000000000..95430e21e5e01 --- /dev/null +++ b/modules/calcite/src/test/java/org/apache/ignite/internal/processors/query/calcite/integration/ResultSetSizeMetricsTest.java @@ -0,0 +1,130 @@ +/* + * 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.ignite.internal.processors.query.calcite.integration; + +import org.apache.ignite.IgniteCache; +import org.apache.ignite.configuration.CacheConfiguration; +import org.apache.ignite.configuration.IgniteConfiguration; +import org.apache.ignite.internal.IgniteEx; +import org.apache.ignite.internal.processors.metric.MetricRegistryImpl; +import org.apache.ignite.internal.processors.metric.impl.HistogramMetricImpl; +import org.apache.ignite.internal.processors.metric.impl.MaxValueMetric; +import org.junit.Test; + +import static org.apache.ignite.internal.processors.query.running.RunningQueryManager.SQL_USER_QUERIES_REG_NAME; + +/** + * Tests for result set size histogram and max result set size metrics. + */ +public class ResultSetSizeMetricsTest extends AbstractMultiEngineIntegrationTest { + /** */ + private static final int FILL_SIZE = 1000; + + /** */ + @Override protected void afterTest() throws Exception { + stopAllGrids(); + } + + /** {@inheritDoc} */ + @Override protected IgniteConfiguration getConfiguration(String igniteInstanceName) throws Exception { + IgniteConfiguration cfg = super.getConfiguration(igniteInstanceName); + + cfg.setCacheConfiguration(new CacheConfiguration<>(DEFAULT_CACHE_NAME) + .setIndexedTypes(Integer.class, Integer.class)); + + return cfg; + } + + /** */ + @Test + public void testResultSetSizeMetrics() throws Exception { + IgniteEx initNode = startGrids(nodeCount()); + + IgniteCache cache = initNode.cache(DEFAULT_CACHE_NAME); + + for (int i = 0; i < FILL_SIZE; i++) + cache.put(i, i); + + // Execute simple queries with result set sizes: 0, 1, 5, 50, 500. + for (int limit : new int[] {0, 1, 5, 50, 500}) + sql(initNode, "SELECT _key FROM \"" + DEFAULT_CACHE_NAME + "\".Integer WHERE _key < ?", limit); + + // Execute queries with aggregation (different reducers on h2) with result set sizes: 10, 100. + for (int limit : new int[] {10, 100}) + sql(initNode, "SELECT DISTINCT _key FROM \"" + DEFAULT_CACHE_NAME + "\".Integer WHERE _key < ?", limit); + + // Verify histogram on the initiating node. + // Bounds: {0, 1, 10, 100, 1_000, 10_000, 100_000, 1_000_000} + // Bucket 0: x <= 0 -> 1 (size 0) + // Bucket 1: x <= 1 -> 1 (size 1) + // Bucket 2: x <= 10 -> 2 (sizes 5, 10) + // Bucket 3: x <= 100 -> 2 (sizes 50, 100) + // Bucket 4: x <= 1000 -> 1 (size 500) + // Buckets 5-8 -> 0 + long[] values = resultSetSizeHistogram(initNode).value(); + + assertEquals(1, values[0]); + assertEquals(1, values[1]); + assertEquals(2, values[2]); + assertEquals(2, values[3]); + assertEquals(1, values[4]); + + for (int i = 5; i < values.length; i++) + assertEquals(0, values[i]); + + // Verify max value on the initiating node. + assertEquals(500L, resultSetSizeMax(initNode).value()); + + // Verify all other server nodes have zero metrics. + for (int i = 0; i < nodeCount(); i++) { + IgniteEx node = grid(i); + + if (node == initNode) + continue; + + long[] nodeVals = resultSetSizeHistogram(node).value(); + + for (long v : nodeVals) + assertEquals(0, v); + + assertEquals("Expected max value 0 on node [" + node.name() + "]", + 0L, resultSetSizeMax(node).value()); + } + } + + /** */ + private HistogramMetricImpl resultSetSizeHistogram(IgniteEx ignite) { + MetricRegistryImpl mreg = ignite.context().metric().registry(SQL_USER_QUERIES_REG_NAME); + + HistogramMetricImpl hist = mreg.findMetric("resultSetSizeHistogram"); + + assertNotNull(hist); + + return hist; + } + + /** */ + private MaxValueMetric resultSetSizeMax(IgniteEx ignite) { + MetricRegistryImpl mreg = ignite.context().metric().registry(SQL_USER_QUERIES_REG_NAME); + + MaxValueMetric max = mreg.findMetric("maxResultSetSize"); + + assertNotNull(max); + + return max; + } +} diff --git a/modules/calcite/src/test/java/org/apache/ignite/testsuites/IntegrationTestSuite.java b/modules/calcite/src/test/java/org/apache/ignite/testsuites/IntegrationTestSuite.java index 2a2cb856731a3..c9605a431cd29 100644 --- a/modules/calcite/src/test/java/org/apache/ignite/testsuites/IntegrationTestSuite.java +++ b/modules/calcite/src/test/java/org/apache/ignite/testsuites/IntegrationTestSuite.java @@ -68,6 +68,7 @@ import org.apache.ignite.internal.processors.query.calcite.integration.QueryMetadataIntegrationTest; import org.apache.ignite.internal.processors.query.calcite.integration.QueryWithPartitionsIntegrationTest; import org.apache.ignite.internal.processors.query.calcite.integration.RecursiveCteIntegrationTest; +import org.apache.ignite.internal.processors.query.calcite.integration.ResultSetSizeMetricsTest; import org.apache.ignite.internal.processors.query.calcite.integration.RunningQueriesIntegrationTest; import org.apache.ignite.internal.processors.query.calcite.integration.ScalarInIntegrationTest; import org.apache.ignite.internal.processors.query.calcite.integration.SelectByKeyFieldTest; @@ -197,6 +198,7 @@ SystemColumnsScanTest.class, BulkOperationDeadlockIntegrationTest.class, SelectForUpdateIntegrationTest.class, + ResultSetSizeMetricsTest.class, }) public class IntegrationTestSuite { } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/query/running/RunningQueryManager.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/query/running/RunningQueryManager.java index 0d1229add8a6f..3990d21b353ab 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/query/running/RunningQueryManager.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/query/running/RunningQueryManager.java @@ -47,7 +47,9 @@ import org.apache.ignite.internal.processors.closure.GridClosureProcessor; import org.apache.ignite.internal.processors.metric.MetricRegistryImpl; import org.apache.ignite.internal.processors.metric.impl.AtomicLongMetric; +import org.apache.ignite.internal.processors.metric.impl.HistogramMetricImpl; import org.apache.ignite.internal.processors.metric.impl.LongAdderMetric; +import org.apache.ignite.internal.processors.metric.impl.MaxValueMetric; import org.apache.ignite.internal.processors.query.GridQueryCancel; import org.apache.ignite.internal.processors.query.GridQueryFinishedInfo; import org.apache.ignite.internal.processors.query.GridQueryStartedInfo; @@ -141,6 +143,12 @@ public class RunningQueryManager { */ private final AtomicLongMetric canceledQrsCnt; + /** Histogram of result set sizes for SQL queries. */ + private final HistogramMetricImpl resultSetSizeHistogram; + + /** Maximum result set size for SQL queries. */ + private final MaxValueMetric maxResultSetSize; + /** Kernal context. */ private final GridKernalContext ctx; @@ -231,6 +239,13 @@ public RunningQueryManager(GridKernalContext ctx) { canceledQrsCnt = userMetrics.longMetric("canceled", "Number of canceled queries that have been started " + "on this node. This metric number included in the general 'failed' metric."); + + resultSetSizeHistogram = userMetrics.histogram("resultSetSizeHistogram", + new long[] {0, 1, 10, 100, 1_000, 10_000, 100_000, 1_000_000}, + "Histogram of result set sizes for SQL queries."); + + maxResultSetSize = userMetrics.maxValueMetric("maxResultSetSize", + "Maximum result set size for SQL queries.", 60_000L, 5); } /** */ @@ -269,6 +284,16 @@ public void start(GridSpinBusyLock busyLock) { }, EventType.EVT_NODE_FAILED, EventType.EVT_NODE_LEFT); } + /** + * Called when a result set is fully fetched. Increments result set size metrics. + * + * @param size Result set size (number of fetched rows). + */ + public void onFullyFetched(long size) { + resultSetSizeHistogram.value(size); + maxResultSetSize.update(size); + } + /** * Registers running query and returns an id associated with the query. * diff --git a/modules/indexing/src/main/java/org/apache/ignite/internal/processors/query/h2/H2ResultSetIterator.java b/modules/indexing/src/main/java/org/apache/ignite/internal/processors/query/h2/H2ResultSetIterator.java index 379b215367713..729fd8b2ef810 100644 --- a/modules/indexing/src/main/java/org/apache/ignite/internal/processors/query/h2/H2ResultSetIterator.java +++ b/modules/indexing/src/main/java/org/apache/ignite/internal/processors/query/h2/H2ResultSetIterator.java @@ -320,6 +320,8 @@ public void onClose() throws IgniteCheckedException { try { resultSetChecker.checkOnClose(); + h2.runningQueryManager().onFullyFetched(resultSetChecker.fetchedSize()); + PerformanceStatisticsProcessor perfStat = ctx.performanceStatistics(); if (perfStat.enabled() && resultSetChecker.fetchedSize() > 0) { diff --git a/modules/indexing/src/main/java/org/apache/ignite/internal/processors/query/h2/twostep/GridReduceQueryExecutor.java b/modules/indexing/src/main/java/org/apache/ignite/internal/processors/query/h2/twostep/GridReduceQueryExecutor.java index a5c896e39fbee..d3695c2c47986 100644 --- a/modules/indexing/src/main/java/org/apache/ignite/internal/processors/query/h2/twostep/GridReduceQueryExecutor.java +++ b/modules/indexing/src/main/java/org/apache/ignite/internal/processors/query/h2/twostep/GridReduceQueryExecutor.java @@ -1400,4 +1400,9 @@ private List prepareMapQueryForSinglePartition(GridCacheTwoSt return Collections.singletonList(originalQry); } + + /** */ + IgniteH2Indexing h2() { + return h2; + } } diff --git a/modules/indexing/src/main/java/org/apache/ignite/internal/processors/query/h2/twostep/ReduceIndexIterator.java b/modules/indexing/src/main/java/org/apache/ignite/internal/processors/query/h2/twostep/ReduceIndexIterator.java index 5924b7ce8263a..74ceceb1b31b5 100644 --- a/modules/indexing/src/main/java/org/apache/ignite/internal/processors/query/h2/twostep/ReduceIndexIterator.java +++ b/modules/indexing/src/main/java/org/apache/ignite/internal/processors/query/h2/twostep/ReduceIndexIterator.java @@ -57,6 +57,9 @@ public class ReduceIndexIterator implements Iterator>, AutoCloseable { /** Whether remote resources were released. */ private boolean released; + /** Fetched rows count. */ + private long fetched; + /** * Constructor. * @@ -95,6 +98,8 @@ public ReduceIndexIterator(GridReduceQueryExecutor rdcExec, if (res == null) throw new NoSuchElementException(); + fetched++; + advance(); return res; @@ -156,6 +161,7 @@ private void releaseIfNeeded() { if (!released) { try { rdcExec.releaseRemoteResources(nodes, run, qryReqId, distributedJoins); + rdcExec.h2().runningQueryManager().onFullyFetched(fetched); } finally { released = true;