Skip to content

Commit 7419a47

Browse files
authored
Optimize tag index timeseries lookup (#18753) (#18781)
* Optimize tag index timeseries lookup * Add opt-in tag index performance test (cherry picked from commit 0c9a5ee)
1 parent 2f36c37 commit 7419a47

3 files changed

Lines changed: 328 additions & 27 deletions

File tree

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/tag/TagManager.java‎

Lines changed: 27 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -244,7 +244,8 @@ private boolean containsIndex(String tagKey, String tagValue) {
244244
return tagValueMap != null && tagValueMap.containsKey(tagValue);
245245
}
246246

247-
private List<IMeasurementMNode<?>> getMatchedTimeseriesInIndex(TagFilter tagFilter) {
247+
List<IMeasurementMNode<?>> getMatchedTimeseriesInIndex(
248+
final TagFilter tagFilter, final PartialPath pathPattern, final boolean isPrefixMatch) {
248249
Map<String, Set<IMeasurementMNode<?>>> value2Node = tagIndex.get(tagFilter.getKey());
249250
if (value2Node == null || value2Node.isEmpty()) {
250251
return Collections.emptyList();
@@ -262,19 +263,20 @@ private List<IMeasurementMNode<?>> getMatchedTimeseriesInIndex(TagFilter tagFilt
262263
}
263264
}
264265
} else {
265-
for (Map.Entry<String, Set<IMeasurementMNode<?>>> entry : value2Node.entrySet()) {
266-
if (entry.getKey() == null || entry.getValue() == null) {
267-
continue;
268-
}
269-
String tagValue = entry.getKey();
270-
if (tagFilter.getValue().equals(tagValue)) {
271-
allMatchedNodes.addAll(entry.getValue());
272-
}
266+
final Set<IMeasurementMNode<?>> matchedNodes = value2Node.get(tagFilter.getValue());
267+
if (matchedNodes != null) {
268+
allMatchedNodes.addAll(matchedNodes);
273269
}
274270
}
275-
// we just sort them by the alphabetical order
271+
272+
// Filter by path before sorting to avoid sorting irrelevant measurements.
276273
allMatchedNodes =
277274
allMatchedNodes.stream()
275+
.filter(
276+
node ->
277+
isPrefixMatch
278+
? pathPattern.prefixMatchFullPath(node.getPartialPath())
279+
: pathPattern.matchFullPath(node.getPartialPath()))
278280
.sorted(Comparator.comparing(IMNode::getFullPath))
279281
.collect(toList());
280282

@@ -287,11 +289,13 @@ public ISchemaReader<ITimeSeriesSchemaInfo> getTimeSeriesReaderWithIndex(
287289
SchemaFilter schemaFilter = plan.getSchemaFilter();
288290
// currently, only one TagFilter is supported
289291
// all IMeasurementMNode in allMatchedNodes satisfied TagFilter
292+
PartialPath pathPattern = plan.getPath();
290293
Iterator<IMeasurementMNode<?>> allMatchedNodes =
291294
getMatchedTimeseriesInIndex(
292-
(TagFilter) SchemaFilter.extract(schemaFilter, SchemaFilterType.TAGS_FILTER).get(0))
295+
(TagFilter) SchemaFilter.extract(schemaFilter, SchemaFilterType.TAGS_FILTER).get(0),
296+
pathPattern,
297+
plan.isPrefixMatch())
293298
.iterator();
294-
PartialPath pathPattern = plan.getPath();
295299
SchemaIterator<ITimeSeriesSchemaInfo> schemaIterator =
296300
new SchemaIterator<ITimeSeriesSchemaInfo>() {
297301
private ITimeSeriesSchemaInfo nextMatched;
@@ -323,21 +327,17 @@ private void getNext() throws IOException {
323327
nextMatched = null;
324328
while (allMatchedNodes.hasNext()) {
325329
IMeasurementMNode<?> node = allMatchedNodes.next();
326-
if (plan.isPrefixMatch()
327-
? pathPattern.prefixMatchFullPath(node.getPartialPath())
328-
: pathPattern.matchFullPath(node.getPartialPath())) {
329-
Pair<Map<String, String>, Map<String, String>> tagAndAttributePair =
330-
readTagFile(node.getOffset());
331-
nextMatched =
332-
new ShowTimeSeriesResult(
333-
node.getFullPath(),
334-
node.getAlias(),
335-
node.getSchema(),
336-
tagAndAttributePair.left,
337-
tagAndAttributePair.right,
338-
node.getParent().getAsDeviceMNode().isAligned());
339-
break;
340-
}
330+
Pair<Map<String, String>, Map<String, String>> tagAndAttributePair =
331+
readTagFile(node.getOffset());
332+
nextMatched =
333+
new ShowTimeSeriesResult(
334+
node.getFullPath(),
335+
node.getAlias(),
336+
node.getSchema(),
337+
tagAndAttributePair.left,
338+
tagAndAttributePair.right,
339+
node.getParent().getAsDeviceMNode().isAligned());
340+
break;
341341
}
342342
}
343343

Lines changed: 258 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,258 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
20+
package org.apache.iotdb.db.schemaengine.schemaregion.tag;
21+
22+
import org.apache.iotdb.commons.path.PartialPath;
23+
import org.apache.iotdb.commons.schema.filter.impl.TagFilter;
24+
import org.apache.iotdb.commons.schema.node.IMNode;
25+
import org.apache.iotdb.commons.schema.node.role.IDeviceMNode;
26+
import org.apache.iotdb.commons.schema.node.role.IMeasurementMNode;
27+
import org.apache.iotdb.commons.schema.node.utils.IMNodeFactory;
28+
import org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.mem.mnode.IMemMNode;
29+
import org.apache.iotdb.db.schemaengine.schemaregion.mtree.loader.MNodeFactoryLoader;
30+
import org.apache.iotdb.db.utils.ManualPerformanceTestUtils;
31+
import org.apache.iotdb.db.utils.ManualPerformanceTestUtils.Measurement;
32+
import org.apache.iotdb.db.utils.ManualPerformanceTestUtils.Summary;
33+
34+
import org.apache.commons.io.FileUtils;
35+
import org.apache.tsfile.enums.TSDataType;
36+
import org.apache.tsfile.file.metadata.enums.CompressionType;
37+
import org.apache.tsfile.file.metadata.enums.TSEncoding;
38+
import org.apache.tsfile.write.schema.MeasurementSchema;
39+
import org.junit.Assert;
40+
import org.junit.Assume;
41+
import org.junit.Test;
42+
43+
import java.io.File;
44+
import java.lang.reflect.Field;
45+
import java.nio.file.Files;
46+
import java.util.ArrayList;
47+
import java.util.Collections;
48+
import java.util.Comparator;
49+
import java.util.List;
50+
import java.util.Locale;
51+
import java.util.Map;
52+
import java.util.Set;
53+
import java.util.stream.Collectors;
54+
55+
public class TagManagerPerformanceTest {
56+
private static final String PREFIX = "iotdb.tag.index.perf.";
57+
private static volatile List<IMeasurementMNode<?>> benchmarkBlackhole;
58+
59+
@Test
60+
public void compareExactLookupAndPathFiltering() throws Exception {
61+
Assume.assumeTrue(
62+
"Manual performance UT: enable with -D" + PREFIX + "enabled=true",
63+
Boolean.getBoolean(PREFIX + "enabled"));
64+
Assume.assumeTrue(
65+
"Current-thread CPU and allocation metrics are required.",
66+
ManualPerformanceTestUtils.enableThreadMetrics());
67+
final int distinctValues = Integer.getInteger(PREFIX + "distinctValues", 10_000);
68+
final int candidateNodes = Integer.getInteger(PREFIX + "candidateNodes", 2_000);
69+
final int matchingNodes = Integer.getInteger(PREFIX + "matchingNodes", 100);
70+
final int warmups = Integer.getInteger(PREFIX + "warmups", 30);
71+
final int exactIterations = Integer.getInteger(PREFIX + "exactIterations", 50_000);
72+
final int pathIterations = Integer.getInteger(PREFIX + "pathIterations", 1_000);
73+
final int rounds = Integer.getInteger(PREFIX + "rounds", 5);
74+
Assert.assertTrue(
75+
distinctValues > 0
76+
&& candidateNodes > 0
77+
&& matchingNodes > 0
78+
&& matchingNodes <= candidateNodes);
79+
Assert.assertTrue(warmups >= 0 && exactIterations > 0 && pathIterations > 0 && rounds > 0);
80+
81+
final File directory = Files.createTempDirectory("tag-index-performance").toFile();
82+
try {
83+
final TagManager manager = new TagManager(directory.getAbsolutePath(), null);
84+
try {
85+
final IMNodeFactory<IMemMNode> factory =
86+
MNodeFactoryLoader.getInstance().getMemMNodeIMNodeFactory();
87+
final IMemMNode root = factory.createInternalMNode(null, "root");
88+
final IDeviceMNode<IMemMNode> matchingDevice =
89+
factory.createDeviceMNode(factory.createInternalMNode(root, "sg"), "d");
90+
final IDeviceMNode<IMemMNode> otherDevice =
91+
factory.createDeviceMNode(factory.createInternalMNode(root, "other"), "d");
92+
final IMeasurementMNode<IMemMNode> shared =
93+
newMeasurement(factory, matchingDevice, "single");
94+
for (int i = 0; i < distinctValues; i++) {
95+
manager.addIndex("cardinality", "v" + i, shared);
96+
}
97+
for (int i = 0; i < candidateNodes; i++) {
98+
manager.addIndex(
99+
"selectivity",
100+
"target",
101+
newMeasurement(factory, i < matchingNodes ? matchingDevice : otherDevice, "s" + i));
102+
}
103+
final Map<String, Map<String, Set<IMeasurementMNode<?>>>> index = getIndex(manager);
104+
final PartialPath pattern = new PartialPath("root.sg.**");
105+
System.out.printf(
106+
Locale.ROOT,
107+
"Tag index benchmark: distinctValues=%d, candidateNodes=%d, matchingNodes=%d, warmups=%d, exactIterations/round=%d, pathIterations/round=%d, rounds=%d%n",
108+
distinctValues,
109+
candidateNodes,
110+
matchingNodes,
111+
warmups,
112+
exactIterations,
113+
pathIterations,
114+
rounds);
115+
benchmark(
116+
"exact value lookup",
117+
manager,
118+
index,
119+
new TagFilter("cardinality", "v0", false),
120+
pattern,
121+
1,
122+
warmups,
123+
exactIterations,
124+
rounds);
125+
benchmark(
126+
"path filtering before sort",
127+
manager,
128+
index,
129+
new TagFilter("selectivity", "target", false),
130+
pattern,
131+
matchingNodes,
132+
warmups,
133+
pathIterations,
134+
rounds);
135+
} finally {
136+
manager.clear();
137+
}
138+
} finally {
139+
FileUtils.deleteDirectory(directory);
140+
}
141+
}
142+
143+
private static IMeasurementMNode<IMemMNode> newMeasurement(
144+
IMNodeFactory<IMemMNode> factory, IDeviceMNode<IMemMNode> parent, String name) {
145+
return factory.createMeasurementMNode(
146+
parent,
147+
name,
148+
new MeasurementSchema(name, TSDataType.INT64, TSEncoding.PLAIN, CompressionType.SNAPPY),
149+
null);
150+
}
151+
152+
@SuppressWarnings("unchecked")
153+
private static Map<String, Map<String, Set<IMeasurementMNode<?>>>> getIndex(TagManager manager)
154+
throws Exception {
155+
final Field field = TagManager.class.getDeclaredField("tagIndex");
156+
field.setAccessible(true);
157+
return (Map<String, Map<String, Set<IMeasurementMNode<?>>>>) field.get(manager);
158+
}
159+
160+
// Reproduce the old value scan and sort, then apply the reader's path filter.
161+
private static List<IMeasurementMNode<?>> legacyLookup(
162+
Map<String, Map<String, Set<IMeasurementMNode<?>>>> index,
163+
TagFilter filter,
164+
PartialPath pattern) {
165+
final Map<String, Set<IMeasurementMNode<?>>> values = index.get(filter.getKey());
166+
if (values == null || values.isEmpty()) {
167+
return Collections.emptyList();
168+
}
169+
final List<IMeasurementMNode<?>> candidates = new ArrayList<>();
170+
for (Map.Entry<String, Set<IMeasurementMNode<?>>> entry : values.entrySet()) {
171+
if (entry.getKey() != null
172+
&& entry.getValue() != null
173+
&& filter.getValue().equals(entry.getKey())) {
174+
candidates.addAll(entry.getValue());
175+
}
176+
}
177+
final List<IMeasurementMNode<?>> sorted =
178+
candidates.stream()
179+
.sorted(Comparator.comparing(IMNode::getFullPath))
180+
.collect(Collectors.toList());
181+
sorted.removeIf(node -> !pattern.matchFullPath(node.getPartialPath()));
182+
return sorted;
183+
}
184+
185+
private static void benchmark(
186+
String label,
187+
TagManager manager,
188+
Map<String, Map<String, Set<IMeasurementMNode<?>>>> index,
189+
TagFilter filter,
190+
PartialPath pattern,
191+
int expectedMatches,
192+
int warmups,
193+
int iterations,
194+
int rounds) {
195+
final List<IMeasurementMNode<?>> legacyResult = legacyLookup(index, filter, pattern);
196+
final List<IMeasurementMNode<?>> optimizedResult =
197+
manager.getMatchedTimeseriesInIndex(filter, pattern, false);
198+
Assert.assertEquals(expectedMatches, optimizedResult.size());
199+
Assert.assertEquals(legacyResult, optimizedResult);
200+
final Runnable legacy = () -> benchmarkBlackhole = legacyLookup(index, filter, pattern);
201+
final Runnable optimized =
202+
() -> benchmarkBlackhole = manager.getMatchedTimeseriesInIndex(filter, pattern, false);
203+
for (int i = 0; i < warmups; i++) {
204+
if ((i & 1) == 0) {
205+
legacy.run();
206+
optimized.run();
207+
} else {
208+
optimized.run();
209+
legacy.run();
210+
}
211+
}
212+
final Measurement[] legacyMeasurements = new Measurement[rounds];
213+
final Measurement[] optimizedMeasurements = new Measurement[rounds];
214+
for (int i = 0; i < rounds; i++) {
215+
if ((i & 1) == 0) {
216+
legacyMeasurements[i] = ManualPerformanceTestUtils.measure(iterations, legacy);
217+
optimizedMeasurements[i] = ManualPerformanceTestUtils.measure(iterations, optimized);
218+
} else {
219+
optimizedMeasurements[i] = ManualPerformanceTestUtils.measure(iterations, optimized);
220+
legacyMeasurements[i] = ManualPerformanceTestUtils.measure(iterations, legacy);
221+
}
222+
}
223+
final Summary oldSummary = ManualPerformanceTestUtils.summarize(legacyMeasurements, iterations);
224+
final Summary newSummary =
225+
ManualPerformanceTestUtils.summarize(optimizedMeasurements, iterations);
226+
System.out.printf(Locale.ROOT, " %s (matches=%d):%n", label, expectedMatches);
227+
printSummary("legacy", oldSummary);
228+
printSummary("optimized", newSummary);
229+
final double allocationReduction =
230+
oldSummary.getAllocatedBytesPerOperation() == 0
231+
? 0
232+
: (oldSummary.getAllocatedBytesPerOperation()
233+
- newSummary.getAllocatedBytesPerOperation())
234+
* 100.0
235+
/ oldSummary.getAllocatedBytesPerOperation();
236+
if (oldSummary.getCpuNanosPerOperation() > 0 && newSummary.getCpuNanosPerOperation() > 0) {
237+
System.out.printf(
238+
Locale.ROOT,
239+
" CPU speedup=%.2fx, allocation reduction=%.1f%%%n",
240+
oldSummary.getCpuNanosPerOperation() / newSummary.getCpuNanosPerOperation(),
241+
allocationReduction);
242+
} else {
243+
System.out.printf(
244+
Locale.ROOT,
245+
" CPU speedup=n/a (timer resolution), allocation reduction=%.1f%%%n",
246+
allocationReduction);
247+
}
248+
}
249+
250+
private static void printSummary(String label, Summary summary) {
251+
System.out.printf(
252+
Locale.ROOT,
253+
" %-9s CPU=%.3f us/query, allocated=%.1f B/query%n",
254+
label,
255+
summary.getCpuNanosPerOperation() / 1_000.0,
256+
summary.getAllocatedBytesPerOperation());
257+
}
258+
}

0 commit comments

Comments
 (0)