From e777e04c3c4e57e274211eff1544c0cf06ca9d24 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Tue, 21 Jul 2026 18:28:05 +0800 Subject: [PATCH 01/10] feat(load): reject managed TsFile directories --- .../apache/iotdb/db/it/IoTDBLoadTsFileIT.java | 22 +++++++++++ .../iotdb/db/i18n/DataNodeQueryMessages.java | 3 ++ .../iotdb/db/i18n/DataNodeQueryMessages.java | 3 ++ .../statement/crud/LoadTsFileStatement.java | 37 +++++++++++++++++-- .../crud/LoadTsFileStatementTest.java | 31 ++++++++++++++++ 5 files changed, 93 insertions(+), 3 deletions(-) diff --git a/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileIT.java b/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileIT.java index 1eda0f0622651..7fc079e6d4a6f 100644 --- a/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileIT.java @@ -887,6 +887,28 @@ public void testLoadWithRelativePathName() throws Exception { } } + @Test + public void testLoadDataNodeInternalDataDirectoryIsRejectedWithoutLeakingPath() throws Exception { + final DataNodeWrapper dataNodeWrapper = EnvFactory.getEnv().getDataNodeWrapper(0); + final File dataDir = new File(dataNodeWrapper.getDataPath()); + + try (final Connection connection = + EnvFactory.getEnv().getConnectionWithSpecifiedDataNode(dataNodeWrapper); + final Statement statement = connection.createStatement()) { + try { + statement.execute(String.format("load \"%s\"", dataDir.getAbsolutePath())); + Assert.fail("Expected LOAD from the DataNode internal data directory to be rejected."); + } catch (final SQLException e) { + Assert.assertTrue( + e.getMessage(), + e.getMessage() + .contains( + "Cannot load files because the specified directory contains IoTDB data.")); + Assert.assertFalse(e.getMessage(), e.getMessage().contains(dataDir.getAbsolutePath())); + } + } + } + @Test public void testLoadWithMods() throws Exception { final long writtenPoint1; diff --git a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java index b540ebb25404c..acb1382646690 100644 --- a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java +++ b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java @@ -3013,6 +3013,9 @@ private DataNodeQueryMessages() {} "Can not find %s on this machine, notice that load can only handle files on this machine."; public static final String QUERY_EXCEPTION_LOAD_TSFILE_SOURCE_PATH_S_IS_OUTSIDE_ALLOWED_DIRECTORIES_85A6019F = "Load TsFile source path %s is outside allowed directories %s."; + public static final String + QUERY_EXCEPTION_CANNOT_LOAD_FILES_BECAUSE_SPECIFIED_DIRECTORY_CONTAINS_IOTDB_DATA_B0A1B93D = + "Cannot load files because the specified directory contains IoTDB data."; public static final String QUERY_EXCEPTION_FAILED_TO_RESOLVE_CANONICAL_PATH_FOR_LOAD_TSFILE_SOURCE_09CC9AC6 = "Failed to resolve canonical path for Load TsFile source %s: %s"; public static final String QUERY_EXCEPTION_DATA_TYPE_IS_NOT_CONSISTENT_INPUT_S_REGISTERED_S_AE9DBDC0 = diff --git a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java index 0e73bd1dd967e..7fe8dcc7781f4 100644 --- a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java +++ b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java @@ -3633,6 +3633,9 @@ private DataNodeQueryMessages() {} public static final String QUERY_EXCEPTION_LOAD_TSFILE_SOURCE_PATH_S_IS_OUTSIDE_ALLOWED_DIRECTORIES_85A6019F = "加载 TsFile 的源路径 %s 位于允许目录 %s 之外。"; + public static final String + QUERY_EXCEPTION_CANNOT_LOAD_FILES_BECAUSE_SPECIFIED_DIRECTORY_CONTAINS_IOTDB_DATA_B0A1B93D = + "指定目录包含 IoTDB 数据,无法加载文件。"; public static final String QUERY_EXCEPTION_FAILED_TO_RESOLVE_CANONICAL_PATH_FOR_LOAD_TSFILE_SOURCE_09CC9AC6 = "无法解析 load TsFile source %s 的 canonical path:%s"; public static final String QUERY_EXCEPTION_DATA_TYPE_IS_NOT_CONSISTENT_INPUT_S_REGISTERED_S_AE9DBDC0 = 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 0d632ef87bf7e..31a4697da4d7f 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 @@ -48,6 +48,8 @@ import java.util.Map; import static org.apache.iotdb.commons.conf.IoTDBConstant.PATH_ROOT; +import static org.apache.iotdb.commons.conf.IoTDBConstant.SEQUENCE_FOLDER_NAME; +import static org.apache.iotdb.commons.conf.IoTDBConstant.UNSEQUENCE_FOLDER_NAME; import static org.apache.iotdb.db.storageengine.load.config.LoadTsFileConfigurator.ASYNC_LOAD_KEY; import static org.apache.iotdb.db.storageengine.load.config.LoadTsFileConfigurator.CONVERT_ON_TYPE_MISMATCH_KEY; import static org.apache.iotdb.db.storageengine.load.config.LoadTsFileConfigurator.DATABASE_LEVEL_KEY; @@ -109,6 +111,8 @@ public static List processTsFile(final File file) throws FileNotFoundExcep public static List processTsFile(final File file, final boolean validateSourcePath) throws FileNotFoundException { + final Path[] internalTsFileDirCanonicalPaths = getInternalTsFileDirCanonicalPaths(); + validateNotLoadingInternalTsFile(file, internalTsFileDirCanonicalPaths); if (validateSourcePath) { validateLoadSourcePath(file); } @@ -124,7 +128,7 @@ public static List processTsFile(final File file, final boolean validateSo .QUERY_EXCEPTION_CAN_NOT_FIND_S_ON_THIS_MACHINE_NOTICE_THAT_LOAD_CAN_ONLY_B7886C0E, file.getPath())); } - tsFiles.addAll(findAllTsFile(file, validateSourcePath)); + tsFiles.addAll(findAllTsFile(file, validateSourcePath, internalTsFileDirCanonicalPaths)); } sortTsFiles(tsFiles); return tsFiles; @@ -146,7 +150,8 @@ protected LoadTsFileStatement() { this.statementType = StatementType.MULTI_BATCH_INSERT; } - private static List findAllTsFile(File file, boolean validateSourcePath) + private static List findAllTsFile( + File file, boolean validateSourcePath, Path[] internalTsFileDirCanonicalPaths) throws FileNotFoundException { final File[] files = file.listFiles(); if (files == null) { @@ -155,13 +160,14 @@ private static List findAllTsFile(File file, boolean validateSourcePath) final List tsFiles = new ArrayList<>(); for (File nowFile : files) { + validateNotLoadingInternalTsFile(nowFile, internalTsFileDirCanonicalPaths); if (validateSourcePath) { validateLoadSourcePath(nowFile); } if (nowFile.getName().endsWith(TsFileConstant.TSFILE_SUFFIX)) { tsFiles.add(nowFile); } else if (nowFile.isDirectory()) { - tsFiles.addAll(findAllTsFile(nowFile, validateSourcePath)); + tsFiles.addAll(findAllTsFile(nowFile, validateSourcePath, internalTsFileDirCanonicalPaths)); } } return tsFiles; @@ -195,6 +201,31 @@ private static void validateLoadSourcePath(final File file) throws FileNotFoundE Arrays.toString(allowedDirs))); } + private static Path[] getInternalTsFileDirCanonicalPaths() throws FileNotFoundException { + final String[] localDataDirs = IoTDBDescriptor.getInstance().getConfig().getLocalDataDirs(); + final Path[] internalTsFileDirCanonicalPaths = new Path[localDataDirs.length * 2]; + for (int i = 0; i < localDataDirs.length; i++) { + internalTsFileDirCanonicalPaths[i * 2] = + canonicalPath(new File(localDataDirs[i], SEQUENCE_FOLDER_NAME)); + internalTsFileDirCanonicalPaths[i * 2 + 1] = + canonicalPath(new File(localDataDirs[i], UNSEQUENCE_FOLDER_NAME)); + } + return internalTsFileDirCanonicalPaths; + } + + private static void validateNotLoadingInternalTsFile( + final File file, final Path[] internalTsFileDirCanonicalPaths) throws FileNotFoundException { + final Path sourcePath = canonicalPath(file); + for (final Path internalTsFileDirCanonicalPath : internalTsFileDirCanonicalPaths) { + if (sourcePath.startsWith(internalTsFileDirCanonicalPath) + || internalTsFileDirCanonicalPath.startsWith(sourcePath)) { + throw new FileNotFoundException( + DataNodeQueryMessages + .QUERY_EXCEPTION_CANNOT_LOAD_FILES_BECAUSE_SPECIFIED_DIRECTORY_CONTAINS_IOTDB_DATA_B0A1B93D); + } + } + } + private static Path canonicalPath(final File file) throws FileNotFoundException { try { return file.getCanonicalFile().toPath(); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java index bfebf51d28191..c3ded839d1a44 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java @@ -109,6 +109,37 @@ public void testLoadSourcePathCheckCanBeDisabled() throws Exception { } } + @Test + public void testLoadInternalTsFileIsRejectedWithoutLeakingPath() throws Exception { + final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); + final String[][] originalTierDataDirs = config.getTierDataDirs(); + final boolean originalCheckEnabled = config.isLoadTsFileSourcePathCheckEnabled(); + final Path dataDir = Files.createTempDirectory("load-tsfile-internal-data"); + final Path internalTsFile = + Files.createDirectories(dataDir.resolve("sequence").resolve("root.db")).resolve("a.tsfile"); + Files.createFile(internalTsFile); + + try { + config.setTierDataDirs(new String[][] {{dataDir.toString()}}); + config.setLoadTsFileSourcePathCheckEnabled(false); + + try { + new LoadTsFileStatement(dataDir.toString()); + Assert.fail("Expected internal IoTDB data directory to be rejected."); + } catch (final FileNotFoundException e) { + Assert.assertEquals( + "Cannot load files because the specified directory contains IoTDB data.", + e.getMessage()); + Assert.assertFalse(e.getMessage().contains(dataDir.toString())); + Assert.assertFalse(e.getMessage().contains(internalTsFile.toString())); + } + } finally { + config.setTierDataDirs(originalTierDataDirs); + config.setLoadTsFileSourcePathCheckEnabled(originalCheckEnabled); + deleteRecursively(dataDir); + } + } + private static void assertLoadSourcePathRejected(final Path sourcePath) { try { new LoadTsFileStatement(sourcePath.toString()); From fdccefc505a2a55a5d4356161ffc2c2742da6f01 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Mon, 27 Jul 2026 10:52:00 +0800 Subject: [PATCH 02/10] perf(load): cache internal TsFile directories --- .../org/apache/iotdb/db/conf/IoTDBConfig.java | 22 +++++++++++++++++++ .../statement/crud/LoadTsFileStatement.java | 17 ++------------ 2 files changed, 24 insertions(+), 15 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java index 7f151e450d2ce..4f616bce44d39 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java @@ -310,6 +310,8 @@ public class IoTDBConfig { private CanonicalPaths loadTsFileAllowedDirCanonicalPaths = canonicalPaths(loadTsFileAllowedDirs); + private volatile CanonicalPaths internalTsFileDirCanonicalPaths = new CanonicalPaths(new Path[0]); + private boolean loadTsFileSourcePathCheckEnabled = false; /** Strategy of multiple directories. */ @@ -1429,6 +1431,7 @@ private void formulateFolders() { queryDir = addDataHomeDir(queryDir); sortTmpDir = addDataHomeDir(sortTmpDir); formulateDataDirs(tierDataDirs); + formulateInternalTsFileDirs(tierDataDirs); } private void formulateDataDirs(String[][] tierDataDirs) { @@ -1480,6 +1483,7 @@ void reloadDataDirs(String[][] newTierDataDirs) throws LoadConfigurationExceptio } } this.tierDataDirs = newTierDataDirs; + formulateInternalTsFileDirs(newTierDataDirs); reloadSystemMetrics(); } @@ -1556,6 +1560,10 @@ public String[] getLocalDataDirs() { .toArray(String[]::new); } + public Path[] getInternalTsFileDirCanonicalPaths() throws FileNotFoundException { + return internalTsFileDirCanonicalPaths.getPaths(); + } + public String[][] getTierDataDirs() { return tierDataDirs; } @@ -1564,6 +1572,7 @@ public String[][] getTierDataDirs() { public void setTierDataDirs(String[][] tierDataDirs) { formulateDataDirs(tierDataDirs); this.tierDataDirs = tierDataDirs; + formulateInternalTsFileDirs(tierDataDirs); // TODO(szywilliam): rewrite the logic here when ratis supports complete snapshot semantic setRatisDataRegionSnapshotDir( tierDataDirs[0][0] + File.separator + IoTDBConstant.SNAPSHOT_FOLDER_NAME); @@ -1663,6 +1672,19 @@ public void formulateLoadTsFileDirs(String[][] tierDataDirs) { this.loadTsFileDirCanonicalPaths = canonicalPaths(newLoadTsFileDirs); } + private void formulateInternalTsFileDirs(final String[][] tierDataDirs) { + final List internalTsFileDirs = new ArrayList<>(); + for (final String[] tierDataDir : tierDataDirs) { + for (final String dataDir : tierDataDir) { + if (FSUtils.isLocal(dataDir)) { + internalTsFileDirs.add(dataDir + File.separator + IoTDBConstant.SEQUENCE_FOLDER_NAME); + internalTsFileDirs.add(dataDir + File.separator + IoTDBConstant.UNSEQUENCE_FOLDER_NAME); + } + } + } + internalTsFileDirCanonicalPaths = canonicalPaths(internalTsFileDirs.toArray(new String[0])); + } + private static CanonicalPaths canonicalPaths(final String[] dirs) { final Path[] paths = new Path[dirs.length]; for (int i = 0; i < dirs.length; i++) { 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 31a4697da4d7f..595241f513abf 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 @@ -48,8 +48,6 @@ import java.util.Map; import static org.apache.iotdb.commons.conf.IoTDBConstant.PATH_ROOT; -import static org.apache.iotdb.commons.conf.IoTDBConstant.SEQUENCE_FOLDER_NAME; -import static org.apache.iotdb.commons.conf.IoTDBConstant.UNSEQUENCE_FOLDER_NAME; import static org.apache.iotdb.db.storageengine.load.config.LoadTsFileConfigurator.ASYNC_LOAD_KEY; import static org.apache.iotdb.db.storageengine.load.config.LoadTsFileConfigurator.CONVERT_ON_TYPE_MISMATCH_KEY; import static org.apache.iotdb.db.storageengine.load.config.LoadTsFileConfigurator.DATABASE_LEVEL_KEY; @@ -111,7 +109,8 @@ public static List processTsFile(final File file) throws FileNotFoundExcep public static List processTsFile(final File file, final boolean validateSourcePath) throws FileNotFoundException { - final Path[] internalTsFileDirCanonicalPaths = getInternalTsFileDirCanonicalPaths(); + final Path[] internalTsFileDirCanonicalPaths = + IoTDBDescriptor.getInstance().getConfig().getInternalTsFileDirCanonicalPaths(); validateNotLoadingInternalTsFile(file, internalTsFileDirCanonicalPaths); if (validateSourcePath) { validateLoadSourcePath(file); @@ -201,18 +200,6 @@ private static void validateLoadSourcePath(final File file) throws FileNotFoundE Arrays.toString(allowedDirs))); } - private static Path[] getInternalTsFileDirCanonicalPaths() throws FileNotFoundException { - final String[] localDataDirs = IoTDBDescriptor.getInstance().getConfig().getLocalDataDirs(); - final Path[] internalTsFileDirCanonicalPaths = new Path[localDataDirs.length * 2]; - for (int i = 0; i < localDataDirs.length; i++) { - internalTsFileDirCanonicalPaths[i * 2] = - canonicalPath(new File(localDataDirs[i], SEQUENCE_FOLDER_NAME)); - internalTsFileDirCanonicalPaths[i * 2 + 1] = - canonicalPath(new File(localDataDirs[i], UNSEQUENCE_FOLDER_NAME)); - } - return internalTsFileDirCanonicalPaths; - } - private static void validateNotLoadingInternalTsFile( final File file, final Path[] internalTsFileDirCanonicalPaths) throws FileNotFoundException { final Path sourcePath = canonicalPath(file); From 6bed0ebba9aa20078f49177861c8df0f2009a15c Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Tue, 28 Jul 2026 17:34:53 +0800 Subject: [PATCH 03/10] fix(load): reject all local data directories --- .../org/apache/iotdb/db/conf/IoTDBConfig.java | 21 +++++++------- .../statement/crud/LoadTsFileStatement.java | 22 +++++++-------- .../crud/LoadTsFileStatementTest.java | 28 +++++++++++++++++-- 3 files changed, 47 insertions(+), 24 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java index 4f616bce44d39..945f0edbbb181 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java @@ -310,7 +310,7 @@ public class IoTDBConfig { private CanonicalPaths loadTsFileAllowedDirCanonicalPaths = canonicalPaths(loadTsFileAllowedDirs); - private volatile CanonicalPaths internalTsFileDirCanonicalPaths = new CanonicalPaths(new Path[0]); + private volatile CanonicalPaths internalDataDirCanonicalPaths = new CanonicalPaths(new Path[0]); private boolean loadTsFileSourcePathCheckEnabled = false; @@ -1431,7 +1431,7 @@ private void formulateFolders() { queryDir = addDataHomeDir(queryDir); sortTmpDir = addDataHomeDir(sortTmpDir); formulateDataDirs(tierDataDirs); - formulateInternalTsFileDirs(tierDataDirs); + formulateInternalDataDirs(tierDataDirs); } private void formulateDataDirs(String[][] tierDataDirs) { @@ -1483,7 +1483,7 @@ void reloadDataDirs(String[][] newTierDataDirs) throws LoadConfigurationExceptio } } this.tierDataDirs = newTierDataDirs; - formulateInternalTsFileDirs(newTierDataDirs); + formulateInternalDataDirs(newTierDataDirs); reloadSystemMetrics(); } @@ -1560,8 +1560,8 @@ public String[] getLocalDataDirs() { .toArray(String[]::new); } - public Path[] getInternalTsFileDirCanonicalPaths() throws FileNotFoundException { - return internalTsFileDirCanonicalPaths.getPaths(); + public Path[] getInternalDataDirCanonicalPaths() throws FileNotFoundException { + return internalDataDirCanonicalPaths.getPaths(); } public String[][] getTierDataDirs() { @@ -1572,7 +1572,7 @@ public String[][] getTierDataDirs() { public void setTierDataDirs(String[][] tierDataDirs) { formulateDataDirs(tierDataDirs); this.tierDataDirs = tierDataDirs; - formulateInternalTsFileDirs(tierDataDirs); + formulateInternalDataDirs(tierDataDirs); // TODO(szywilliam): rewrite the logic here when ratis supports complete snapshot semantic setRatisDataRegionSnapshotDir( tierDataDirs[0][0] + File.separator + IoTDBConstant.SNAPSHOT_FOLDER_NAME); @@ -1672,17 +1672,16 @@ public void formulateLoadTsFileDirs(String[][] tierDataDirs) { this.loadTsFileDirCanonicalPaths = canonicalPaths(newLoadTsFileDirs); } - private void formulateInternalTsFileDirs(final String[][] tierDataDirs) { - final List internalTsFileDirs = new ArrayList<>(); + private void formulateInternalDataDirs(final String[][] tierDataDirs) { + final List internalDataDirs = new ArrayList<>(); for (final String[] tierDataDir : tierDataDirs) { for (final String dataDir : tierDataDir) { if (FSUtils.isLocal(dataDir)) { - internalTsFileDirs.add(dataDir + File.separator + IoTDBConstant.SEQUENCE_FOLDER_NAME); - internalTsFileDirs.add(dataDir + File.separator + IoTDBConstant.UNSEQUENCE_FOLDER_NAME); + internalDataDirs.add(dataDir); } } } - internalTsFileDirCanonicalPaths = canonicalPaths(internalTsFileDirs.toArray(new String[0])); + internalDataDirCanonicalPaths = canonicalPaths(internalDataDirs.toArray(new String[0])); } private static CanonicalPaths canonicalPaths(final String[] dirs) { 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 595241f513abf..d76e161d6a6dd 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 @@ -109,9 +109,9 @@ public static List processTsFile(final File file) throws FileNotFoundExcep public static List processTsFile(final File file, final boolean validateSourcePath) throws FileNotFoundException { - final Path[] internalTsFileDirCanonicalPaths = - IoTDBDescriptor.getInstance().getConfig().getInternalTsFileDirCanonicalPaths(); - validateNotLoadingInternalTsFile(file, internalTsFileDirCanonicalPaths); + final Path[] internalDataDirCanonicalPaths = + IoTDBDescriptor.getInstance().getConfig().getInternalDataDirCanonicalPaths(); + validateNotLoadingInternalTsFile(file, internalDataDirCanonicalPaths); if (validateSourcePath) { validateLoadSourcePath(file); } @@ -127,7 +127,7 @@ public static List processTsFile(final File file, final boolean validateSo .QUERY_EXCEPTION_CAN_NOT_FIND_S_ON_THIS_MACHINE_NOTICE_THAT_LOAD_CAN_ONLY_B7886C0E, file.getPath())); } - tsFiles.addAll(findAllTsFile(file, validateSourcePath, internalTsFileDirCanonicalPaths)); + tsFiles.addAll(findAllTsFile(file, validateSourcePath, internalDataDirCanonicalPaths)); } sortTsFiles(tsFiles); return tsFiles; @@ -150,7 +150,7 @@ protected LoadTsFileStatement() { } private static List findAllTsFile( - File file, boolean validateSourcePath, Path[] internalTsFileDirCanonicalPaths) + File file, boolean validateSourcePath, Path[] internalDataDirCanonicalPaths) throws FileNotFoundException { final File[] files = file.listFiles(); if (files == null) { @@ -159,14 +159,14 @@ private static List findAllTsFile( final List tsFiles = new ArrayList<>(); for (File nowFile : files) { - validateNotLoadingInternalTsFile(nowFile, internalTsFileDirCanonicalPaths); + validateNotLoadingInternalTsFile(nowFile, internalDataDirCanonicalPaths); if (validateSourcePath) { validateLoadSourcePath(nowFile); } if (nowFile.getName().endsWith(TsFileConstant.TSFILE_SUFFIX)) { tsFiles.add(nowFile); } else if (nowFile.isDirectory()) { - tsFiles.addAll(findAllTsFile(nowFile, validateSourcePath, internalTsFileDirCanonicalPaths)); + tsFiles.addAll(findAllTsFile(nowFile, validateSourcePath, internalDataDirCanonicalPaths)); } } return tsFiles; @@ -201,11 +201,11 @@ private static void validateLoadSourcePath(final File file) throws FileNotFoundE } private static void validateNotLoadingInternalTsFile( - final File file, final Path[] internalTsFileDirCanonicalPaths) throws FileNotFoundException { + final File file, final Path[] internalDataDirCanonicalPaths) throws FileNotFoundException { final Path sourcePath = canonicalPath(file); - for (final Path internalTsFileDirCanonicalPath : internalTsFileDirCanonicalPaths) { - if (sourcePath.startsWith(internalTsFileDirCanonicalPath) - || internalTsFileDirCanonicalPath.startsWith(sourcePath)) { + for (final Path internalDataDirCanonicalPath : internalDataDirCanonicalPaths) { + if (sourcePath.startsWith(internalDataDirCanonicalPath) + || internalDataDirCanonicalPath.startsWith(sourcePath)) { throw new FileNotFoundException( DataNodeQueryMessages .QUERY_EXCEPTION_CANNOT_LOAD_FILES_BECAUSE_SPECIFIED_DIRECTORY_CONTAINS_IOTDB_DATA_B0A1B93D); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java index c3ded839d1a44..0cc3c627dd6f4 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java @@ -116,7 +116,7 @@ public void testLoadInternalTsFileIsRejectedWithoutLeakingPath() throws Exceptio final boolean originalCheckEnabled = config.isLoadTsFileSourcePathCheckEnabled(); final Path dataDir = Files.createTempDirectory("load-tsfile-internal-data"); final Path internalTsFile = - Files.createDirectories(dataDir.resolve("sequence").resolve("root.db")).resolve("a.tsfile"); + Files.createDirectories(dataDir.resolve("pipe-hardlink")).resolve("a.tsfile"); Files.createFile(internalTsFile); try { @@ -124,7 +124,7 @@ public void testLoadInternalTsFileIsRejectedWithoutLeakingPath() throws Exceptio config.setLoadTsFileSourcePathCheckEnabled(false); try { - new LoadTsFileStatement(dataDir.toString()); + new LoadTsFileStatement(internalTsFile.toString()); Assert.fail("Expected internal IoTDB data directory to be rejected."); } catch (final FileNotFoundException e) { Assert.assertEquals( @@ -140,6 +140,30 @@ public void testLoadInternalTsFileIsRejectedWithoutLeakingPath() throws Exceptio } } + @Test + public void testLoadPipeReceiverTsFileOutsideDataDirIsAllowed() throws Exception { + final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); + final String[][] originalTierDataDirs = config.getTierDataDirs(); + final boolean originalCheckEnabled = config.isLoadTsFileSourcePathCheckEnabled(); + final Path dataDir = Files.createTempDirectory("load-tsfile-internal-data"); + final Path pipeReceiverDir = Files.createTempDirectory("load-tsfile-pipe-receiver"); + final Path pipeReceiverTsFile = Files.createFile(pipeReceiverDir.resolve("a.tsfile")); + + try { + config.setTierDataDirs(new String[][] {{dataDir.toString()}}); + config.setLoadTsFileSourcePathCheckEnabled(false); + + final LoadTsFileStatement statement = new LoadTsFileStatement(pipeReceiverTsFile.toString()); + Assert.assertEquals(1, statement.getTsFiles().size()); + Assert.assertEquals(pipeReceiverTsFile.toFile(), statement.getTsFiles().get(0)); + } finally { + config.setTierDataDirs(originalTierDataDirs); + config.setLoadTsFileSourcePathCheckEnabled(originalCheckEnabled); + deleteRecursively(dataDir); + deleteRecursively(pipeReceiverDir); + } + } + private static void assertLoadSourcePathRejected(final Path sourcePath) { try { new LoadTsFileStatement(sourcePath.toString()); From 1c38947e7ea0e2b1f663c3c4f1577ea16b673fef Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Wed, 29 Jul 2026 09:57:58 +0800 Subject: [PATCH 04/10] test(load): cover pipe receiver outside data directory --- .../statement/crud/LoadTsFileStatementTest.java | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java index 0cc3c627dd6f4..2d29134658bc3 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java @@ -114,7 +114,8 @@ public void testLoadInternalTsFileIsRejectedWithoutLeakingPath() throws Exceptio final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); final String[][] originalTierDataDirs = config.getTierDataDirs(); final boolean originalCheckEnabled = config.isLoadTsFileSourcePathCheckEnabled(); - final Path dataDir = Files.createTempDirectory("load-tsfile-internal-data"); + final Path dataNodeDir = Files.createTempDirectory("load-tsfile-datanode"); + final Path dataDir = dataNodeDir.resolve("data"); final Path internalTsFile = Files.createDirectories(dataDir.resolve("pipe-hardlink")).resolve("a.tsfile"); Files.createFile(internalTsFile); @@ -136,7 +137,7 @@ public void testLoadInternalTsFileIsRejectedWithoutLeakingPath() throws Exceptio } finally { config.setTierDataDirs(originalTierDataDirs); config.setLoadTsFileSourcePathCheckEnabled(originalCheckEnabled); - deleteRecursively(dataDir); + deleteRecursively(dataNodeDir); } } @@ -145,8 +146,10 @@ public void testLoadPipeReceiverTsFileOutsideDataDirIsAllowed() throws Exception final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); final String[][] originalTierDataDirs = config.getTierDataDirs(); final boolean originalCheckEnabled = config.isLoadTsFileSourcePathCheckEnabled(); - final Path dataDir = Files.createTempDirectory("load-tsfile-internal-data"); - final Path pipeReceiverDir = Files.createTempDirectory("load-tsfile-pipe-receiver"); + final Path dataNodeDir = Files.createTempDirectory("load-tsfile-datanode"); + final Path dataDir = dataNodeDir.resolve("data"); + final Path pipeReceiverDir = + Files.createDirectories(dataNodeDir.resolve("system").resolve("pipe").resolve("receiver")); final Path pipeReceiverTsFile = Files.createFile(pipeReceiverDir.resolve("a.tsfile")); try { @@ -159,8 +162,7 @@ public void testLoadPipeReceiverTsFileOutsideDataDirIsAllowed() throws Exception } finally { config.setTierDataDirs(originalTierDataDirs); config.setLoadTsFileSourcePathCheckEnabled(originalCheckEnabled); - deleteRecursively(dataDir); - deleteRecursively(pipeReceiverDir); + deleteRecursively(dataNodeDir); } } From 828791b7db2400348f14ae8c0e7f255b29c4a040 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Wed, 29 Jul 2026 10:05:11 +0800 Subject: [PATCH 05/10] fix(load): reject IoTDB data directory --- .../org/apache/iotdb/db/conf/IoTDBConfig.java | 1 + .../protocol/legacy/loader/TsFileLoader.java | 2 +- .../thrift/IoTDBDataNodeReceiver.java | 2 +- .../statement/crud/LoadTsFileStatement.java | 45 ++++++++++++++----- .../crud/LoadTsFileStatementTest.java | 3 ++ 5 files changed, 41 insertions(+), 12 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java index 945f0edbbb181..b33dff6ad609d 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java @@ -1674,6 +1674,7 @@ public void formulateLoadTsFileDirs(String[][] tierDataDirs) { private void formulateInternalDataDirs(final String[][] tierDataDirs) { final List internalDataDirs = new ArrayList<>(); + internalDataDirs.add(addDataHomeDir("data")); for (final String[] tierDataDir : tierDataDirs) { for (final String dataDir : tierDataDir) { if (FSUtils.isLocal(dataDir)) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/legacy/loader/TsFileLoader.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/legacy/loader/TsFileLoader.java index bf6f717be01b9..b91547137fbdf 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/legacy/loader/TsFileLoader.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/legacy/loader/TsFileLoader.java @@ -53,7 +53,7 @@ public TsFileLoader(File tsFile, String database) { @Override public void load(final SessionInfo sessionInfo) { try { - LoadTsFileStatement statement = LoadTsFileStatement.createUnchecked(tsFile.getAbsolutePath()); + LoadTsFileStatement statement = LoadTsFileStatement.createForPipe(tsFile.getAbsolutePath()); statement.setDeleteAfterLoad(true); statement.setConvertOnTypeMismatch(true); statement.setDatabaseLevel(parseSgLevel()); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java index 4cca55d705f63..dabd2dd2c56cc 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java @@ -676,7 +676,7 @@ static LoadTsFileStatement buildLoadTsFileStatementForSync( final boolean validateTsFile, final boolean shouldConvertDataTypeOnTypeMismatch) throws FileNotFoundException { - final LoadTsFileStatement statement = LoadTsFileStatement.createUnchecked(fileAbsolutePath); + final LoadTsFileStatement statement = LoadTsFileStatement.createForPipe(fileAbsolutePath); statement.setDeleteAfterLoad(true); statement.setConvertOnTypeMismatch(shouldConvertDataTypeOnTypeMismatch); statement.setVerifySchema(validateTsFile || shouldConvertDataTypeOnTypeMismatch); 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 d76e161d6a6dd..e81ead9d12a25 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 @@ -78,14 +78,19 @@ public class LoadTsFileStatement extends Statement { private boolean needDecode4TimeColumn; public LoadTsFileStatement(String filePath) throws FileNotFoundException { - this(filePath, true); + this(filePath, true, true); } public static LoadTsFileStatement createUnchecked(String filePath) throws FileNotFoundException { - return new LoadTsFileStatement(filePath, false); + return new LoadTsFileStatement(filePath, false, true); } - private LoadTsFileStatement(String filePath, boolean validateSourcePath) + public static LoadTsFileStatement createForPipe(String filePath) throws FileNotFoundException { + return new LoadTsFileStatement(filePath, false, false); + } + + private LoadTsFileStatement( + String filePath, boolean validateSourcePath, boolean validateInternalDataDir) throws FileNotFoundException { this.file = new File(filePath).getAbsoluteFile(); this.databaseLevel = IoTDBDescriptor.getInstance().getConfig().getDefaultDatabaseLevel(); @@ -96,7 +101,7 @@ private LoadTsFileStatement(String filePath, boolean validateSourcePath) IoTDBDescriptor.getInstance().getConfig().getLoadTabletConversionThresholdBytes(); this.autoCreateDatabase = IoTDBDescriptor.getInstance().getConfig().isAutoCreateSchemaEnabled(); - this.tsFiles = processTsFile(file, validateSourcePath); + this.tsFiles = processTsFile(file, validateSourcePath, validateInternalDataDir); this.resources = new ArrayList<>(); this.writePointCountList = new ArrayList<>(); this.isTableModel = new ArrayList<>(Collections.nCopies(this.tsFiles.size(), false)); @@ -104,14 +109,22 @@ private LoadTsFileStatement(String filePath, boolean validateSourcePath) } public static List processTsFile(final File file) throws FileNotFoundException { - return processTsFile(file, true); + return processTsFile(file, true, true); } public static List processTsFile(final File file, final boolean validateSourcePath) throws FileNotFoundException { + return processTsFile(file, validateSourcePath, true); + } + + private static List processTsFile( + final File file, final boolean validateSourcePath, final boolean validateInternalDataDir) + throws FileNotFoundException { final Path[] internalDataDirCanonicalPaths = IoTDBDescriptor.getInstance().getConfig().getInternalDataDirCanonicalPaths(); - validateNotLoadingInternalTsFile(file, internalDataDirCanonicalPaths); + if (validateInternalDataDir) { + validateNotLoadingInternalTsFile(file, internalDataDirCanonicalPaths); + } if (validateSourcePath) { validateLoadSourcePath(file); } @@ -127,7 +140,9 @@ public static List processTsFile(final File file, final boolean validateSo .QUERY_EXCEPTION_CAN_NOT_FIND_S_ON_THIS_MACHINE_NOTICE_THAT_LOAD_CAN_ONLY_B7886C0E, file.getPath())); } - tsFiles.addAll(findAllTsFile(file, validateSourcePath, internalDataDirCanonicalPaths)); + tsFiles.addAll( + findAllTsFile( + file, validateSourcePath, validateInternalDataDir, internalDataDirCanonicalPaths)); } sortTsFiles(tsFiles); return tsFiles; @@ -150,7 +165,10 @@ protected LoadTsFileStatement() { } private static List findAllTsFile( - File file, boolean validateSourcePath, Path[] internalDataDirCanonicalPaths) + File file, + boolean validateSourcePath, + boolean validateInternalDataDir, + Path[] internalDataDirCanonicalPaths) throws FileNotFoundException { final File[] files = file.listFiles(); if (files == null) { @@ -159,14 +177,21 @@ private static List findAllTsFile( final List tsFiles = new ArrayList<>(); for (File nowFile : files) { - validateNotLoadingInternalTsFile(nowFile, internalDataDirCanonicalPaths); + if (validateInternalDataDir) { + validateNotLoadingInternalTsFile(nowFile, internalDataDirCanonicalPaths); + } if (validateSourcePath) { validateLoadSourcePath(nowFile); } if (nowFile.getName().endsWith(TsFileConstant.TSFILE_SUFFIX)) { tsFiles.add(nowFile); } else if (nowFile.isDirectory()) { - tsFiles.addAll(findAllTsFile(nowFile, validateSourcePath, internalDataDirCanonicalPaths)); + tsFiles.addAll( + findAllTsFile( + nowFile, + validateSourcePath, + validateInternalDataDir, + internalDataDirCanonicalPaths)); } } return tsFiles; diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java index 2d29134658bc3..fac3227f4b5e3 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java @@ -134,6 +134,9 @@ public void testLoadInternalTsFileIsRejectedWithoutLeakingPath() throws Exceptio Assert.assertFalse(e.getMessage().contains(dataDir.toString())); Assert.assertFalse(e.getMessage().contains(internalTsFile.toString())); } + + Assert.assertEquals( + 1, LoadTsFileStatement.createForPipe(internalTsFile.toString()).getTsFiles().size()); } finally { config.setTierDataDirs(originalTierDataDirs); config.setLoadTsFileSourcePathCheckEnabled(originalCheckEnabled); From 962a829ad4eefee3739240d84793d0fef13c5f2f Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Wed, 29 Jul 2026 11:10:48 +0800 Subject: [PATCH 06/10] fix(pipe): preserve load source during type conversion --- .../plan/analyze/load/LoadTsFileAnalyzer.java | 22 ++++++++++++++----- .../plan/relational/sql/ast/LoadTsFile.java | 19 +++++++++++----- .../statement/crud/LoadTsFileStatement.java | 4 ++++ 3 files changed, 34 insertions(+), 11 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzer.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzer.java index 48d391658403a..3df5116753dad 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzer.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzer.java @@ -488,15 +488,20 @@ private boolean handleSingleMiniFile(final int i) throws FileNotFoundException { isTableModelTsFile.get(i) ? loadTsFileDataTypeConverter .convertForTableModel( - LoadTsFile.createUnchecked( - null, tsFiles.get(i).getPath(), Collections.emptyMap()) + (isGeneratedByPipe + ? LoadTsFile.createForPipe( + null, tsFiles.get(i).getPath(), Collections.emptyMap()) + : LoadTsFile.createUnchecked( + null, tsFiles.get(i).getPath(), Collections.emptyMap())) .setDatabase(databaseForTableData) .setDeleteAfterLoad(isDeleteAfterLoad) .setConvertOnTypeMismatch(isConvertOnTypeMismatch)) .orElse(null) : loadTsFileDataTypeConverter .convertForTreeModel( - LoadTsFileStatement.createUnchecked(tsFiles.get(i).getPath()) + (isGeneratedByPipe + ? LoadTsFileStatement.createForPipe(tsFiles.get(i).getPath()) + : LoadTsFileStatement.createUnchecked(tsFiles.get(i).getPath())) .setDeleteAfterLoad(isDeleteAfterLoad) .setConvertOnTypeMismatch(isConvertOnTypeMismatch)) .orElse(null); @@ -781,15 +786,20 @@ private void executeTabletConversionOnException( isTableModelTsFile.get(i) ? loadTsFileDataTypeConverter .convertForTableModel( - LoadTsFile.createUnchecked( - null, tsFiles.get(i).getPath(), Collections.emptyMap()) + (isGeneratedByPipe + ? LoadTsFile.createForPipe( + null, tsFiles.get(i).getPath(), Collections.emptyMap()) + : LoadTsFile.createUnchecked( + null, tsFiles.get(i).getPath(), Collections.emptyMap())) .setDatabase(databaseForTableData) .setDeleteAfterLoad(isDeleteAfterLoad) .setConvertOnTypeMismatch(isConvertOnTypeMismatch)) .orElse(null) : loadTsFileDataTypeConverter .convertForTreeModel( - LoadTsFileStatement.createUnchecked(tsFiles.get(i).getPath()) + (isGeneratedByPipe + ? LoadTsFileStatement.createForPipe(tsFiles.get(i).getPath()) + : LoadTsFileStatement.createUnchecked(tsFiles.get(i).getPath())) .setDeleteAfterLoad(isDeleteAfterLoad) .setConvertOnTypeMismatch(isConvertOnTypeMismatch)) .orElse(null); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java index 28ae4d7f08163..ab34915d3538c 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java @@ -73,19 +73,25 @@ public class LoadTsFile extends Statement { private boolean needDecode4TimeColumn; public LoadTsFile(NodeLocation location, String filePath, Map loadAttributes) { - this(location, filePath, loadAttributes, true); + this(location, filePath, loadAttributes, true, true); } public static LoadTsFile createUnchecked( NodeLocation location, String filePath, Map loadAttributes) { - return new LoadTsFile(location, filePath, loadAttributes, false); + return new LoadTsFile(location, filePath, loadAttributes, false, true); + } + + public static LoadTsFile createForPipe( + NodeLocation location, String filePath, Map loadAttributes) { + return new LoadTsFile(location, filePath, loadAttributes, false, false); } private LoadTsFile( NodeLocation location, String filePath, Map loadAttributes, - boolean validateSourcePath) { + boolean validateSourcePath, + boolean validateInternalDataDir) { super(location); this.filePath = requireNonNull(filePath, DataNodeQueryMessages.EXCEPTION_FILEPATH_IS_NULL_84CE8A66); @@ -103,8 +109,11 @@ private LoadTsFile( try { this.tsFiles = - org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement.processTsFile( - new File(filePath), validateSourcePath); + validateInternalDataDir + ? org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement + .processTsFile(new File(filePath), validateSourcePath) + : org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement + .processTsFileForPipe(new File(filePath)); this.resources = new ArrayList<>(); this.writePointCountList = new ArrayList<>(); this.isTableModel = new ArrayList<>(Collections.nCopies(this.tsFiles.size(), true)); 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 e81ead9d12a25..f6c50e5fc56e8 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 @@ -117,6 +117,10 @@ public static List processTsFile(final File file, final boolean validateSo return processTsFile(file, validateSourcePath, true); } + public static List processTsFileForPipe(final File file) throws FileNotFoundException { + return processTsFile(file, false, false); + } + private static List processTsFile( final File file, final boolean validateSourcePath, final boolean validateInternalDataDir) throws FileNotFoundException { From 9c32e1bfa75a71c06c5e2e747d8f2f41d51704cd Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Wed, 29 Jul 2026 14:46:20 +0800 Subject: [PATCH 07/10] fix(pipe): preserve load source during scheduler retry --- .../plan/scheduler/load/LoadTsFileScheduler.java | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java index d3beefc0802b6..6f3c1f428f7f5 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java @@ -625,7 +625,10 @@ private void convertFailedTsFilesToTabletsAndRetry() { failedNode.isTableModel() ? loadTsFileDataTypeConverter .convertForTableModel( - LoadTsFile.createUnchecked(null, filePath, Collections.emptyMap()) + (isGeneratedByPipe + ? LoadTsFile.createForPipe(null, filePath, Collections.emptyMap()) + : LoadTsFile.createUnchecked( + null, filePath, Collections.emptyMap())) .setDatabase(failedNode.getDatabase()) .setDeleteAfterLoad(failedNode.isDeleteAfterLoad()) .setConvertOnTypeMismatch(true)) @@ -684,7 +687,9 @@ private LoadTsFileStatement buildRetryTreeLoadStatement( final String filePath, final boolean deleteAfterLoad, final String database) throws FileNotFoundException { final LoadTsFileStatement statement = - LoadTsFileStatement.createUnchecked(filePath) + (isGeneratedByPipe + ? LoadTsFileStatement.createForPipe(filePath) + : LoadTsFileStatement.createUnchecked(filePath)) .setDeleteAfterLoad(deleteAfterLoad) .setConvertOnTypeMismatch(true); if (database != null) { From fcba830c2e64333dcf36f69d99686eed83406506 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Wed, 29 Jul 2026 14:51:11 +0800 Subject: [PATCH 08/10] fix(pipe): preserve load source during active load --- .../db/queryengine/plan/relational/sql/ast/LoadTsFile.java | 4 +++- .../queryengine/plan/statement/crud/LoadTsFileStatement.java | 4 +++- .../db/storageengine/load/active/ActiveLoadTsFileLoader.java | 4 +++- 3 files changed, 9 insertions(+), 3 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java index ab34915d3538c..c92eb6dce43f6 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java @@ -307,7 +307,9 @@ public List getSubStatements() { final Map properties = this.loadAttributes; final LoadTsFile subStatement = - LoadTsFile.createUnchecked(getLocation().orElse(null), filePath, properties); + isGeneratedByPipe + ? LoadTsFile.createForPipe(getLocation().orElse(null), filePath, properties) + : LoadTsFile.createUnchecked(getLocation().orElse(null), filePath, properties); // Copy all configuration properties subStatement.databaseLevel = this.databaseLevel; 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 f6c50e5fc56e8..dcd2252d165d9 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 @@ -529,7 +529,9 @@ public List getPaths() { loadAttributes.put(PIPE_GENERATED_KEY, String.valueOf(true)); } - return LoadTsFile.createUnchecked(null, file.getAbsolutePath(), loadAttributes); + return isGeneratedByPipe + ? LoadTsFile.createForPipe(null, file.getAbsolutePath(), loadAttributes) + : LoadTsFile.createUnchecked(null, file.getAbsolutePath(), loadAttributes); } @Override diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java index ff5222e575440..0451c1812450b 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java @@ -248,7 +248,9 @@ private TSStatus loadTsFile( throws FileNotFoundException { final File tsFile = new File(entry.getFile()); final LoadTsFileStatement statement = - LoadTsFileStatement.createUnchecked(tsFile.getAbsolutePath()); + entry.isGeneratedByPipe() + ? LoadTsFileStatement.createForPipe(tsFile.getAbsolutePath()) + : LoadTsFileStatement.createUnchecked(tsFile.getAbsolutePath()); final List files = statement.getTsFiles(); statement.setDeleteAfterLoad(true); From 6bb600b95ccc7292a32415e38cfadb5ae0177b1f Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Wed, 29 Jul 2026 14:52:05 +0800 Subject: [PATCH 09/10] fix(pipe): keep active load source validation --- .../db/storageengine/load/active/ActiveLoadTsFileLoader.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java index 0451c1812450b..ff5222e575440 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java @@ -248,9 +248,7 @@ private TSStatus loadTsFile( throws FileNotFoundException { final File tsFile = new File(entry.getFile()); final LoadTsFileStatement statement = - entry.isGeneratedByPipe() - ? LoadTsFileStatement.createForPipe(tsFile.getAbsolutePath()) - : LoadTsFileStatement.createUnchecked(tsFile.getAbsolutePath()); + LoadTsFileStatement.createUnchecked(tsFile.getAbsolutePath()); final List files = statement.getTsFiles(); statement.setDeleteAfterLoad(true); From bc4f3e8f4efa1f1fbfc9d1cfa0438f8a06a5156f Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Fri, 31 Jul 2026 10:31:56 +0800 Subject: [PATCH 10/10] style(load): replace FQCN with LoadTsFileStatement import Co-authored-by: Cursor --- .../db/queryengine/plan/relational/sql/ast/LoadTsFile.java | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java index c92eb6dce43f6..4ef173d3d25d5 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java @@ -27,6 +27,7 @@ import org.apache.iotdb.commons.queryengine.plan.relational.sql.ast.Statement; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.i18n.DataNodeQueryMessages; +import org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement; import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource; import org.apache.iotdb.db.storageengine.load.config.LoadTsFileConfigurator; @@ -110,10 +111,8 @@ private LoadTsFile( try { this.tsFiles = validateInternalDataDir - ? org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement - .processTsFile(new File(filePath), validateSourcePath) - : org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement - .processTsFileForPipe(new File(filePath)); + ? LoadTsFileStatement.processTsFile(new File(filePath), validateSourcePath) + : LoadTsFileStatement.processTsFileForPipe(new File(filePath)); this.resources = new ArrayList<>(); this.writePointCountList = new ArrayList<>(); this.isTableModel = new ArrayList<>(Collections.nCopies(this.tsFiles.size(), true));