This is an automated email from the ASF dual-hosted git repository.

Caideyipi pushed a commit to branch dev/1.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/dev/1.3 by this push:
     new 3f54368ac9e [to dev/1.3] Fix pipe tree database creation on receiver 
(#17998)
3f54368ac9e is described below

commit 3f54368ac9e1ab2b53949a47670aabbae61f92ef
Author: Caideyipi <[email protected]>
AuthorDate: Wed Aug 19 10:38:19 2026 +0800

    [to dev/1.3] Fix pipe tree database creation on receiver (#17998)
    
    * Fix pipe tree database creation on receiver (#17991)
    
    (cherry picked from commit 28c4e68a6c40fb8210eab994f50c647d40ba17a3)
    
    * Fix LoadTsFileSchedulerTest import order
---
 .../protocol/thrift/IoTDBDataNodeReceiver.java     | 106 +++++++++++++++++++++
 .../queryengine/plan/planner/TreeModelPlanner.java |   7 +-
 .../plan/scheduler/load/LoadTsFileScheduler.java   |  26 ++++-
 .../protocol/thrift/IoTDBDataNodeReceiverTest.java |  19 ++++
 .../scheduler/load/LoadTsFileSchedulerTest.java    |  34 ++++++-
 5 files changed, 186 insertions(+), 6 deletions(-)

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 b7da3775787..9439d81f6bc 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
@@ -22,6 +22,8 @@ package org.apache.iotdb.db.pipe.receiver.protocol.thrift;
 import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.commons.conf.IoTDBConstant;
 import org.apache.iotdb.commons.exception.IllegalPathException;
+import org.apache.iotdb.commons.exception.IoTDBException;
+import org.apache.iotdb.commons.exception.IoTDBRuntimeException;
 import 
org.apache.iotdb.commons.exception.pipe.PipeRuntimeOutOfMemoryCriticalException;
 import org.apache.iotdb.commons.log.LoggerPeriodicalLogReducer;
 import org.apache.iotdb.commons.path.PartialPath;
@@ -73,7 +75,9 @@ import 
org.apache.iotdb.db.queryengine.common.header.ColumnHeaderConstant;
 import org.apache.iotdb.db.queryengine.plan.Coordinator;
 import org.apache.iotdb.db.queryengine.plan.analyze.ClusterPartitionFetcher;
 import 
org.apache.iotdb.db.queryengine.plan.analyze.schema.ClusterSchemaFetcher;
+import org.apache.iotdb.db.queryengine.plan.execution.config.ConfigTaskResult;
 import 
org.apache.iotdb.db.queryengine.plan.execution.config.executor.ClusterConfigTaskExecutor;
+import 
org.apache.iotdb.db.queryengine.plan.execution.config.metadata.DatabaseSchemaTask;
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.write.view.AlterLogicalViewNode;
 import org.apache.iotdb.db.queryengine.plan.statement.Statement;
 import org.apache.iotdb.db.queryengine.plan.statement.StatementType;
@@ -96,6 +100,7 @@ import org.apache.iotdb.rpc.TSStatusCode;
 import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq;
 import org.apache.iotdb.service.rpc.thrift.TPipeTransferResp;
 
+import com.google.common.util.concurrent.ListenableFuture;
 import org.apache.tsfile.utils.Pair;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -112,6 +117,8 @@ import java.util.Map;
 import java.util.Objects;
 import java.util.Optional;
 import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ExecutionException;
 import java.util.concurrent.atomic.AtomicLong;
 import java.util.concurrent.atomic.AtomicReference;
 import java.util.stream.Collectors;
@@ -153,6 +160,15 @@ public class IoTDBDataNodeReceiver extends 
IoTDBFileReceiver {
   private static final SessionManager SESSION_MANAGER = 
SessionManager.getInstance();
 
   private PipeMemoryBlock allocatedMemoryBlock;
+  private final Set<String> autoCreatedTreeDatabases = 
ConcurrentHashMap.newKeySet();
+  private final Set<String> conflictedTreeDatabases = 
ConcurrentHashMap.newKeySet();
+
+  private enum TreeDatabaseCreationResult {
+    SKIPPED,
+    CREATED_OR_EXISTED,
+    CONFLICTED
+  }
+
   private final List<PipeMemoryBlock> allocatedSliceMemoryBlocks = new 
ArrayList<>();
 
   static {
@@ -1045,6 +1061,11 @@ public class IoTDBDataNodeReceiver extends 
IoTDBFileReceiver {
       return RpcUtils.getStatus(status.getCode(), status.getMessage());
     }
 
+    if (autoCreateTreeDatabaseIfNecessary(getTreeDatabaseName(statement))
+        == TreeDatabaseCreationResult.CONFLICTED) {
+      clearTreeDatabaseName(statement);
+    }
+
     return Coordinator.getInstance()
         .executeForTreeModel(
             shouldMarkAsPipeRequest.get() ? new 
PipeEnrichedStatement(statement) : statement,
@@ -1059,6 +1080,91 @@ public class IoTDBDataNodeReceiver extends 
IoTDBFileReceiver {
         .status;
   }
 
+  private TreeDatabaseCreationResult autoCreateTreeDatabaseIfNecessary(final 
String database) {
+    if (database == null
+        || LoadTsFileStatement.getDatabaseLevelByTreeDatabase(database) == null
+        || 
!IoTDBDescriptor.getInstance().getConfig().isAutoCreateSchemaEnabled()) {
+      return TreeDatabaseCreationResult.SKIPPED;
+    }
+    if (autoCreatedTreeDatabases.contains(database)) {
+      return TreeDatabaseCreationResult.CREATED_OR_EXISTED;
+    }
+    if (conflictedTreeDatabases.contains(database)) {
+      return TreeDatabaseCreationResult.CONFLICTED;
+    }
+
+    try {
+      final DatabaseSchemaStatement statement =
+          new 
DatabaseSchemaStatement(DatabaseSchemaStatement.DatabaseSchemaStatementType.CREATE);
+      statement.setDatabasePath(new PartialPath(database));
+      statement.setEnablePrintExceptionLog(false);
+
+      final TSStatus permissionStatus = 
statement.checkPermissionBeforeProcess(username);
+      if (permissionStatus.getCode() != 
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+        throw new PipeException(permissionStatus.getMessage());
+      }
+
+      final DatabaseSchemaTask task = new DatabaseSchemaTask(statement);
+      final ListenableFuture<ConfigTaskResult> future =
+          task.execute(ClusterConfigTaskExecutor.getInstance());
+      final ConfigTaskResult result = future.get();
+      final int statusCode = result.getStatusCode().getStatusCode();
+      if (statusCode == TSStatusCode.SUCCESS_STATUS.getStatusCode()
+          || statusCode == 
TSStatusCode.DATABASE_ALREADY_EXISTS.getStatusCode()) {
+        autoCreatedTreeDatabases.add(database);
+        return TreeDatabaseCreationResult.CREATED_OR_EXISTED;
+      }
+      if (statusCode == TSStatusCode.DATABASE_CONFLICT.getStatusCode()) {
+        conflictedTreeDatabases.add(database);
+        return TreeDatabaseCreationResult.CONFLICTED;
+      }
+      throw new PipeException(
+          String.format(
+              "Auto create tree database failed: %s, status: %s",
+              database, result.getStatus() == null ? result.getStatusCode() : 
result.getStatus()));
+    } catch (final IllegalPathException e) {
+      throw new PipeException(String.format("Illegal tree database %s.", 
database), e);
+    } catch (final ExecutionException | InterruptedException e) {
+      if (e instanceof InterruptedException) {
+        Thread.currentThread().interrupt();
+      }
+      final Throwable rootCause = getRootCause(e);
+      final int errorCode;
+      if (rootCause instanceof IoTDBException) {
+        errorCode = ((IoTDBException) rootCause).getErrorCode();
+      } else if (rootCause instanceof IoTDBRuntimeException) {
+        errorCode = ((IoTDBRuntimeException) rootCause).getErrorCode();
+      } else {
+        errorCode = TSStatusCode.INTERNAL_SERVER_ERROR.getStatusCode();
+      }
+      if (errorCode == TSStatusCode.DATABASE_ALREADY_EXISTS.getStatusCode()) {
+        autoCreatedTreeDatabases.add(database);
+        return TreeDatabaseCreationResult.CREATED_OR_EXISTED;
+      }
+      if (errorCode == TSStatusCode.DATABASE_CONFLICT.getStatusCode()) {
+        conflictedTreeDatabases.add(database);
+        return TreeDatabaseCreationResult.CONFLICTED;
+      }
+      throw new PipeException("Auto create tree database failed because " + 
e.getMessage(), e);
+    }
+  }
+
+  private String getTreeDatabaseName(final Statement statement) {
+    if (statement instanceof LoadTsFileStatement) {
+      return ((LoadTsFileStatement) statement).getDatabase();
+    }
+    return null;
+  }
+
+  static void clearTreeDatabaseName(final Statement statement) {
+    if (statement instanceof LoadTsFileStatement) {
+      final LoadTsFileStatement loadTsFileStatement = (LoadTsFileStatement) 
statement;
+      loadTsFileStatement.setDatabase(null);
+      loadTsFileStatement.setDatabaseLevel(
+          IoTDBDescriptor.getInstance().getConfig().getDefaultDatabaseLevel());
+    }
+  }
+
   @Override
   protected TSStatus login() {
     final IClientSession session = SESSION_MANAGER.getCurrSession();
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/TreeModelPlanner.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/TreeModelPlanner.java
index 62aa7385707..9b4ebb4924b 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/TreeModelPlanner.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/TreeModelPlanner.java
@@ -130,6 +130,10 @@ public class TreeModelPlanner implements IPlanner {
             && ((PipeEnrichedStatement) statement).getInnerStatement()
                 instanceof LoadTsFileStatement;
     if (statement instanceof LoadTsFileStatement || isPipeEnrichedTsFileLoad) {
+      final LoadTsFileStatement loadTsFileStatement =
+          statement instanceof LoadTsFileStatement
+              ? (LoadTsFileStatement) statement
+              : (LoadTsFileStatement) ((PipeEnrichedStatement) 
statement).getInnerStatement();
       scheduler =
           new LoadTsFileScheduler(
               distributedPlan,
@@ -137,7 +141,8 @@ public class TreeModelPlanner implements IPlanner {
               stateMachine,
               syncInternalServiceClientManager,
               partitionFetcher,
-              isPipeEnrichedTsFileLoad);
+              isPipeEnrichedTsFileLoad,
+              loadTsFileStatement.getDatabase());
     } else {
       scheduler =
           new ClusterScheduler(
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 58159003a74..9c6946be6ab 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
@@ -82,6 +82,7 @@ import org.slf4j.LoggerFactory;
 
 import java.io.DataOutputStream;
 import java.io.File;
+import java.io.FileNotFoundException;
 import java.io.IOException;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
@@ -136,6 +137,7 @@ public class LoadTsFileScheduler implements IScheduler {
   private final PlanFragmentId fragmentId;
   private final Set<TRegionReplicaSet> allReplicaSets;
   private final boolean isGeneratedByPipe;
+  private final String treeDatabaseForRetry;
   private final Map<TTimePartitionSlot, ProgressIndex> 
timePartitionSlotToProgressIndex;
   private final LoadTsFileDataCacheMemoryBlock block;
 
@@ -145,7 +147,8 @@ public class LoadTsFileScheduler implements IScheduler {
       QueryStateMachine stateMachine,
       IClientManager<TEndPoint, SyncDataNodeInternalServiceClient> 
internalServiceClientManager,
       IPartitionFetcher partitionFetcher,
-      boolean isGeneratedByPipe) {
+      boolean isGeneratedByPipe,
+      String treeDatabaseForRetry) {
     this.queryContext = queryContext;
     this.stateMachine = stateMachine;
     this.tsFileNodeList = new ArrayList<>();
@@ -155,6 +158,7 @@ public class LoadTsFileScheduler implements IScheduler {
     this.partitionFetcher = new DataPartitionBatchFetcher(partitionFetcher);
     this.allReplicaSets = new HashSet<>();
     this.isGeneratedByPipe = isGeneratedByPipe;
+    this.treeDatabaseForRetry = treeDatabaseForRetry;
     this.timePartitionSlotToProgressIndex = new HashMap<>();
     this.block = 
LoadTsFileMemoryManager.getInstance().allocateDataCacheMemoryBlock();
 
@@ -566,9 +570,7 @@ public class LoadTsFileScheduler implements IScheduler {
         final TSStatus status =
             loadTsFileDataTypeConverter
                 .convertForTreeModel(
-                    LoadTsFileStatement.createUnchecked(filePath)
-                        .setDeleteAfterLoad(failedNode.isDeleteAfterLoad())
-                        .setConvertOnTypeMismatch(true))
+                    buildRetryTreeLoadStatement(filePath, 
failedNode.isDeleteAfterLoad()))
                 .orElse(null);
 
         if (loadTsFileDataTypeConverter.isSuccessful(status)) {
@@ -607,6 +609,22 @@ public class LoadTsFileScheduler implements IScheduler {
     }
   }
 
+  private LoadTsFileStatement buildRetryTreeLoadStatement(
+      final String filePath, final boolean deleteAfterLoad) throws 
FileNotFoundException {
+    final LoadTsFileStatement statement =
+        LoadTsFileStatement.createUnchecked(filePath)
+            .setDeleteAfterLoad(deleteAfterLoad)
+            .setConvertOnTypeMismatch(true);
+    if (treeDatabaseForRetry != null) {
+      statement.setDatabase(treeDatabaseForRetry);
+      statement.updateDatabaseLevelByTreeDatabase();
+    }
+    if (isGeneratedByPipe) {
+      statement.markIsGeneratedByPipe();
+    }
+    return statement;
+  }
+
   @Override
   public void stop(Throwable t) {
     dispatcher.abort();
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiverTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiverTest.java
index 82dea37f52e..d616c77eabb 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiverTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiverTest.java
@@ -232,4 +232,23 @@ public class IoTDBDataNodeReceiverTest {
       Files.deleteIfExists(tsFile);
     }
   }
+
+  @Test
+  public void testClearTreeDatabaseNameForLoadTsFileStatement() throws 
Exception {
+    final Path tsFile = Files.createTempFile("pipe-load-clear-tree-database", 
".tsfile");
+    try {
+      final LoadTsFileStatement statement =
+          IoTDBDataNodeReceiver.buildLoadTsFileStatementForSync(
+              "root.test.sg_0", tsFile.toString(), true, true);
+
+      IoTDBDataNodeReceiver.clearTreeDatabaseName(statement);
+
+      Assert.assertNull(statement.getDatabase());
+      Assert.assertEquals(
+          IoTDBDescriptor.getInstance().getConfig().getDefaultDatabaseLevel(),
+          statement.getDatabaseLevel());
+    } finally {
+      Files.deleteIfExists(tsFile);
+    }
+  }
 }
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileSchedulerTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileSchedulerTest.java
index d1151aef7da..f378bce7693 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileSchedulerTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileSchedulerTest.java
@@ -28,6 +28,7 @@ import 
org.apache.iotdb.db.queryengine.plan.planner.plan.DistributedQueryPlan;
 import org.apache.iotdb.db.queryengine.plan.planner.plan.PlanFragment;
 import org.apache.iotdb.db.queryengine.plan.planner.plan.SubPlan;
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.load.LoadSingleTsFileNode;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement;
 import 
org.apache.iotdb.db.storageengine.load.memory.LoadTsFileDataCacheMemoryBlock;
 
 import org.junit.Assert;
@@ -36,9 +37,11 @@ import org.junit.Test;
 import org.mockito.Mock;
 import org.mockito.MockitoAnnotations;
 
+import java.io.File;
 import java.lang.reflect.Constructor;
 import java.lang.reflect.Field;
 import java.lang.reflect.Method;
+import java.util.Collections;
 
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.spy;
@@ -56,6 +59,7 @@ public class LoadTsFileSchedulerTest {
     when(distributedQueryPlan.getRootSubPlan()).thenReturn(subPlan);
     when(subPlan.getPlanFragment()).thenReturn(planFragment);
     when(planFragment.getId()).thenReturn(new PlanFragmentId("test", 0));
+    
when(distributedQueryPlan.getInstances()).thenReturn(Collections.emptyList());
   }
 
   @Test
@@ -68,12 +72,40 @@ public class LoadTsFileSchedulerTest {
                 mock(QueryStateMachine.class),
                 mock(IClientManager.class),
                 mock(IPartitionFetcher.class),
-                false));
+                false,
+                null));
     t.start();
     Assert.assertNull(t.getTotalCpuTime());
     Assert.assertNull(t.getFragmentInfo());
   }
 
+  @Test
+  public void testBuildRetryTreeLoadStatementUpdatesDatabaseLevel() throws 
Exception {
+    final LoadTsFileScheduler scheduler =
+        new LoadTsFileScheduler(
+            distributedQueryPlan,
+            mock(MPPQueryContext.class),
+            mock(QueryStateMachine.class),
+            mock(IClientManager.class),
+            mock(IPartitionFetcher.class),
+            true,
+            "root.test.sg_0");
+    final Method method =
+        LoadTsFileScheduler.class.getDeclaredMethod(
+            "buildRetryTreeLoadStatement", String.class, boolean.class);
+    method.setAccessible(true);
+
+    final File tsFile = File.createTempFile("test", ".tsfile");
+    tsFile.deleteOnExit();
+
+    final LoadTsFileStatement statement =
+        (LoadTsFileStatement) method.invoke(scheduler, 
tsFile.getAbsolutePath(), true);
+
+    Assert.assertEquals("root.test.sg_0", statement.getDatabase());
+    Assert.assertEquals(2, statement.getDatabaseLevel());
+    Assert.assertTrue(statement.isGeneratedByPipe());
+  }
+
   @Test
   public void testTsFileDataManagerClearReleasesCachedMemory() throws 
Exception {
     final Constructor<LoadTsFileDataCacheMemoryBlock> memoryBlockConstructor =

Reply via email to