From 152f64e11ccdf0819de22c7064fb8e90855b059c Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Wed, 30 Sep 2026 15:44:14 +0800 Subject: [PATCH] Optimize tag index timeseries lookup (#18753) * Optimize tag index timeseries lookup * Add opt-in tag index performance test (cherry picked from commit 0c9a5ee09164b8cd35beb34e4dad75af54c6feae) --- .../schemaregion/tag/TagManager.java | 54 ++-- .../tag/TagManagerPerformanceTest.java | 258 ++++++++++++++++++ .../schemaregion/tag/TagManagerTest.java | 43 +++ 3 files changed, 328 insertions(+), 27 deletions(-) create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManagerPerformanceTest.java diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManager.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManager.java index 805421bb50485..1d8919dadcc7c 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManager.java @@ -244,7 +244,8 @@ private boolean containsIndex(String tagKey, String tagValue) { return tagValueMap != null && tagValueMap.containsKey(tagValue); } - private List> getMatchedTimeseriesInIndex(TagFilter tagFilter) { + List> getMatchedTimeseriesInIndex( + final TagFilter tagFilter, final PartialPath pathPattern, final boolean isPrefixMatch) { Map>> value2Node = tagIndex.get(tagFilter.getKey()); if (value2Node == null || value2Node.isEmpty()) { return Collections.emptyList(); @@ -262,19 +263,20 @@ private List> getMatchedTimeseriesInIndex(TagFilter tagFilt } } } else { - for (Map.Entry>> entry : value2Node.entrySet()) { - if (entry.getKey() == null || entry.getValue() == null) { - continue; - } - String tagValue = entry.getKey(); - if (tagFilter.getValue().equals(tagValue)) { - allMatchedNodes.addAll(entry.getValue()); - } + final Set> matchedNodes = value2Node.get(tagFilter.getValue()); + if (matchedNodes != null) { + allMatchedNodes.addAll(matchedNodes); } } - // we just sort them by the alphabetical order + + // Filter by path before sorting to avoid sorting irrelevant measurements. allMatchedNodes = allMatchedNodes.stream() + .filter( + node -> + isPrefixMatch + ? pathPattern.prefixMatchFullPath(node.getPartialPath()) + : pathPattern.matchFullPath(node.getPartialPath())) .sorted(Comparator.comparing(IMNode::getFullPath)) .collect(toList()); @@ -287,11 +289,13 @@ public ISchemaReader getTimeSeriesReaderWithIndex( SchemaFilter schemaFilter = plan.getSchemaFilter(); // currently, only one TagFilter is supported // all IMeasurementMNode in allMatchedNodes satisfied TagFilter + PartialPath pathPattern = plan.getPath(); Iterator> allMatchedNodes = getMatchedTimeseriesInIndex( - (TagFilter) SchemaFilter.extract(schemaFilter, SchemaFilterType.TAGS_FILTER).get(0)) + (TagFilter) SchemaFilter.extract(schemaFilter, SchemaFilterType.TAGS_FILTER).get(0), + pathPattern, + plan.isPrefixMatch()) .iterator(); - PartialPath pathPattern = plan.getPath(); SchemaIterator schemaIterator = new SchemaIterator() { private ITimeSeriesSchemaInfo nextMatched; @@ -323,21 +327,17 @@ private void getNext() throws IOException { nextMatched = null; while (allMatchedNodes.hasNext()) { IMeasurementMNode node = allMatchedNodes.next(); - if (plan.isPrefixMatch() - ? pathPattern.prefixMatchFullPath(node.getPartialPath()) - : pathPattern.matchFullPath(node.getPartialPath())) { - Pair, Map> tagAndAttributePair = - readTagFile(node.getOffset()); - nextMatched = - new ShowTimeSeriesResult( - node.getFullPath(), - node.getAlias(), - node.getSchema(), - tagAndAttributePair.left, - tagAndAttributePair.right, - node.getParent().getAsDeviceMNode().isAligned()); - break; - } + Pair, Map> tagAndAttributePair = + readTagFile(node.getOffset()); + nextMatched = + new ShowTimeSeriesResult( + node.getFullPath(), + node.getAlias(), + node.getSchema(), + tagAndAttributePair.left, + tagAndAttributePair.right, + node.getParent().getAsDeviceMNode().isAligned()); + break; } } diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManagerPerformanceTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManagerPerformanceTest.java new file mode 100644 index 0000000000000..58fb783c73c33 --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManagerPerformanceTest.java @@ -0,0 +1,258 @@ +/* + * 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.iotdb.db.schemaengine.schemaregion.tag; + +import org.apache.iotdb.commons.path.PartialPath; +import org.apache.iotdb.commons.schema.filter.impl.TagFilter; +import org.apache.iotdb.commons.schema.node.IMNode; +import org.apache.iotdb.commons.schema.node.role.IDeviceMNode; +import org.apache.iotdb.commons.schema.node.role.IMeasurementMNode; +import org.apache.iotdb.commons.schema.node.utils.IMNodeFactory; +import org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.mem.mnode.IMemMNode; +import org.apache.iotdb.db.schemaengine.schemaregion.mtree.loader.MNodeFactoryLoader; +import org.apache.iotdb.db.utils.ManualPerformanceTestUtils; +import org.apache.iotdb.db.utils.ManualPerformanceTestUtils.Measurement; +import org.apache.iotdb.db.utils.ManualPerformanceTestUtils.Summary; + +import org.apache.commons.io.FileUtils; +import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.file.metadata.enums.CompressionType; +import org.apache.tsfile.file.metadata.enums.TSEncoding; +import org.apache.tsfile.write.schema.MeasurementSchema; +import org.junit.Assert; +import org.junit.Assume; +import org.junit.Test; + +import java.io.File; +import java.lang.reflect.Field; +import java.nio.file.Files; +import java.util.ArrayList; +import java.util.Collections; +import java.util.Comparator; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Set; +import java.util.stream.Collectors; + +public class TagManagerPerformanceTest { + private static final String PREFIX = "iotdb.tag.index.perf."; + private static volatile List> benchmarkBlackhole; + + @Test + public void compareExactLookupAndPathFiltering() throws Exception { + Assume.assumeTrue( + "Manual performance UT: enable with -D" + PREFIX + "enabled=true", + Boolean.getBoolean(PREFIX + "enabled")); + Assume.assumeTrue( + "Current-thread CPU and allocation metrics are required.", + ManualPerformanceTestUtils.enableThreadMetrics()); + final int distinctValues = Integer.getInteger(PREFIX + "distinctValues", 10_000); + final int candidateNodes = Integer.getInteger(PREFIX + "candidateNodes", 2_000); + final int matchingNodes = Integer.getInteger(PREFIX + "matchingNodes", 100); + final int warmups = Integer.getInteger(PREFIX + "warmups", 30); + final int exactIterations = Integer.getInteger(PREFIX + "exactIterations", 50_000); + final int pathIterations = Integer.getInteger(PREFIX + "pathIterations", 1_000); + final int rounds = Integer.getInteger(PREFIX + "rounds", 5); + Assert.assertTrue( + distinctValues > 0 + && candidateNodes > 0 + && matchingNodes > 0 + && matchingNodes <= candidateNodes); + Assert.assertTrue(warmups >= 0 && exactIterations > 0 && pathIterations > 0 && rounds > 0); + + final File directory = Files.createTempDirectory("tag-index-performance").toFile(); + try { + final TagManager manager = new TagManager(directory.getAbsolutePath(), null); + try { + final IMNodeFactory factory = + MNodeFactoryLoader.getInstance().getMemMNodeIMNodeFactory(); + final IMemMNode root = factory.createInternalMNode(null, "root"); + final IDeviceMNode matchingDevice = + factory.createDeviceMNode(factory.createInternalMNode(root, "sg"), "d"); + final IDeviceMNode otherDevice = + factory.createDeviceMNode(factory.createInternalMNode(root, "other"), "d"); + final IMeasurementMNode shared = + newMeasurement(factory, matchingDevice, "single"); + for (int i = 0; i < distinctValues; i++) { + manager.addIndex("cardinality", "v" + i, shared); + } + for (int i = 0; i < candidateNodes; i++) { + manager.addIndex( + "selectivity", + "target", + newMeasurement(factory, i < matchingNodes ? matchingDevice : otherDevice, "s" + i)); + } + final Map>>> index = getIndex(manager); + final PartialPath pattern = new PartialPath("root.sg.**"); + System.out.printf( + Locale.ROOT, + "Tag index benchmark: distinctValues=%d, candidateNodes=%d, matchingNodes=%d, warmups=%d, exactIterations/round=%d, pathIterations/round=%d, rounds=%d%n", + distinctValues, + candidateNodes, + matchingNodes, + warmups, + exactIterations, + pathIterations, + rounds); + benchmark( + "exact value lookup", + manager, + index, + new TagFilter("cardinality", "v0", false), + pattern, + 1, + warmups, + exactIterations, + rounds); + benchmark( + "path filtering before sort", + manager, + index, + new TagFilter("selectivity", "target", false), + pattern, + matchingNodes, + warmups, + pathIterations, + rounds); + } finally { + manager.clear(); + } + } finally { + FileUtils.deleteDirectory(directory); + } + } + + private static IMeasurementMNode newMeasurement( + IMNodeFactory factory, IDeviceMNode parent, String name) { + return factory.createMeasurementMNode( + parent, + name, + new MeasurementSchema(name, TSDataType.INT64, TSEncoding.PLAIN, CompressionType.SNAPPY), + null); + } + + @SuppressWarnings("unchecked") + private static Map>>> getIndex(TagManager manager) + throws Exception { + final Field field = TagManager.class.getDeclaredField("tagIndex"); + field.setAccessible(true); + return (Map>>>) field.get(manager); + } + + // Reproduce the old value scan and sort, then apply the reader's path filter. + private static List> legacyLookup( + Map>>> index, + TagFilter filter, + PartialPath pattern) { + final Map>> values = index.get(filter.getKey()); + if (values == null || values.isEmpty()) { + return Collections.emptyList(); + } + final List> candidates = new ArrayList<>(); + for (Map.Entry>> entry : values.entrySet()) { + if (entry.getKey() != null + && entry.getValue() != null + && filter.getValue().equals(entry.getKey())) { + candidates.addAll(entry.getValue()); + } + } + final List> sorted = + candidates.stream() + .sorted(Comparator.comparing(IMNode::getFullPath)) + .collect(Collectors.toList()); + sorted.removeIf(node -> !pattern.matchFullPath(node.getPartialPath())); + return sorted; + } + + private static void benchmark( + String label, + TagManager manager, + Map>>> index, + TagFilter filter, + PartialPath pattern, + int expectedMatches, + int warmups, + int iterations, + int rounds) { + final List> legacyResult = legacyLookup(index, filter, pattern); + final List> optimizedResult = + manager.getMatchedTimeseriesInIndex(filter, pattern, false); + Assert.assertEquals(expectedMatches, optimizedResult.size()); + Assert.assertEquals(legacyResult, optimizedResult); + final Runnable legacy = () -> benchmarkBlackhole = legacyLookup(index, filter, pattern); + final Runnable optimized = + () -> benchmarkBlackhole = manager.getMatchedTimeseriesInIndex(filter, pattern, false); + for (int i = 0; i < warmups; i++) { + if ((i & 1) == 0) { + legacy.run(); + optimized.run(); + } else { + optimized.run(); + legacy.run(); + } + } + final Measurement[] legacyMeasurements = new Measurement[rounds]; + final Measurement[] optimizedMeasurements = new Measurement[rounds]; + for (int i = 0; i < rounds; i++) { + if ((i & 1) == 0) { + legacyMeasurements[i] = ManualPerformanceTestUtils.measure(iterations, legacy); + optimizedMeasurements[i] = ManualPerformanceTestUtils.measure(iterations, optimized); + } else { + optimizedMeasurements[i] = ManualPerformanceTestUtils.measure(iterations, optimized); + legacyMeasurements[i] = ManualPerformanceTestUtils.measure(iterations, legacy); + } + } + final Summary oldSummary = ManualPerformanceTestUtils.summarize(legacyMeasurements, iterations); + final Summary newSummary = + ManualPerformanceTestUtils.summarize(optimizedMeasurements, iterations); + System.out.printf(Locale.ROOT, " %s (matches=%d):%n", label, expectedMatches); + printSummary("legacy", oldSummary); + printSummary("optimized", newSummary); + final double allocationReduction = + oldSummary.getAllocatedBytesPerOperation() == 0 + ? 0 + : (oldSummary.getAllocatedBytesPerOperation() + - newSummary.getAllocatedBytesPerOperation()) + * 100.0 + / oldSummary.getAllocatedBytesPerOperation(); + if (oldSummary.getCpuNanosPerOperation() > 0 && newSummary.getCpuNanosPerOperation() > 0) { + System.out.printf( + Locale.ROOT, + " CPU speedup=%.2fx, allocation reduction=%.1f%%%n", + oldSummary.getCpuNanosPerOperation() / newSummary.getCpuNanosPerOperation(), + allocationReduction); + } else { + System.out.printf( + Locale.ROOT, + " CPU speedup=n/a (timer resolution), allocation reduction=%.1f%%%n", + allocationReduction); + } + } + + private static void printSummary(String label, Summary summary) { + System.out.printf( + Locale.ROOT, + " %-9s CPU=%.3f us/query, allocated=%.1f B/query%n", + label, + summary.getCpuNanosPerOperation() / 1_000.0, + summary.getAllocatedBytesPerOperation()); + } +} diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManagerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManagerTest.java index 1d78b4a548360..51b6456eabb76 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManagerTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManagerTest.java @@ -18,6 +18,8 @@ */ package org.apache.iotdb.db.schemaengine.schemaregion.tag; +import org.apache.iotdb.commons.path.PartialPath; +import org.apache.iotdb.commons.schema.filter.impl.TagFilter; import org.apache.iotdb.commons.schema.node.role.IMeasurementMNode; import org.apache.iotdb.commons.schema.node.utils.IMNodeFactory; import org.apache.iotdb.db.schemaengine.rescon.MemSchemaEngineStatistics; @@ -34,11 +36,18 @@ import org.junit.After; import org.junit.Assert; import org.junit.Test; +import org.mockito.Mockito; import java.io.File; +import java.lang.reflect.Field; import java.nio.file.Files; import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashMap; +import java.util.HashSet; import java.util.List; +import java.util.Map; +import java.util.Set; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -188,6 +197,33 @@ public void concurrentAddAndRemoveIndexEventuallyReleasesAllMemory() throws Exce Assert.assertEquals(0, regionStatistics.getRegionMemoryUsage()); } + @Test + @SuppressWarnings("unchecked") + public void preciseQueryUsesDirectLookupAndFiltersPathBeforeSorting() throws Exception { + initTagManager(); + final IMeasurementMNode first = mockMeasurement("root.sg.d.s1"); + final IMeasurementMNode second = mockMeasurement("root.sg.d.s2"); + final IMeasurementMNode outside = mockMeasurement("root.other.d.s0"); + + final Map>> valueMap = Mockito.spy(new HashMap<>()); + valueMap.put("target", new HashSet<>(Arrays.asList(second, outside, first))); + valueMap.put("unrelated", new HashSet<>()); + + final Field tagIndexField = TagManager.class.getDeclaredField("tagIndex"); + tagIndexField.setAccessible(true); + final Map>>> tagIndex = + (Map>>>) tagIndexField.get(tagManager); + tagIndex.put("key", valueMap); + + final List> result = + tagManager.getMatchedTimeseriesInIndex( + new TagFilter("key", "target", false), new PartialPath("root.sg.**"), false); + + Assert.assertEquals(Arrays.asList(first, second), result); + Mockito.verify(valueMap).get("target"); + Mockito.verify(valueMap, Mockito.never()).entrySet(); + } + private void initTagManager() throws Exception { tempDir = Files.createTempDirectory("tag-manager").toFile(); regionStatistics = new MemSchemaRegionStatistics(0, new MemSchemaEngineStatistics()); @@ -205,6 +241,13 @@ private static IMeasurementMNode newMeasurementMNode(final String mea null); } + private static IMeasurementMNode mockMeasurement(final String path) throws Exception { + final IMeasurementMNode node = Mockito.mock(IMeasurementMNode.class); + Mockito.when(node.getFullPath()).thenReturn(path); + Mockito.when(node.getPartialPath()).thenReturn(new PartialPath(path)); + return node; + } + private static long indexMemory( final String tagKey, final String tagValue, final int measurementCount) { return RamUsageEstimator.sizeOf(tagKey)