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 =