This is an automated email from the ASF dual-hosted git repository.
jt2594838 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new b9683c0c106 fix(load): validate TsFile before async move (#18531)
b9683c0c106 is described below
commit b9683c0c10664d90ae85102deaeb554a9aab0806
Author: Zhenyu Luo <[email protected]>
AuthorDate: Mon Aug 31 11:22:57 2026 +0800
fix(load): validate TsFile before async move (#18531)
* fix(load): validate TsFile before async move
* fix(load): report invalid tsfile in async pipe path
---
.../apache/iotdb/db/it/IoTDBLoadTsFileAuthIT.java | 29 +++++++++++
.../plan/statement/crud/LoadTsFileStatement.java | 2 +-
.../iotdb/db/storageengine/load/util/LoadUtil.java | 38 ++++++++++++++
.../db/storageengine/load/util/LoadUtilTest.java | 59 ++++++++++++++++++++++
4 files changed, 127 insertions(+), 1 deletion(-)
diff --git
a/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileAuthIT.java
b/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileAuthIT.java
index d842941e0cc..0424eb189cb 100644
---
a/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileAuthIT.java
+++
b/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileAuthIT.java
@@ -42,6 +42,7 @@ import org.junit.experimental.categories.Category;
import org.junit.runner.RunWith;
import java.io.File;
+import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.sql.Connection;
import java.sql.ResultSet;
@@ -180,6 +181,34 @@ public class IoTDBLoadTsFileAuthIT {
assertCountEventually(10, TimeUnit.SECONDS.toMillis(60));
}
+ @Test
+ public void testAsyncLoadShouldNotMoveFileWithoutTsFileSuffix() throws
Exception {
+ assertAsyncLoadDoesNotMoveInvalidTsFile(new File(tmpDir,
"not-a-tsfile.txt"), "Can not find");
+ }
+
+ @Test
+ public void testAsyncLoadShouldNotMoveInvalidTsFile() throws Exception {
+ assertAsyncLoadDoesNotMoveInvalidTsFile(new File(tmpDir,
"invalid.tsfile"), "Loading file");
+ }
+
+ private static void assertAsyncLoadDoesNotMoveInvalidTsFile(
+ final File sourceFile, final String expectedErrorMessage) throws
Exception {
+ final String originalContent = "ordinary file content";
+ Files.write(sourceFile.toPath(),
originalContent.getBytes(StandardCharsets.UTF_8));
+
+ assertNonQueryTestFail(
+ String.format(
+ "load \"%s\" with ('async'='true', 'on-success'='delete')",
+ sourceFile.getAbsolutePath()),
+ expectedErrorMessage);
+
+ Assert.assertTrue(
+ "Non-TsFile source must remain in its original location",
sourceFile.isFile());
+ Assert.assertEquals(
+ originalContent,
+ new String(Files.readAllBytes(sourceFile.toPath()),
StandardCharsets.UTF_8));
+ }
+
private static void prepareSchemaAndTsFile(final File tsFile) throws
Exception {
prepareSchema(MEASUREMENT.getType());
generateTsFile(tsFile);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatement.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatement.java
index ca0a0a3e2bf..cacbe36cdbb 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatement.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatement.java
@@ -145,7 +145,7 @@ public class LoadTsFileStatement extends Statement {
}
final List<File> tsFiles = new ArrayList<>();
- if (file.isFile()) {
+ if (file.isFile() &&
file.getName().endsWith(TsFileConstant.TSFILE_SUFFIX)) {
tsFiles.add(file);
} else {
if (file.listFiles() == null) {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/util/LoadUtil.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/util/LoadUtil.java
index ea1440bd316..8d653002dbc 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/util/LoadUtil.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/util/LoadUtil.java
@@ -26,6 +26,7 @@ import org.apache.iotdb.commons.utils.FileUtils;
import org.apache.iotdb.commons.utils.RetryUtils;
import org.apache.iotdb.db.auth.AuthorityChecker;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.i18n.DataNodeQueryMessages;
import org.apache.iotdb.db.i18n.StorageEngineMessages;
import org.apache.iotdb.db.protocol.session.IClientSession;
import org.apache.iotdb.db.protocol.session.SessionManager;
@@ -35,7 +36,9 @@ import
org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
import org.apache.iotdb.db.storageengine.load.active.ActiveLoadPathHelper;
import org.apache.iotdb.db.storageengine.load.disk.ILoadDiskSelector;
+import org.apache.tsfile.common.conf.TSFileConfig;
import org.apache.tsfile.common.constant.TsFileConstant;
+import org.apache.tsfile.read.TsFileSequenceReader;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -68,6 +71,13 @@ public class LoadUtil {
}
try {
+ // Validate the complete batch before moving any source file. Otherwise
a later malformed
+ // file could leave earlier files queued in the active-load directory.
+ for (final File file : tsFiles) {
+ if (file != null && !isValidTsFile(file)) {
+ return false;
+ }
+ }
for (File file : tsFiles) {
if (!loadTsFilesToActiveDir(loadAttributes, file, isDeleteAfterLoad)) {
return false;
@@ -119,6 +129,11 @@ public class LoadUtil {
return true;
}
+ // Validate before moving the source so ordinary or malformed files remain
in place.
+ if (!isValidTsFile(file)) {
+ return false;
+ }
+
final File targetFilePath;
try {
targetFilePath =
@@ -148,6 +163,19 @@ public class LoadUtil {
return true;
}
+ private static boolean isValidTsFile(final File file) {
+ if (!file.isFile() ||
!file.getName().endsWith(TsFileConstant.TSFILE_SUFFIX)) {
+ return false;
+ }
+ try (final TsFileSequenceReader reader =
+ new TsFileSequenceReader(file.getAbsolutePath(), false)) {
+ return TSFileConfig.MAGIC_STRING.equals(reader.readHeadMagic())
+ && TSFileConfig.MAGIC_STRING.equals(reader.readTailMagic());
+ } catch (Exception e) {
+ return false;
+ }
+ }
+
private static Map<String, String> appendCurrentUserIfAbsent(
final Map<String, String> loadAttributes) {
final Map<String, String> attributes =
@@ -196,6 +224,16 @@ public class LoadUtil {
for (final String file : files) {
sourceFiles.add(new File(file));
}
+ // The main TsFile must be valid before any TsFile or sidecar is
transferred. Pipe acknowledges
+ // this method immediately, so deferring validation to the active loader
is too late.
+ for (final File sourceFile : sourceFiles) {
+ if (isTsFile(sourceFile) && !isValidTsFile(sourceFile)) {
+ throw new IOException(
+ String.format(
+
DataNodeQueryMessages.THE_FILE_S_IS_NOT_A_VALID_TSFILE_PLEASE_CHECK_THE_INPUT_FILE,
+ sourceFile.getAbsolutePath()));
+ }
+ }
sourceFiles.sort(Comparator.comparing(LoadUtil::isTsFile));
transferFilesToActiveDir(
targetDir,
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/util/LoadUtilTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/util/LoadUtilTest.java
index bffa2c5578c..5531734e105 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/util/LoadUtilTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/util/LoadUtilTest.java
@@ -19,11 +19,15 @@
package org.apache.iotdb.db.storageengine.load.util;
+import org.apache.iotdb.db.conf.IoTDBConfig;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.i18n.DataNodeQueryMessages;
import
org.apache.iotdb.db.storageengine.dataregion.modification.ModificationFile;
import
org.apache.iotdb.db.storageengine.dataregion.modification.v1.ModificationFileV1;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
import org.apache.iotdb.db.storageengine.load.active.ActiveLoadPathHelper;
+import org.apache.tsfile.write.writer.TsFileIOWriter;
import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
@@ -40,6 +44,8 @@ import java.util.Set;
public class LoadUtilTest {
+ private final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
+ private String[] originalListeningDirs;
private File tempDir;
private File sourceDir;
private File targetDir;
@@ -51,13 +57,58 @@ public class LoadUtilTest {
targetDir = new File(tempDir, "target");
Assert.assertTrue(sourceDir.mkdirs());
Assert.assertTrue(targetDir.mkdirs());
+ originalListeningDirs = config.getLoadActiveListeningDirs();
+ config.setLoadActiveListeningDirs(new String[]
{targetDir.getAbsolutePath()});
+ LoadUtil.updateLoadDiskSelector();
}
@After
public void tearDown() {
+ config.setLoadActiveListeningDirs(originalListeningDirs);
+ LoadUtil.updateLoadDiskSelector();
deleteRecursively(tempDir);
}
+ @Test
+ public void testAsyncLoadValidatesCompleteBatchBeforeTransfer() throws
Exception {
+ final File validTsFile = createCompletedTsFile("valid.tsfile");
+ final File invalidTsFile = new File(sourceDir, "invalid.tsfile");
+ Files.write(invalidTsFile.toPath(),
"invalid".getBytes(StandardCharsets.UTF_8));
+
+ Assert.assertFalse(
+ LoadUtil.loadTsFileAsyncToActiveDir(Arrays.asList(validTsFile,
invalidTsFile), null, true));
+ Assert.assertTrue(validTsFile.exists());
+ Assert.assertTrue(invalidTsFile.exists());
+ Assert.assertEquals(0, targetDir.listFiles().length);
+ }
+
+ @Test
+ public void testPipeAsyncLoadReportsInvalidTsFileAndKeepsSources() throws
Exception {
+ final File invalidTsFile = new File(sourceDir, "invalid.tsfile");
+ final File resourceFile =
+ new File(invalidTsFile.getAbsolutePath() +
TsFileResource.RESOURCE_SUFFIX);
+ Files.write(invalidTsFile.toPath(),
"invalid".getBytes(StandardCharsets.UTF_8));
+ Files.write(resourceFile.toPath(),
"resource".getBytes(StandardCharsets.UTF_8));
+
+ try {
+ LoadUtil.loadFilesToActiveDir(
+ null,
+ Arrays.asList(resourceFile.getAbsolutePath(),
invalidTsFile.getAbsolutePath()),
+ true);
+ Assert.fail("Expected invalid TsFile error");
+ } catch (final IOException e) {
+ Assert.assertEquals(
+ String.format(
+
DataNodeQueryMessages.THE_FILE_S_IS_NOT_A_VALID_TSFILE_PLEASE_CHECK_THE_INPUT_FILE,
+ invalidTsFile.getAbsolutePath()),
+ e.getMessage());
+ }
+
+ Assert.assertTrue(invalidTsFile.exists());
+ Assert.assertTrue(resourceFile.exists());
+ Assert.assertEquals(0, targetDir.listFiles().length);
+ }
+
@Test
public void
testTransferFilesKeepsSameNamedGroupsIsolatedAndDeletesSourcesAfterHandoff()
throws Exception {
@@ -224,6 +275,14 @@ public class LoadUtilTest {
return sourceFiles;
}
+ private File createCompletedTsFile(final String fileName) throws Exception {
+ final File tsFile = new File(sourceDir, fileName);
+ try (final TsFileIOWriter writer = new TsFileIOWriter(tsFile)) {
+ writer.endFile();
+ }
+ return tsFile;
+ }
+
private static void deleteRecursively(final File file) {
if (file == null || !file.exists()) {
return;