Skip to content

Commit 94694d9

Browse files
authored
Avoid redundant TsFile table schema registration (#18731)
1 parent 825bf9e commit 94694d9

6 files changed

Lines changed: 303 additions & 7 deletions

File tree

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/table/DataNodeTableCache.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -435,6 +435,7 @@ public void invalid(String database, final String tableName, final String column
435435
}
436436
}
437437

438+
@Override
438439
public long getInstanceVersion() {
439440
return instanceVersion.get();
440441
}

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/table/ITableCache.java‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,8 @@
2929

3030
public interface ITableCache {
3131

32+
long getInstanceVersion();
33+
3234
void init(final byte[] tableInitializationBytes);
3335

3436
void preUpdateTable(final String database, final TsTable table, final @Nullable String oldName);

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -98,6 +98,7 @@
9898
import org.apache.iotdb.db.queryengine.plan.relational.metadata.fetcher.cache.TableDeviceSchemaCache;
9999
import org.apache.iotdb.db.queryengine.plan.relational.metadata.fetcher.cache.TreeDeviceSchemaCacheManager;
100100
import org.apache.iotdb.db.schemaengine.table.DataNodeTableCache;
101+
import org.apache.iotdb.db.schemaengine.table.ITableCache;
101102
import org.apache.iotdb.db.service.SettleService;
102103
import org.apache.iotdb.db.service.metrics.CompactionMetrics;
103104
import org.apache.iotdb.db.service.metrics.DataNodeExceptionMetrics;
@@ -1730,14 +1731,15 @@ private TsFileProcessor insertTabletWithTypeConsistencyCheck(
17301731
private void registerToTsFile(InsertNode node, TsFileProcessor tsFileProcessor) {
17311732
final String tableName = node.getTableName();
17321733
if (tableName != null) {
1734+
final ITableCache tableCache = DataNodeTableCache.getInstance();
17331735
tsFileProcessor.registerToTsFile(
17341736
tableName,
1737+
tableCache.getInstanceVersion(),
17351738
t -> {
17361739
final String database = getDatabaseName();
17371740

17381741
TsTable tsTable =
1739-
DataNodeTableCache.getInstance()
1740-
.getTable(database, t, false, LeaseFencedRetryPolicy.RETRY_UNTIL_SUCCESS);
1742+
tableCache.getTable(database, t, false, LeaseFencedRetryPolicy.RETRY_UNTIL_SUCCESS);
17411743
if (tsTable == null) {
17421744
// There is a high probability that the leader node has been executed and is currently
17431745
// located in the follower node.

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java‎

Lines changed: 13 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -113,6 +113,7 @@
113113
import java.util.Map;
114114
import java.util.Set;
115115
import java.util.concurrent.CompletableFuture;
116+
import java.util.concurrent.ConcurrentHashMap;
116117
import java.util.concurrent.ConcurrentLinkedDeque;
117118
import java.util.concurrent.CopyOnWriteArrayList;
118119
import java.util.concurrent.ExecutionException;
@@ -156,6 +157,8 @@ public class TsFileProcessor {
156157
/** Writer for restore tsfile and flushing. */
157158
private RestorableTsFileIOWriter writer;
158159

160+
private final Map<String, Long> registeredTableSchemaVersions = new ConcurrentHashMap<>();
161+
159162
/** Tsfile resource for index this tsfile. */
160163
private final TsFileResource tsFileResource;
161164

@@ -2736,11 +2739,16 @@ public ConcurrentLinkedDeque<IMemTable> getFlushingMemTable() {
27362739
}
27372740

27382741
public void registerToTsFile(
2739-
String tableName, Function<String, TableSchema> tableSchemaFunction) {
2740-
getWriter()
2741-
.getSchema()
2742-
.getTableSchemaMap()
2743-
.put(tableName, tableSchemaFunction.apply(tableName));
2742+
String tableName, long tableCacheVersion, Function<String, TableSchema> tableSchemaFunction) {
2743+
registeredTableSchemaVersions.compute(
2744+
tableName,
2745+
(name, registeredVersion) -> {
2746+
if (registeredVersion != null && registeredVersion == tableCacheVersion) {
2747+
return registeredVersion;
2748+
}
2749+
getWriter().getSchema().getTableSchemaMap().put(name, tableSchemaFunction.apply(name));
2750+
return tableCacheVersion;
2751+
});
27442752
}
27452753

27462754
public void writeLock() {

‎iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java‎

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,7 @@
5454
import org.apache.tsfile.enums.TSDataType;
5555
import org.apache.tsfile.external.commons.io.FileUtils;
5656
import org.apache.tsfile.file.metadata.IDeviceID;
57+
import org.apache.tsfile.file.metadata.TableSchema;
5758
import org.apache.tsfile.file.metadata.enums.CompressionType;
5859
import org.apache.tsfile.file.metadata.enums.TSEncoding;
5960
import org.apache.tsfile.read.TimeValuePair;
@@ -145,6 +146,58 @@ public void tearDown() throws Exception {
145146
config.setTargetChunkSize(defaultTargetChunkSize);
146147
}
147148

149+
@Test
150+
public void testRegisterToTsFileCachesTableSchemaByVersion() throws IOException {
151+
processor =
152+
new TsFileProcessor(
153+
storageGroup,
154+
SystemFileFactory.INSTANCE.getFile(filePath),
155+
sgInfo,
156+
this::closeTsFileProcessor,
157+
(tsFileProcessor, updateMap, systemFlushTime) -> {},
158+
true);
159+
160+
final AtomicInteger supplierCalls = new AtomicInteger();
161+
final AtomicInteger secondTableSupplierCalls = new AtomicInteger();
162+
final TableSchema firstSchema = new TableSchema("table1");
163+
final TableSchema updatedSchema = new TableSchema("table1");
164+
final TableSchema secondTableSchema = new TableSchema("table2");
165+
166+
processor.registerToTsFile(
167+
"table1",
168+
1,
169+
tableName -> supplierCalls.incrementAndGet() == 1 ? firstSchema : updatedSchema);
170+
processor.registerToTsFile(
171+
"table1",
172+
1,
173+
tableName -> supplierCalls.incrementAndGet() == 1 ? firstSchema : updatedSchema);
174+
175+
Assert.assertEquals(1, supplierCalls.get());
176+
Assert.assertSame(
177+
firstSchema, processor.getWriter().getSchema().getTableSchemaMap().get("table1"));
178+
179+
processor.registerToTsFile(
180+
"table2",
181+
1,
182+
tableName -> {
183+
secondTableSupplierCalls.incrementAndGet();
184+
return secondTableSchema;
185+
});
186+
187+
Assert.assertEquals(1, secondTableSupplierCalls.get());
188+
Assert.assertSame(
189+
secondTableSchema, processor.getWriter().getSchema().getTableSchemaMap().get("table2"));
190+
191+
processor.registerToTsFile(
192+
"table1",
193+
2,
194+
tableName -> supplierCalls.incrementAndGet() == 1 ? firstSchema : updatedSchema);
195+
196+
Assert.assertEquals(2, supplierCalls.get());
197+
Assert.assertSame(
198+
updatedSchema, processor.getWriter().getSchema().getTableSchemaMap().get("table1"));
199+
}
200+
148201
@Test
149202
public void testWriteAndFlush()
150203
throws IOException, WriteProcessException, MetadataException, ExecutionException {
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,230 @@
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.storageengine.dataregion.memtable;
21+
22+
import org.apache.iotdb.commons.file.SystemFileFactory;
23+
import org.apache.iotdb.db.storageengine.dataregion.DataRegionInfo;
24+
import org.apache.iotdb.db.storageengine.dataregion.DataRegionTest;
25+
import org.apache.iotdb.db.storageengine.rescon.memory.SystemInfo;
26+
import org.apache.iotdb.db.utils.EnvironmentUtils;
27+
import org.apache.iotdb.db.utils.ManualPerformanceTestUtils;
28+
import org.apache.iotdb.db.utils.ManualPerformanceTestUtils.Measurement;
29+
import org.apache.iotdb.db.utils.ManualPerformanceTestUtils.Summary;
30+
import org.apache.iotdb.db.utils.constant.TestConstant;
31+
32+
import org.apache.tsfile.file.metadata.TableSchema;
33+
import org.junit.Assert;
34+
import org.junit.Assume;
35+
import org.junit.Test;
36+
37+
import java.io.File;
38+
import java.util.Locale;
39+
import java.util.concurrent.atomic.AtomicLong;
40+
import java.util.function.Function;
41+
42+
public class TsFileTableSchemaRegistrationPerformanceTest {
43+
44+
private static final String ENABLED_PROPERTY = "iotdb.tsfile.table.schema.perf.enabled";
45+
private static final String WARMUP_ITERATIONS_PROPERTY =
46+
"iotdb.tsfile.table.schema.perf.warmup.iterations";
47+
private static final String ITERATIONS_PROPERTY = "iotdb.tsfile.table.schema.perf.iterations";
48+
private static final String ROUNDS_PROPERTY = "iotdb.tsfile.table.schema.perf.rounds";
49+
private static volatile int benchmarkBlackhole;
50+
51+
@Test
52+
public void tableSchemaRegistrationBenchmark() throws Exception {
53+
Assume.assumeTrue(
54+
String.format(
55+
"Manual performance UT. Enable with -D%s=true, optionally tune -D%s, -D%s and -D%s.",
56+
ENABLED_PROPERTY, WARMUP_ITERATIONS_PROPERTY, ITERATIONS_PROPERTY, ROUNDS_PROPERTY),
57+
Boolean.getBoolean(ENABLED_PROPERTY));
58+
Assume.assumeTrue(
59+
"Current-thread CPU time and allocation metrics are required.",
60+
ManualPerformanceTestUtils.enableThreadMetrics());
61+
62+
final int warmupIterations = Integer.getInteger(WARMUP_ITERATIONS_PROPERTY, 20_000);
63+
final int iterations = Integer.getInteger(ITERATIONS_PROPERTY, 100_000);
64+
final int rounds = Integer.getInteger(ROUNDS_PROPERTY, 7);
65+
Assert.assertTrue(warmupIterations > 0);
66+
Assert.assertTrue(iterations > 0);
67+
Assert.assertTrue(rounds > 0);
68+
69+
final String storageGroup = "root.vehicle";
70+
final String systemDir = TestConstant.OUTPUT_DATA_DIR.concat("info");
71+
final String filePath = TestConstant.getTestTsFilePath(storageGroup, 0, 0, 0);
72+
TsFileProcessor processor = null;
73+
EnvironmentUtils.envSetUp();
74+
try {
75+
final File file = SystemFileFactory.INSTANCE.getFile(filePath);
76+
if (!file.getParentFile().exists()) {
77+
Assert.assertTrue(file.getParentFile().mkdirs());
78+
}
79+
final DataRegionInfo regionInfo =
80+
new DataRegionInfo(new DataRegionTest.DummyDataRegion(systemDir, storageGroup));
81+
processor =
82+
new TsFileProcessor(
83+
storageGroup,
84+
file,
85+
regionInfo,
86+
tsFileProcessor -> {},
87+
(tsFileProcessor, updateMap, systemFlushTime) -> {},
88+
true);
89+
final TsFileProcessorInfo processorInfo = new TsFileProcessorInfo(regionInfo);
90+
processor.setTsFileProcessorInfo(processorInfo);
91+
regionInfo.initTsFileProcessorInfo(processor);
92+
SystemInfo.getInstance().reportStorageGroupStatus(regionInfo, processor);
93+
94+
final AtomicLong legacySupplierCalls = new AtomicLong();
95+
final AtomicLong cachedSupplierCalls = new AtomicLong();
96+
final Function<String, TableSchema> legacySupplier =
97+
tableName -> {
98+
legacySupplierCalls.incrementAndGet();
99+
return new TableSchema(tableName);
100+
};
101+
final Function<String, TableSchema> cachedSupplier =
102+
tableName -> {
103+
cachedSupplierCalls.incrementAndGet();
104+
return new TableSchema(tableName);
105+
};
106+
107+
runLegacy(processor, warmupIterations, legacySupplier);
108+
runCached(processor, warmupIterations, cachedSupplier);
109+
legacySupplierCalls.set(0);
110+
cachedSupplierCalls.set(0);
111+
112+
final Measurement[] legacyMeasurements = new Measurement[rounds];
113+
final Measurement[] cachedMeasurements = new Measurement[rounds];
114+
for (int round = 0; round < rounds; ++round) {
115+
if ((round & 1) == 0) {
116+
legacyMeasurements[round] = measureLegacy(processor, iterations, legacySupplier);
117+
cachedMeasurements[round] = measureCached(processor, iterations, cachedSupplier);
118+
} else {
119+
cachedMeasurements[round] = measureCached(processor, iterations, cachedSupplier);
120+
legacyMeasurements[round] = measureLegacy(processor, iterations, legacySupplier);
121+
}
122+
}
123+
124+
final Summary legacySummary =
125+
ManualPerformanceTestUtils.summarize(legacyMeasurements, iterations);
126+
final Summary cachedSummary =
127+
ManualPerformanceTestUtils.summarize(cachedMeasurements, iterations);
128+
Assert.assertEquals((long) iterations * rounds, legacySupplierCalls.get());
129+
Assert.assertEquals(0, cachedSupplierCalls.get());
130+
printResult(
131+
warmupIterations,
132+
iterations,
133+
rounds,
134+
legacySummary,
135+
cachedSummary,
136+
legacySupplierCalls.get(),
137+
cachedSupplierCalls.get());
138+
} finally {
139+
if (processor != null) {
140+
processor.syncClose();
141+
}
142+
EnvironmentUtils.cleanEnv();
143+
EnvironmentUtils.cleanDir(TestConstant.OUTPUT_DATA_DIR);
144+
}
145+
}
146+
147+
private static Measurement measureLegacy(
148+
TsFileProcessor processor,
149+
int iterations,
150+
Function<String, TableSchema> tableSchemaFunction) {
151+
return ManualPerformanceTestUtils.measure(
152+
iterations, () -> legacyRegister(processor, tableSchemaFunction));
153+
}
154+
155+
private static Measurement measureCached(
156+
TsFileProcessor processor,
157+
int iterations,
158+
Function<String, TableSchema> tableSchemaFunction) {
159+
return ManualPerformanceTestUtils.measure(
160+
iterations, () -> cachedRegister(processor, tableSchemaFunction));
161+
}
162+
163+
private static void runLegacy(
164+
TsFileProcessor processor,
165+
int iterations,
166+
Function<String, TableSchema> tableSchemaFunction) {
167+
for (int i = 0; i < iterations; ++i) {
168+
legacyRegister(processor, tableSchemaFunction);
169+
}
170+
}
171+
172+
private static void runCached(
173+
TsFileProcessor processor,
174+
int iterations,
175+
Function<String, TableSchema> tableSchemaFunction) {
176+
for (int i = 0; i < iterations; ++i) {
177+
cachedRegister(processor, tableSchemaFunction);
178+
}
179+
}
180+
181+
private static void legacyRegister(
182+
TsFileProcessor processor, Function<String, TableSchema> tableSchemaFunction) {
183+
final String tableName = "legacy_table";
184+
processor
185+
.getWriter()
186+
.getSchema()
187+
.getTableSchemaMap()
188+
.put(tableName, tableSchemaFunction.apply(tableName));
189+
benchmarkBlackhole = processor.getWriter().getSchema().getTableSchemaMap().size();
190+
}
191+
192+
private static void cachedRegister(
193+
TsFileProcessor processor, Function<String, TableSchema> tableSchemaFunction) {
194+
processor.registerToTsFile("cached_table", 1, tableSchemaFunction);
195+
benchmarkBlackhole = processor.getWriter().getSchema().getTableSchemaMap().size();
196+
}
197+
198+
private static void printResult(
199+
int warmupIterations,
200+
int iterations,
201+
int rounds,
202+
Summary legacySummary,
203+
Summary cachedSummary,
204+
long legacySupplierCalls,
205+
long cachedSupplierCalls) {
206+
System.out.printf(
207+
Locale.ROOT,
208+
"TsFile table schema registration benchmark: warmup=%d, iterations=%d, rounds=%d%n"
209+
+ "legacy: cpu=%.2f ns/op, allocation=%.2f bytes/op, supplierCalls=%d%n"
210+
+ "cached: cpu=%.2f ns/op, allocation=%.2f bytes/op, supplierCalls=%d%n"
211+
+ "reduction: cpu=%.2f%%, allocation=%.2f%%%n",
212+
warmupIterations,
213+
iterations,
214+
rounds,
215+
legacySummary.getCpuNanosPerOperation(),
216+
legacySummary.getAllocatedBytesPerOperation(),
217+
legacySupplierCalls,
218+
cachedSummary.getCpuNanosPerOperation(),
219+
cachedSummary.getAllocatedBytesPerOperation(),
220+
cachedSupplierCalls,
221+
reduction(legacySummary.getCpuNanosPerOperation(), cachedSummary.getCpuNanosPerOperation()),
222+
reduction(
223+
legacySummary.getAllocatedBytesPerOperation(),
224+
cachedSummary.getAllocatedBytesPerOperation()));
225+
}
226+
227+
private static double reduction(double before, double after) {
228+
return before == 0 ? 0 : (before - after) * 100 / before;
229+
}
230+
}

0 commit comments

Comments
 (0)