From 7304c1364306589bad5932a1bfe0328502401935 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Thu, 27 Aug 2026 14:18:07 +0800 Subject: [PATCH 1/2] fix(load): validate TsFile before async move --- .../iotdb/db/it/IoTDBLoadTsFileAuthIT.java | 29 +++++++++++++++++++ .../statement/crud/LoadTsFileStatement.java | 2 +- .../db/storageengine/load/util/LoadUtil.java | 20 +++++++++++++ 3 files changed, 50 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 d842941e0cc9c..0424eb189cbf5 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.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 void testAsyncLoadShouldCheckWriteDataPermissionWithStoredUser() throws E 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 0b0dea7631121..ab8555b64539a 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 @@ -137,7 +137,7 @@ private static List processTsFile( } final List 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 4a67b362a5b53..3abfdaa59fefb 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 @@ -35,7 +35,9 @@ 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; @@ -118,6 +120,11 @@ private static boolean loadTsFilesToActiveDir( 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 = @@ -145,6 +152,19 @@ private static boolean loadTsFilesToActiveDir( 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 appendCurrentUserIfAbsent( final Map loadAttributes) { final Map attributes = From 44a881860afdcc4f35ea7a60e4d920a5dd2d9247 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Fri, 28 Aug 2026 15:29:52 +0800 Subject: [PATCH 2/2] fix(load): report invalid tsfile in async pipe path --- .../db/storageengine/load/util/LoadUtil.java | 18 ++++++ .../storageengine/load/util/LoadUtilTest.java | 59 +++++++++++++++++++ 2 files changed, 77 insertions(+) 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 3abfdaa59fefb..a9a187ad4f125 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.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; @@ -69,6 +70,13 @@ public static boolean loadTsFileAsyncToActiveDir( } 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; @@ -212,6 +220,16 @@ public static boolean loadFilesToActiveDir( 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, sourceFiles, isDeleteAfterLoad); return true; 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 f9537835da783..3848f93370466 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,10 +19,14 @@ 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.tsfile.write.writer.TsFileIOWriter; import org.junit.After; import org.junit.Assert; import org.junit.Before; @@ -39,6 +43,8 @@ public class LoadUtilTest { + private final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); + private String[] originalListeningDirs; private File tempDir; private File sourceDir; private File targetDir; @@ -50,13 +56,58 @@ public void setUp() throws Exception { 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 { @@ -126,6 +177,14 @@ private List createTsFileAndCompanions() throws Exception { 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;