Skip to content

Commit 9ff0043

Browse files
committed
Limit concurrent load tsfile type conversions (#18014)
* Limit concurrent load tsfile type conversions * Keep load conversion retryable on memory pressure * Use i18n message for active load retry log * Add active load retry coverage for temporary unavailable * Fix load conversion memory test compile (cherry picked from commit 4d05a85)
1 parent 818b6fc commit 9ff0043

13 files changed

Lines changed: 659 additions & 5 deletions

File tree

‎integration-test/src/main/java/org/apache/iotdb/it/env/cluster/config/MppDataNodeConfig.java‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -89,6 +89,31 @@ public DataNodeConfig setLoadTsFileAnalyzeSchemaMemorySizeInBytes(
8989
return this;
9090
}
9191

92+
@Override
93+
public DataNodeConfig setMaxAllocateMemoryRatioForLoad(double maxAllocateMemoryRatioForLoad) {
94+
properties.setProperty(
95+
"max_allocate_memory_ratio_for_load", String.valueOf(maxAllocateMemoryRatioForLoad));
96+
return this;
97+
}
98+
99+
@Override
100+
public DataNodeConfig setLoadTsFileTabletConversionBatchMemorySizeInBytes(
101+
long loadTsFileTabletConversionBatchMemorySizeInBytes) {
102+
properties.setProperty(
103+
"load_tsfile_tablet_conversion_batch_memory_size_in_bytes",
104+
String.valueOf(loadTsFileTabletConversionBatchMemorySizeInBytes));
105+
return this;
106+
}
107+
108+
@Override
109+
public DataNodeConfig setLoadActiveListeningCheckIntervalSeconds(
110+
long loadActiveListeningCheckIntervalSeconds) {
111+
properties.setProperty(
112+
"load_active_listening_check_interval_seconds",
113+
String.valueOf(loadActiveListeningCheckIntervalSeconds));
114+
return this;
115+
}
116+
92117
@Override
93118
public DataNodeConfig setLoadLastCacheStrategy(String strategyName) {
94119
setProperty("last_cache_operation_on_load", strategyName);

‎integration-test/src/main/java/org/apache/iotdb/it/env/remote/config/RemoteDataNodeConfig.java‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,23 @@ public DataNodeConfig setLoadTsFileAnalyzeSchemaMemorySizeInBytes(
5454
return this;
5555
}
5656

57+
@Override
58+
public DataNodeConfig setMaxAllocateMemoryRatioForLoad(double maxAllocateMemoryRatioForLoad) {
59+
return this;
60+
}
61+
62+
@Override
63+
public DataNodeConfig setLoadTsFileTabletConversionBatchMemorySizeInBytes(
64+
long loadTsFileTabletConversionBatchMemorySizeInBytes) {
65+
return this;
66+
}
67+
68+
@Override
69+
public DataNodeConfig setLoadActiveListeningCheckIntervalSeconds(
70+
long loadActiveListeningCheckIntervalSeconds) {
71+
return this;
72+
}
73+
5774
@Override
5875
public DataNodeConfig setLoadLastCacheStrategy(String strategyName) {
5976
return this;

‎integration-test/src/main/java/org/apache/iotdb/itbase/env/DataNodeConfig.java‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,14 @@ public interface DataNodeConfig {
3636
DataNodeConfig setLoadTsFileAnalyzeSchemaMemorySizeInBytes(
3737
long loadTsFileAnalyzeSchemaMemorySizeInBytes);
3838

39+
DataNodeConfig setMaxAllocateMemoryRatioForLoad(double maxAllocateMemoryRatioForLoad);
40+
41+
DataNodeConfig setLoadTsFileTabletConversionBatchMemorySizeInBytes(
42+
long loadTsFileTabletConversionBatchMemorySizeInBytes);
43+
44+
DataNodeConfig setLoadActiveListeningCheckIntervalSeconds(
45+
long loadActiveListeningCheckIntervalSeconds);
46+
3947
DataNodeConfig setLoadLastCacheStrategy(String strategyName);
4048

4149
DataNodeConfig setCacheLastValuesForLoad(boolean cacheLastValuesForLoad);
Lines changed: 226 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,226 @@
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.it;
21+
22+
import org.apache.iotdb.it.env.EnvFactory;
23+
import org.apache.iotdb.it.env.cluster.node.DataNodeWrapper;
24+
import org.apache.iotdb.it.framework.IoTDBTestRunner;
25+
import org.apache.iotdb.it.utils.TsFileGenerator;
26+
import org.apache.iotdb.itbase.category.ClusterIT;
27+
import org.apache.iotdb.itbase.category.LocalStandaloneIT;
28+
29+
import org.apache.tsfile.enums.TSDataType;
30+
import org.apache.tsfile.file.metadata.enums.TSEncoding;
31+
import org.apache.tsfile.write.schema.MeasurementSchema;
32+
import org.junit.After;
33+
import org.junit.Assert;
34+
import org.junit.Before;
35+
import org.junit.Test;
36+
import org.junit.experimental.categories.Category;
37+
import org.junit.runner.RunWith;
38+
39+
import java.io.File;
40+
import java.nio.file.Files;
41+
import java.sql.Connection;
42+
import java.sql.Statement;
43+
import java.util.Collections;
44+
import java.util.concurrent.TimeUnit;
45+
46+
@RunWith(IoTDBTestRunner.class)
47+
@Category({LocalStandaloneIT.class, ClusterIT.class})
48+
public class IoTDBLoadTsFileActiveRetryIT {
49+
50+
private static final String DATABASE = "root.sg.test_0";
51+
private static final String DEVICE = DATABASE + ".d_0";
52+
private static final String MEASUREMENT = "sensor_00";
53+
private static final long UNALLOCATABLE_TABLET_CONVERSION_BATCH_MEMORY_SIZE_IN_BYTES =
54+
Long.MAX_VALUE / 4;
55+
private static final MeasurementSchema TSFILE_SCHEMA =
56+
new MeasurementSchema(MEASUREMENT, TSDataType.INT32, TSEncoding.RLE);
57+
58+
private File tmpDir;
59+
60+
@Before
61+
public void setUp() throws Exception {
62+
tmpDir = new File(Files.createTempDirectory("load-active-retry").toUri());
63+
EnvFactory.getEnv().getConfig().getCommonConfig().setPipeMemoryManagementEnabled(false);
64+
EnvFactory.getEnv()
65+
.getConfig()
66+
.getDataNodeConfig()
67+
.setMaxAllocateMemoryRatioForLoad(1.0)
68+
.setLoadTsFileAnalyzeSchemaMemorySizeInBytes(1)
69+
.setLoadTsFileTabletConversionBatchMemorySizeInBytes(
70+
UNALLOCATABLE_TABLET_CONVERSION_BATCH_MEMORY_SIZE_IN_BYTES)
71+
.setLoadActiveListeningCheckIntervalSeconds(1);
72+
73+
EnvFactory.getEnv().initClusterEnvironment();
74+
}
75+
76+
@After
77+
public void tearDown() throws Exception {
78+
try (final Connection connection = EnvFactory.getEnv().getConnection();
79+
final Statement statement = connection.createStatement()) {
80+
statement.execute("delete database " + DATABASE);
81+
} catch (final Exception ignored) {
82+
// ignore cleanup failure
83+
} finally {
84+
EnvFactory.getEnv().cleanClusterEnvironment();
85+
deleteRecursively(tmpDir);
86+
}
87+
}
88+
89+
@Test
90+
public void testActiveLoadTemporaryUnavailableShouldKeepFileForRetry() throws Exception {
91+
final DataNodeWrapper dataNodeWrapper = EnvFactory.getEnv().getDataNodeWrapper(0);
92+
final File retryTsFile = new File(tmpDir, "1-0-0-0.tsfile");
93+
final File permanentFailureTsFile = new File(tmpDir, "2-0-0-0.tsfile");
94+
generateTsFile(retryTsFile);
95+
generateTsFile(permanentFailureTsFile);
96+
97+
try (final Connection connection =
98+
EnvFactory.getEnv().getConnectionWithSpecifiedDataNode(dataNodeWrapper);
99+
final Statement statement = connection.createStatement()) {
100+
statement.execute("create database " + DATABASE);
101+
statement.execute(
102+
String.format(
103+
"create timeseries %s.%s %s", DEVICE, MEASUREMENT, TSDataType.INT64.name()));
104+
105+
statement.execute(
106+
String.format(
107+
"load \"%s\" with ('database-level'='3', 'async'='true', 'on-success'='none', "
108+
+ "'convert-on-type-mismatch'='true')",
109+
retryTsFile.getAbsolutePath()));
110+
statement.execute(
111+
String.format(
112+
"load \"%s\" with ('database-level'='3', 'async'='true', 'on-success'='none', "
113+
+ "'convert-on-type-mismatch'='false')",
114+
permanentFailureTsFile.getAbsolutePath()));
115+
116+
final File activeDir = getActiveLoadDir(dataNodeWrapper);
117+
final File failDir = getActiveLoadFailDir(dataNodeWrapper);
118+
final File activeTsFile = waitForFile(activeDir, retryTsFile.getName(), 30_000L);
119+
120+
Assert.assertNotNull(
121+
"Async load should copy tsfile into active load directory", activeTsFile);
122+
123+
Assert.assertNotNull(
124+
"Permanent active load failure should be moved to fail dir",
125+
waitForFile(failDir, permanentFailureTsFile.getName(), TimeUnit.SECONDS.toMillis(60)));
126+
127+
assertFileKeptForRetry(
128+
activeDir, failDir, retryTsFile.getName(), TimeUnit.SECONDS.toMillis(12));
129+
}
130+
}
131+
132+
private void generateTsFile(final File tsFile) throws Exception {
133+
try (final TsFileGenerator generator = new TsFileGenerator(tsFile)) {
134+
generator.registerTimeseries(DEVICE, Collections.singletonList(TSFILE_SCHEMA));
135+
generator.generateData(DEVICE, 10, 1, false);
136+
}
137+
}
138+
139+
private File getActiveLoadDir(final DataNodeWrapper dataNodeWrapper) {
140+
return new File(
141+
dataNodeWrapper.getNodePath()
142+
+ File.separator
143+
+ "ext"
144+
+ File.separator
145+
+ "load"
146+
+ File.separator
147+
+ "pending");
148+
}
149+
150+
private File getActiveLoadFailDir(final DataNodeWrapper dataNodeWrapper) {
151+
return new File(
152+
dataNodeWrapper.getNodePath()
153+
+ File.separator
154+
+ "ext"
155+
+ File.separator
156+
+ "load"
157+
+ File.separator
158+
+ "failed");
159+
}
160+
161+
private File waitForFile(final File root, final String fileName, final long timeoutMs)
162+
throws InterruptedException {
163+
final long deadline = System.currentTimeMillis() + timeoutMs;
164+
while (System.currentTimeMillis() < deadline) {
165+
final File file = findFile(root, fileName);
166+
if (file != null) {
167+
return file;
168+
}
169+
Thread.sleep(500L);
170+
}
171+
return null;
172+
}
173+
174+
private boolean containsFile(final File root, final String fileName) {
175+
return findFile(root, fileName) != null;
176+
}
177+
178+
private void assertFileKeptForRetry(
179+
final File activeDir, final File failDir, final String fileName, final long observationMs)
180+
throws InterruptedException {
181+
final long deadline = System.currentTimeMillis() + observationMs;
182+
while (System.currentTimeMillis() < deadline) {
183+
Assert.assertTrue(
184+
"Temporary unavailable active load should keep tsfile for retry",
185+
containsFile(activeDir, fileName));
186+
Assert.assertFalse(
187+
"Temporary unavailable active load must not move tsfile to fail dir",
188+
containsFile(failDir, fileName));
189+
Thread.sleep(500L);
190+
}
191+
}
192+
193+
private File findFile(final File root, final String fileName) {
194+
if (root == null || !root.exists()) {
195+
return null;
196+
}
197+
if (root.isFile()) {
198+
return root.getName().equals(fileName) ? root : null;
199+
}
200+
201+
final File[] children = root.listFiles();
202+
if (children == null) {
203+
return null;
204+
}
205+
for (final File child : children) {
206+
final File file = findFile(child, fileName);
207+
if (file != null) {
208+
return file;
209+
}
210+
}
211+
return null;
212+
}
213+
214+
private void deleteRecursively(final File file) {
215+
if (file == null || !file.exists()) {
216+
return;
217+
}
218+
final File[] children = file.listFiles();
219+
if (children != null) {
220+
for (final File child : children) {
221+
deleteRecursively(child);
222+
}
223+
}
224+
Assert.assertTrue(file.delete());
225+
}
226+
}

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1158,6 +1158,8 @@ public class IoTDBConfig {
11581158

11591159
private long loadTsFileTabletConversionBatchMemorySizeInBytes = 4096 * 1024;
11601160

1161+
private int loadTsFileTabletConversionThreadCount = 5;
1162+
11611163
private long loadChunkMetadataMemorySizeInBytes = 33554432; // 32MB
11621164

11631165
private long loadMemoryAllocateRetryIntervalMs = 1000L;
@@ -4141,6 +4143,14 @@ public void setLoadTsFileTabletConversionBatchMemorySizeInBytes(
41414143
loadTsFileTabletConversionBatchMemorySizeInBytes;
41424144
}
41434145

4146+
public int getLoadTsFileTabletConversionThreadCount() {
4147+
return loadTsFileTabletConversionThreadCount;
4148+
}
4149+
4150+
public void setLoadTsFileTabletConversionThreadCount(int loadTsFileTabletConversionThreadCount) {
4151+
this.loadTsFileTabletConversionThreadCount = loadTsFileTabletConversionThreadCount;
4152+
}
4153+
41444154
public long getLoadChunkMetadataMemorySizeInBytes() {
41454155
return loadChunkMetadataMemorySizeInBytes;
41464156
}

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2426,6 +2426,11 @@ private void loadLoadTsFileProps(TrimProperties properties) throws IOException {
24262426
properties.getProperty(
24272427
"load_tsfile_tablet_conversion_batch_memory_size_in_bytes",
24282428
String.valueOf(conf.getLoadTsFileTabletConversionBatchMemorySizeInBytes()))));
2429+
conf.setLoadTsFileTabletConversionThreadCount(
2430+
Integer.parseInt(
2431+
properties.getProperty(
2432+
"load_tsfile_tablet_conversion_thread_count",
2433+
String.valueOf(conf.getLoadTsFileTabletConversionThreadCount()))));
24292434
conf.setLoadChunkMetadataMemorySizeInBytes(
24302435
Long.parseLong(
24312436
Optional.ofNullable(

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadFailedMessageHandler.java‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,9 @@
1919

2020
package org.apache.iotdb.db.storageengine.load.active;
2121

22+
import org.apache.iotdb.common.rpc.thrift.TSStatus;
2223
import org.apache.iotdb.commons.conf.CommonDescriptor;
24+
import org.apache.iotdb.rpc.TSStatusCode;
2325

2426
import org.slf4j.Logger;
2527
import org.slf4j.LoggerFactory;
@@ -96,6 +98,20 @@ private interface ExceptionMessageHandler {
9698
void handle(final ActiveLoadPendingQueue.ActiveLoadEntry activeLoadEntry);
9799
}
98100

101+
public static boolean isStatusShouldRetry(
102+
final ActiveLoadPendingQueue.ActiveLoadEntry entry, final TSStatus status) {
103+
if (status != null
104+
&& status.getCode() == TSStatusCode.LOAD_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode()) {
105+
LOGGER.info(
106+
"Rejecting auto load tsfile {} (isGeneratedByPipe = {}) due to temporary unavailability, will retry later. Status: {}",
107+
entry.getFile(),
108+
entry.isGeneratedByPipe(),
109+
status);
110+
return true;
111+
}
112+
return isExceptionMessageShouldRetry(entry, status == null ? null : status.getMessage());
113+
}
114+
99115
public static boolean isExceptionMessageShouldRetry(
100116
final ActiveLoadPendingQueue.ActiveLoadEntry entry, final String message) {
101117
if (CommonDescriptor.getInstance().getConfig().isReadOnly()) {

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -280,7 +280,7 @@ private TSStatus executeStatement(final Statement statement, final IClientSessio
280280

281281
private void handleLoadFailure(
282282
final ActiveLoadPendingQueue.ActiveLoadEntry entry, final TSStatus status) {
283-
if (!ActiveLoadFailedMessageHandler.isExceptionMessageShouldRetry(entry, status.getMessage())) {
283+
if (!ActiveLoadFailedMessageHandler.isStatusShouldRetry(entry, status)) {
284284
LOGGER.warn(
285285
"Failed to auto load tsfile {} (isGeneratedByPipe = {}), status: {}. File will be moved to fail directory.",
286286
entry.getFile(),

0 commit comments

Comments
 (0)