From 56446792bd2ca24b2a40511e49663e1a082d4649 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Thu, 17 Sep 2026 11:40:05 +0800 Subject: [PATCH 1/3] prefer local DataNode for LOAD-created DataRegions --- .../read/partition/GetDataPartitionPlan.java | 14 +++- .../GetOrCreateDataPartitionPlan.java | 14 +++- .../confignode/manager/load/LoadManager.java | 9 +++ .../manager/load/balancer/RegionBalancer.java | 64 +++++++++++++++++++ .../manager/partition/PartitionManager.java | 42 ++++++++---- .../org/apache/iotdb/db/conf/IoTDBConfig.java | 11 ++++ .../apache/iotdb/db/conf/IoTDBDescriptor.java | 10 +++ .../plan/analyze/ClusterPartitionFetcher.java | 11 ++++ .../plan/analyze/IPartitionFetcher.java | 7 ++ .../plan/analyze/load/LoadTsFileAnalyzer.java | 9 ++- .../load/LoadTsFileTableSchemaCache.java | 12 +++- .../TreeSchemaAutoCreatorAndVerifier.java | 4 ++ .../config/metadata/DatabaseSchemaTask.java | 3 + .../scheduler/load/LoadTsFileScheduler.java | 6 +- .../metadata/DatabaseSchemaStatement.java | 11 ++++ .../analyze/load/LoadTsFileAnalyzerTest.java | 2 +- .../metadata/DatabaseSchemaTaskTest.java | 13 ++++ .../src/main/thrift/confignode.thrift | 2 + 18 files changed, 224 insertions(+), 20 deletions(-) diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/partition/GetDataPartitionPlan.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/partition/GetDataPartitionPlan.java index cce8e14611490..a043ce5bdf165 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/partition/GetDataPartitionPlan.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/partition/GetDataPartitionPlan.java @@ -34,6 +34,7 @@ public class GetDataPartitionPlan extends ConfigPhysicalReadPlan { // Map>> protected Map> partitionSlotsMap; + protected int preferredDataNodeId = -1; public GetDataPartitionPlan(final ConfigPhysicalPlanType configPhysicalPlanType) { super(configPhysicalPlanType); @@ -49,6 +50,14 @@ public Map> getPartitionSlotsMa return partitionSlotsMap; } + public int getPreferredDataNodeId() { + return preferredDataNodeId; + } + + public void setPreferredDataNodeId(final int preferredDataNodeId) { + this.preferredDataNodeId = preferredDataNodeId; + } + /** * Convert TDataPartitionReq to GetDataPartitionPlan. * @@ -68,11 +77,12 @@ public boolean equals(final Object o) { return false; } final GetDataPartitionPlan that = (GetDataPartitionPlan) o; - return partitionSlotsMap.equals(that.partitionSlotsMap); + return partitionSlotsMap.equals(that.partitionSlotsMap) + && preferredDataNodeId == that.preferredDataNodeId; } @Override public int hashCode() { - return Objects.hash(partitionSlotsMap); + return Objects.hash(partitionSlotsMap, preferredDataNodeId); } } diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/partition/GetOrCreateDataPartitionPlan.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/partition/GetOrCreateDataPartitionPlan.java index 800b7e0e19a49..d2ad5fe6c9da4 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/partition/GetOrCreateDataPartitionPlan.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/read/partition/GetOrCreateDataPartitionPlan.java @@ -35,6 +35,13 @@ public GetOrCreateDataPartitionPlan( this.partitionSlotsMap = partitionSlotsMap; } + public GetOrCreateDataPartitionPlan( + final Map> partitionSlotsMap, + final int preferredDataNodeId) { + this(partitionSlotsMap); + this.preferredDataNodeId = preferredDataNodeId; + } + /** * Convert TDataPartitionReq to GetOrCreateDataPartitionPlan. * @@ -43,6 +50,11 @@ public GetOrCreateDataPartitionPlan( */ public static GetOrCreateDataPartitionPlan convertFromRpcTDataPartitionReq( final TDataPartitionReq req) { - return new GetOrCreateDataPartitionPlan(new ConcurrentHashMap<>(req.getPartitionSlotsMap())); + final GetOrCreateDataPartitionPlan plan = + new GetOrCreateDataPartitionPlan(new ConcurrentHashMap<>(req.getPartitionSlotsMap())); + if (req.isSetPreferredDataNodeId()) { + plan.setPreferredDataNodeId(req.getPreferredDataNodeId()); + } + return plan; } } diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/LoadManager.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/LoadManager.java index 5a424ea5a5e73..99d6d147434aa 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/LoadManager.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/LoadManager.java @@ -126,6 +126,15 @@ public CreateRegionGroupsPlan allocateRegionGroups( return regionBalancer.genRegionGroupsAllocationPlan(allotmentMap, consensusGroupType); } + public CreateRegionGroupsPlan allocateRegionGroups( + final Map allotmentMap, + final TConsensusGroupType consensusGroupType, + final Map preferredDataNodeMap) + throws NotEnoughDataNodeException, DatabaseNotExistsException { + return regionBalancer.genRegionGroupsAllocationPlan( + allotmentMap, consensusGroupType, preferredDataNodeMap); + } + /** * Allocate SchemaPartitions. * diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/RegionBalancer.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/RegionBalancer.java index 73583151f9819..04f90dfc49989 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/RegionBalancer.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/RegionBalancer.java @@ -22,6 +22,7 @@ import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId; import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType; import org.apache.iotdb.common.rpc.thrift.TDataNodeConfiguration; +import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation; import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet; import org.apache.iotdb.commons.cluster.NodeStatus; import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor; @@ -38,7 +39,10 @@ import org.apache.iotdb.confignode.manager.node.NodeManager; import org.apache.iotdb.confignode.manager.partition.PartitionManager; import org.apache.iotdb.confignode.manager.schema.ClusterSchemaManager; +import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema; +import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -82,6 +86,14 @@ public RegionBalancer(IManager configManager) { public CreateRegionGroupsPlan genRegionGroupsAllocationPlan( final Map allotmentMap, final TConsensusGroupType consensusGroupType) throws NotEnoughDataNodeException, DatabaseNotExistsException { + return genRegionGroupsAllocationPlan(allotmentMap, consensusGroupType, null); + } + + public CreateRegionGroupsPlan genRegionGroupsAllocationPlan( + final Map allotmentMap, + final TConsensusGroupType consensusGroupType, + final Map preferredDataNodeMap) + throws NotEnoughDataNodeException, DatabaseNotExistsException { // Some new RegionGroups will have to occupy unknown DataNodes if the number of online // DataNodes is insufficient (Unknown DataNodes are intentionally kept as candidates). @@ -118,6 +130,14 @@ public CreateRegionGroupsPlan genRegionGroupsAllocationPlan( final int allotment = entry.getValue(); final int replicationFactor = getClusterSchemaManager().getReplicationFactor(database, consensusGroupType); + final Integer requestPreferredDataNode = + preferredDataNodeMap == null ? null : preferredDataNodeMap.get(database); + final int preferredDataNodeId = + TConsensusGroupType.DataRegion.equals(consensusGroupType) + ? (requestPreferredDataNode != null + ? requestPreferredDataNode + : getPreferredDataNodeId(database)) + : -1; // Only considering the specified Database when doing allocation final List databaseAllocatedRegionGroups = getPartitionManager().getAllReplicaSets(database, consensusGroupType); @@ -144,6 +164,7 @@ public CreateRegionGroupsPlan genRegionGroupsAllocationPlan( replicationFactor, new TConsensusGroupId( consensusGroupType, getPartitionManager().generateNextRegionGroupId())); + preferDataNodeIfPossible(newRegionGroup, preferredDataNodeId, availableDataNodeMap); createRegionGroupsPlan.addRegionGroup(database, newRegionGroup); // Mark the new RegionGroup as allocated @@ -175,6 +196,49 @@ private ProcedureManager getProcedureManager() { return configManager.getProcedureManager(); } + private int getPreferredDataNodeId(final String database) { + try { + final TDatabaseSchema databaseSchema = + getClusterSchemaManager().getDatabaseSchemaByName(database); + return databaseSchema.isSetPreferredDataNodeId() + ? databaseSchema.getPreferredDataNodeId() + : -1; + } catch (final DatabaseNotExistsException e) { + return -1; + } + } + + private static void preferDataNodeIfPossible( + final TRegionReplicaSet regionReplicaSet, + final int preferredDataNodeId, + final Map availableDataNodeMap) { + if (preferredDataNodeId < 0 || !availableDataNodeMap.containsKey(preferredDataNodeId)) { + return; + } + + final List locations = + new ArrayList<>(regionReplicaSet.getDataNodeLocations()); + if (locations.isEmpty()) { + return; + } + + int preferredIndex = -1; + for (int i = 0; i < locations.size(); i++) { + if (locations.get(i).getDataNodeId() == preferredDataNodeId) { + preferredIndex = i; + break; + } + } + + if (preferredIndex > 0) { + Collections.swap(locations, 0, preferredIndex); + } else if (preferredIndex < 0) { + locations.set( + locations.size() - 1, availableDataNodeMap.get(preferredDataNodeId).getLocation()); + } + regionReplicaSet.setDataNodeLocations(locations); + } + public enum RegionGroupAllocatePolicy { GREEDY, GCR, diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java index 28655923cca92..d50f530498373 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java @@ -446,7 +446,9 @@ public DataPartitionResp getOrCreateDataPartition(final GetOrCreateDataPartition storageGroup, unassignedDataPartitionSlots.size())); TSStatus status = extendRegionGroupIfNecessary( - unassignedDataPartitionSlotsCountMap, TConsensusGroupType.DataRegion); + unassignedDataPartitionSlotsCountMap, + TConsensusGroupType.DataRegion, + req.getPreferredDataNodeId()); if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) { // Return an error code if Region extension failed resp.setStatus(status); @@ -607,6 +609,13 @@ private TSStatus consensusWritePartitionResult(ConfigPhysicalPlan plan) { private TSStatus extendRegionGroupIfNecessary( final Map unassignedPartitionSlotsCountMap, final TConsensusGroupType consensusGroupType) { + return extendRegionGroupIfNecessary(unassignedPartitionSlotsCountMap, consensusGroupType, -1); + } + + private TSStatus extendRegionGroupIfNecessary( + final Map unassignedPartitionSlotsCountMap, + final TConsensusGroupType consensusGroupType, + final int preferredDataNodeId) { final TSStatus result = new TSStatus(); @@ -615,21 +624,21 @@ private TSStatus extendRegionGroupIfNecessary( switch (CONF.getSchemaRegionGroupExtensionPolicy()) { case CUSTOM: return customExtendRegionGroupIfNecessary( - unassignedPartitionSlotsCountMap, consensusGroupType); + unassignedPartitionSlotsCountMap, consensusGroupType, preferredDataNodeId); case AUTO: default: return autoExtendRegionGroupIfNecessary( - unassignedPartitionSlotsCountMap, consensusGroupType); + unassignedPartitionSlotsCountMap, consensusGroupType, preferredDataNodeId); } } else { switch (CONF.getDataRegionGroupExtensionPolicy()) { case CUSTOM: return customExtendRegionGroupIfNecessary( - unassignedPartitionSlotsCountMap, consensusGroupType); + unassignedPartitionSlotsCountMap, consensusGroupType, preferredDataNodeId); case AUTO: default: return autoExtendRegionGroupIfNecessary( - unassignedPartitionSlotsCountMap, consensusGroupType); + unassignedPartitionSlotsCountMap, consensusGroupType, preferredDataNodeId); } } } catch (NotEnoughDataNodeException e) { @@ -647,7 +656,8 @@ private TSStatus extendRegionGroupIfNecessary( private TSStatus customExtendRegionGroupIfNecessary( final Map unassignedPartitionSlotsCountMap, - final TConsensusGroupType consensusGroupType) + final TConsensusGroupType consensusGroupType, + final int preferredDataNodeId) throws DatabaseNotExistsException, NotEnoughDataNodeException { // Map @@ -666,12 +676,13 @@ private TSStatus customExtendRegionGroupIfNecessary( } } - return generateAndAllocateRegionGroups(allotmentMap, consensusGroupType); + return generateAndAllocateRegionGroups(allotmentMap, consensusGroupType, preferredDataNodeId); } private TSStatus autoExtendRegionGroupIfNecessary( final Map unassignedPartitionSlotsCountMap, - final TConsensusGroupType consensusGroupType) + final TConsensusGroupType consensusGroupType, + final int preferredDataNodeId) throws NotEnoughDataNodeException, DatabaseNotExistsException { // Map @@ -735,15 +746,24 @@ private TSStatus autoExtendRegionGroupIfNecessary( } } - return generateAndAllocateRegionGroups(allotmentMap, consensusGroupType); + return generateAndAllocateRegionGroups(allotmentMap, consensusGroupType, preferredDataNodeId); } private TSStatus generateAndAllocateRegionGroups( - final Map allotmentMap, final TConsensusGroupType consensusGroupType) + final Map allotmentMap, + final TConsensusGroupType consensusGroupType, + final int preferredDataNodeId) throws NotEnoughDataNodeException, DatabaseNotExistsException { if (!allotmentMap.isEmpty()) { + final Map preferredDataNodeMap = new ConcurrentHashMap<>(); + if (preferredDataNodeId >= 0) { + allotmentMap + .keySet() + .forEach(database -> preferredDataNodeMap.put(database, preferredDataNodeId)); + } final CreateRegionGroupsPlan createRegionGroupsPlan = - getLoadManager().allocateRegionGroups(allotmentMap, consensusGroupType); + getLoadManager() + .allocateRegionGroups(allotmentMap, consensusGroupType, preferredDataNodeMap); LOGGER.info(ManagerMessages.CREATEREGIONGROUPS_STARTING_TO_CREATE_THE_FOLLOWING_REGIONGROUPS); createRegionGroupsPlan.planLog(LOGGER); return getProcedureManager().createRegionGroups(consensusGroupType, createRegionGroupsPlan); 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 80f9f1f1b30f8..b448382ac13f4 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 @@ -1156,6 +1156,9 @@ public class IoTDBConfig { /** Load related */ private double maxAllocateMemoryRatioForLoad = 0.8; + /** Hidden hot-reload switch that prefers placing LOAD-created DataRegions on this DataNode. */ + private volatile boolean loadTsFilePreferLocalNode = true; + private int loadTsFileAnalyzeSchemaBatchReadTimeSeriesMetadataCount = 4096; private int loadTsFileAnalyzeSchemaBatchFlushTimeSeriesNumber = 4096; private int loadTsFileAnalyzeSchemaBatchFlushTableDeviceNumber = 4096; // For table model @@ -4289,6 +4292,14 @@ public void setLoadTsFileRetryCountOnRegionChange(int loadTsFileRetryCountOnRegi this.loadTsFileRetryCountOnRegionChange = loadTsFileRetryCountOnRegionChange; } + public boolean isLoadTsFilePreferLocalNode() { + return loadTsFilePreferLocalNode; + } + + public void setLoadTsFilePreferLocalNode(final boolean loadTsFilePreferLocalNode) { + this.loadTsFilePreferLocalNode = loadTsFilePreferLocalNode; + } + public double getLoadWriteThroughputBytesPerSecond() { return loadWriteThroughputBytesPerSecond; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java index 565f0b7832a69..44e14899baf75 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java @@ -2642,6 +2642,11 @@ private void loadLoadTsFileProps(TrimProperties properties) { properties.getProperty( "load_tsfile_source_path_check_enable", Boolean.toString(conf.isLoadTsFileSourcePathCheckEnabled())))); + conf.setLoadTsFilePreferLocalNode( + Boolean.parseBoolean( + properties.getProperty( + "load_tsfile_prefer_local_node", + Boolean.toString(conf.isLoadTsFilePreferLocalNode())))); conf.setLoadTabletConversionThresholdBytes( Long.parseLong( @@ -2791,6 +2796,11 @@ private void loadLoadTsFileHotModifiedProp(TrimProperties properties) throws IOE properties.getProperty( "load_tsfile_source_path_check_enable", Boolean.toString(conf.isLoadTsFileSourcePathCheckEnabled())))); + conf.setLoadTsFilePreferLocalNode( + Boolean.parseBoolean( + properties.getProperty( + "load_tsfile_prefer_local_node", + Boolean.toString(conf.isLoadTsFilePreferLocalNode())))); conf.setLoadActiveListeningEnable( Boolean.parseBoolean( diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/ClusterPartitionFetcher.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/ClusterPartitionFetcher.java index 91c814dbaa475..f5c5f5aeb0697 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/ClusterPartitionFetcher.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/ClusterPartitionFetcher.java @@ -319,6 +319,14 @@ public DataPartition getOrCreateDataPartition( @Override public DataPartition getOrCreateDataPartition( final List dataPartitionQueryParams, final String userName) { + return getOrCreateDataPartition(dataPartitionQueryParams, userName, -1); + } + + @Override + public DataPartition getOrCreateDataPartition( + final List dataPartitionQueryParams, + final String userName, + final int preferredDataNodeId) { final Map> splitDataPartitionQueryParams = splitDataPartitionQueryParam( dataPartitionQueryParams, config.isAutoCreateSchemaEnabled(), userName); @@ -330,6 +338,9 @@ public DataPartition getOrCreateDataPartition( try (final ConfigNodeClient client = configNodeClientManager.borrowClient(ConfigNodeInfo.CONFIG_REGION_ID)) { final TDataPartitionReq req = constructDataPartitionReq(splitDataPartitionQueryParams); + if (preferredDataNodeId >= 0) { + req.setPreferredDataNodeId(preferredDataNodeId); + } final TDataPartitionTableResp dataPartitionTableResp = client.getOrCreateDataPartitionTable(req); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/IPartitionFetcher.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/IPartitionFetcher.java index 6d6f59133af9c..9347e71097e60 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/IPartitionFetcher.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/IPartitionFetcher.java @@ -100,6 +100,13 @@ DataPartition getOrCreateDataPartition( DataPartition getOrCreateDataPartition( final List dataPartitionQueryParams, final String userName); + default DataPartition getOrCreateDataPartition( + final List dataPartitionQueryParams, + final String userName, + final int preferredDataNodeId) { + return getOrCreateDataPartition(dataPartitionQueryParams, userName); + } + /** Get schema partition and matched nodes according to path pattern tree. */ default SchemaNodeManagementPartition getSchemaNodeManagementPartition( PathPatternTree patternTree, PathPatternTree scope, boolean needAuditDB) { 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 52ab6f5ebd204..b660d8439c44c 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 @@ -27,6 +27,7 @@ import org.apache.iotdb.commons.queryengine.common.SqlDialect; import org.apache.iotdb.commons.queryengine.utils.TimestampPrecisionUtils; import org.apache.iotdb.commons.utils.RetryUtils; +import org.apache.iotdb.db.conf.IoTDBConfig; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.exception.load.LoadAnalyzeException; import org.apache.iotdb.db.exception.load.LoadAnalyzeInvalidPathException; @@ -214,6 +215,11 @@ protected boolean isConvertOnTypeMismatch() { return isConvertOnTypeMismatch; } + protected int getPreferredDataNodeId() { + final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); + return config.isLoadTsFilePreferLocalNode() ? config.getDataNodeId() : -1; + } + public IAnalysis analyzeFileByFile(IAnalysis analysis) { if (!checkBeforeAnalyzeFileByFile(analysis)) { return analysis; @@ -708,7 +714,8 @@ private TreeSchemaAutoCreatorAndVerifier getOrCreateTreeSchemaVerifier() { private LoadTsFileTableSchemaCache getOrCreateTableSchemaCache() { if (tableSchemaCache == null) { tableSchemaCache = - new LoadTsFileTableSchemaCache(metadata, context, isAutoCreateDatabase, isVerifySchema); + new LoadTsFileTableSchemaCache( + metadata, context, isAutoCreateDatabase, isVerifySchema, getPreferredDataNodeId()); } return tableSchemaCache; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileTableSchemaCache.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileTableSchemaCache.java index 322a86acfd3b3..f538b309ccd92 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileTableSchemaCache.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileTableSchemaCache.java @@ -100,6 +100,7 @@ public class LoadTsFileTableSchemaCache { private final Metadata metadata; private final MPPQueryContext context; private final boolean shouldVerifyDataType; + private final int preferredDataNodeId; private Map> currentBatchTable2Devices; @@ -121,7 +122,8 @@ public LoadTsFileTableSchemaCache( final Metadata metadata, final MPPQueryContext context, final boolean needToCreateDatabase, - final boolean shouldVerifyDataType) + final boolean shouldVerifyDataType, + final int preferredDataNodeId) throws LoadRuntimeOutOfMemoryException { this.block = LoadTsFileMemoryManager.getInstance() @@ -129,6 +131,7 @@ public LoadTsFileTableSchemaCache( this.metadata = metadata; this.context = context; this.shouldVerifyDataType = shouldVerifyDataType; + this.preferredDataNodeId = preferredDataNodeId; this.currentBatchTable2Devices = new HashMap<>(); this.currentModifications = PatternTreeMapFactory.getModsPatternTreeMap(); this.needToCreateDatabase = needToCreateDatabase; @@ -349,8 +352,11 @@ private void autoCreateTableDatabaseIfAbsent(final String database) throws LoadA AuthorityChecker.getAccessControl() .checkCanCreateDatabase(context.getSession().getUserName(), database, context); - final CreateDBTask task = - new CreateDBTask(new TDatabaseSchema(database).setIsTableModel(true), true); + final TDatabaseSchema schema = new TDatabaseSchema(database).setIsTableModel(true); + if (preferredDataNodeId >= 0) { + schema.setPreferredDataNodeId(preferredDataNodeId); + } + final CreateDBTask task = new CreateDBTask(schema, true); try { final ListenableFuture future = task.execute(ClusterConfigTaskExecutor.getInstance()); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/TreeSchemaAutoCreatorAndVerifier.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/TreeSchemaAutoCreatorAndVerifier.java index 3791109240592..ed8fca06dead1 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/TreeSchemaAutoCreatorAndVerifier.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/TreeSchemaAutoCreatorAndVerifier.java @@ -407,6 +407,10 @@ private void autoCreateDatabase() final DatabaseSchemaStatement statement = new DatabaseSchemaStatement(DatabaseSchemaStatement.DatabaseSchemaStatementType.CREATE); statement.setDatabasePath(databasePath); + final int preferredDataNodeId = loadTsFileAnalyzer.getPreferredDataNodeId(); + if (preferredDataNodeId >= 0) { + statement.setPreferredDataNodeId(preferredDataNodeId); + } // do not print exception log because it is not an error statement.setEnablePrintExceptionLog(false); executeSetDatabaseStatement(statement); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTask.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTask.java index 581b1ed2b9980..5d9dd31b08bf3 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTask.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTask.java @@ -76,6 +76,9 @@ public static TDatabaseSchema constructDatabaseSchema( if (databaseSchemaStatement.isSetNeedLastCache()) { databaseSchema.setNeedLastCache(databaseSchemaStatement.isNeedLastCache()); } + if (databaseSchemaStatement.getPreferredDataNodeId() != null) { + databaseSchema.setPreferredDataNodeId(databaseSchemaStatement.getPreferredDataNodeId()); + } databaseSchema.setIsTableModel(false); return databaseSchema; } 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 a116f3ac1d84b..f2a0c528a093f 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 @@ -947,9 +947,12 @@ private void clear() { private static class DataPartitionBatchFetcher { private final IPartitionFetcher fetcher; private String database; + private final int preferredDataNodeId; public DataPartitionBatchFetcher(IPartitionFetcher fetcher) { this.fetcher = fetcher; + final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); + this.preferredDataNodeId = config.isLoadTsFilePreferLocalNode() ? config.getDataNodeId() : -1; } public void setDatabase(String database) { @@ -965,7 +968,8 @@ public List queryDataPartition( List> subSlotList = slotList.subList(i, Math.min(size, i + TRANSMIT_LIMIT)); DataPartition dataPartition = - fetcher.getOrCreateDataPartition(toQueryParam(subSlotList), userName); + fetcher.getOrCreateDataPartition( + toQueryParam(subSlotList), userName, preferredDataNodeId); for (final Pair pair : subSlotList) { // database is an explicit database hint for table-model loads and // pipe-generated tree-model loads. diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/DatabaseSchemaStatement.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/DatabaseSchemaStatement.java index dc18d56055649..76da63a975938 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/DatabaseSchemaStatement.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/DatabaseSchemaStatement.java @@ -41,6 +41,7 @@ public class DatabaseSchemaStatement extends Statement implements IConfigStateme private boolean enablePrintExceptionLog = true; private boolean needLastCache = true; private boolean isNeedLastCacheSet = false; + private Integer preferredDataNodeId = null; // Deprecated private Integer schemaReplicationFactor = null; @@ -133,6 +134,14 @@ public void setNeedLastCache(final boolean needLastCache) { this.isNeedLastCacheSet = true; } + public Integer getPreferredDataNodeId() { + return preferredDataNodeId; + } + + public void setPreferredDataNodeId(final Integer preferredDataNodeId) { + this.preferredDataNodeId = preferredDataNodeId; + } + @Override public R accept(final StatementVisitor visitor, final C context) { switch (subType) { @@ -171,6 +180,8 @@ public String toString() { + maxSchemaRegionGroupNum + ", maxDataRegionGroupNum=" + maxDataRegionGroupNum + + ", preferredDataNodeId=" + + preferredDataNodeId + '}'; } diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java index 1e9aae0842826..abc4fddac7e6d 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java @@ -447,7 +447,7 @@ private LoadTsFileTreeSchemaCache getTreeSchemaCache( private LoadTsFileTableSchemaCache createTableSchemaCache(final boolean shouldVerifyDataType) throws LoadRuntimeOutOfMemoryException { return new LoadTsFileTableSchemaCache( - null, new MPPQueryContext(new QueryId("load_test")), false, shouldVerifyDataType); + null, new MPPQueryContext(new QueryId("load_test")), false, shouldVerifyDataType, -1); } private Method getVerifyTableDataTypeMethod() throws NoSuchMethodException { diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTaskTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTaskTest.java index 41703a4583493..4ed9f534bb2bb 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTaskTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTaskTest.java @@ -54,4 +54,17 @@ public void testConstructDatabaseSchemaSetsNeedLastCacheWhenPresent() throws Exc assertTrue(databaseSchema.isSetNeedLastCache()); assertEquals(false, databaseSchema.isNeedLastCache()); } + + @Test + public void testConstructDatabaseSchemaSetsPreferredDataNodeWhenPresent() throws Exception { + final DatabaseSchemaStatement statement = + new DatabaseSchemaStatement(DatabaseSchemaStatement.DatabaseSchemaStatementType.ALTER); + statement.setDatabasePath(new PartialPath("root.sg")); + statement.setPreferredDataNodeId(7); + + final TDatabaseSchema databaseSchema = DatabaseSchemaTask.constructDatabaseSchema(statement); + + assertTrue(databaseSchema.isSetPreferredDataNodeId()); + assertEquals(7, databaseSchema.getPreferredDataNodeId()); + } } diff --git a/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift b/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift index 512e2e05f05fc..9ff194dfa9c83 100644 --- a/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift +++ b/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift @@ -241,6 +241,7 @@ struct TDatabaseSchema { 10: optional i64 timePartitionOrigin 11: optional bool isTableModel 12: optional bool needLastCache + 13: optional i32 preferredDataNodeId } // Schema @@ -279,6 +280,7 @@ struct TTimeSlotList { struct TDataPartitionReq { // map> 1: required map> partitionSlotsMap + 2: optional i32 preferredDataNodeId } struct TDataPartitionTableResp { From daf75e4a224905ad4166b7163952992322f127dd Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Thu, 17 Sep 2026 11:49:13 +0800 Subject: [PATCH 2/3] Fix LoadTsFileTableSchemaCache constructor in tests --- .../queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java index abc4fddac7e6d..8a2d1f786de1e 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java @@ -495,7 +495,7 @@ private static class TrackingLoadTsFileTableSchemaCache extends LoadTsFileTableS private final Set> verifiedDevices = new HashSet<>(); private TrackingLoadTsFileTableSchemaCache() throws LoadRuntimeOutOfMemoryException { - super(null, new MPPQueryContext(new QueryId("load_test")), false, true); + super(null, new MPPQueryContext(new QueryId("load_test")), false, true, -1); } @Override From 1b5eb12238bb78c8002d8f765586a06c91ba96fa Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Thu, 17 Sep 2026 14:45:03 +0800 Subject: [PATCH 3/3] Keep LOAD DataNode preference request-scoped --- .../manager/load/balancer/RegionBalancer.java | 17 +------------- .../plan/analyze/ClusterPartitionFetcher.java | 21 +++++++++++++++++- .../plan/analyze/load/LoadTsFileAnalyzer.java | 9 +------- .../load/LoadTsFileTableSchemaCache.java | 12 +++------- .../TreeSchemaAutoCreatorAndVerifier.java | 4 ---- .../config/metadata/DatabaseSchemaTask.java | 3 --- .../metadata/DatabaseSchemaStatement.java | 11 ---------- .../LoadTsFileDataTypeConverter.java | 22 ++++++++++++++++--- .../analyze/load/LoadTsFileAnalyzerTest.java | 4 ++-- .../metadata/DatabaseSchemaTaskTest.java | 13 ----------- .../src/main/thrift/confignode.thrift | 1 - 11 files changed, 46 insertions(+), 71 deletions(-) diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/RegionBalancer.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/RegionBalancer.java index 04f90dfc49989..e27aa3b1fe4e7 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/RegionBalancer.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/RegionBalancer.java @@ -39,7 +39,6 @@ import org.apache.iotdb.confignode.manager.node.NodeManager; import org.apache.iotdb.confignode.manager.partition.PartitionManager; import org.apache.iotdb.confignode.manager.schema.ClusterSchemaManager; -import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema; import java.util.ArrayList; import java.util.Collections; @@ -134,9 +133,7 @@ public CreateRegionGroupsPlan genRegionGroupsAllocationPlan( preferredDataNodeMap == null ? null : preferredDataNodeMap.get(database); final int preferredDataNodeId = TConsensusGroupType.DataRegion.equals(consensusGroupType) - ? (requestPreferredDataNode != null - ? requestPreferredDataNode - : getPreferredDataNodeId(database)) + ? (requestPreferredDataNode != null ? requestPreferredDataNode : -1) : -1; // Only considering the specified Database when doing allocation final List databaseAllocatedRegionGroups = @@ -196,18 +193,6 @@ private ProcedureManager getProcedureManager() { return configManager.getProcedureManager(); } - private int getPreferredDataNodeId(final String database) { - try { - final TDatabaseSchema databaseSchema = - getClusterSchemaManager().getDatabaseSchemaByName(database); - return databaseSchema.isSetPreferredDataNodeId() - ? databaseSchema.getPreferredDataNodeId() - : -1; - } catch (final DatabaseNotExistsException e) { - return -1; - } - } - private static void preferDataNodeIfPossible( final TRegionReplicaSet regionReplicaSet, final int preferredDataNodeId, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/ClusterPartitionFetcher.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/ClusterPartitionFetcher.java index f5c5f5aeb0697..7260ca5cab6e0 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/ClusterPartitionFetcher.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/ClusterPartitionFetcher.java @@ -75,6 +75,7 @@ public class ClusterPartitionFetcher implements IPartitionFetcher { private static final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); + private static final ThreadLocal LOAD_PREFERRED_DATA_NODE_ID = new ThreadLocal<>(); private final SeriesPartitionExecutor partitionExecutor; @@ -94,6 +95,23 @@ public static ClusterPartitionFetcher getInstance() { return ClusterPartitionFetcherHolder.INSTANCE; } + public static void setLoadPreferredDataNodeId(final int preferredDataNodeId) { + if (preferredDataNodeId >= 0) { + LOAD_PREFERRED_DATA_NODE_ID.set(preferredDataNodeId); + } else { + LOAD_PREFERRED_DATA_NODE_ID.remove(); + } + } + + public static void clearLoadPreferredDataNodeId() { + LOAD_PREFERRED_DATA_NODE_ID.remove(); + } + + private static int getLoadPreferredDataNodeId() { + final Integer preferredDataNodeId = LOAD_PREFERRED_DATA_NODE_ID.get(); + return preferredDataNodeId == null ? -1 : preferredDataNodeId; + } + private ClusterPartitionFetcher() { this.partitionExecutor = SeriesPartitionExecutor.getSeriesPartitionExecutor( @@ -319,7 +337,8 @@ public DataPartition getOrCreateDataPartition( @Override public DataPartition getOrCreateDataPartition( final List dataPartitionQueryParams, final String userName) { - return getOrCreateDataPartition(dataPartitionQueryParams, userName, -1); + return getOrCreateDataPartition( + dataPartitionQueryParams, userName, getLoadPreferredDataNodeId()); } @Override 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 b660d8439c44c..52ab6f5ebd204 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 @@ -27,7 +27,6 @@ import org.apache.iotdb.commons.queryengine.common.SqlDialect; import org.apache.iotdb.commons.queryengine.utils.TimestampPrecisionUtils; import org.apache.iotdb.commons.utils.RetryUtils; -import org.apache.iotdb.db.conf.IoTDBConfig; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.exception.load.LoadAnalyzeException; import org.apache.iotdb.db.exception.load.LoadAnalyzeInvalidPathException; @@ -215,11 +214,6 @@ protected boolean isConvertOnTypeMismatch() { return isConvertOnTypeMismatch; } - protected int getPreferredDataNodeId() { - final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); - return config.isLoadTsFilePreferLocalNode() ? config.getDataNodeId() : -1; - } - public IAnalysis analyzeFileByFile(IAnalysis analysis) { if (!checkBeforeAnalyzeFileByFile(analysis)) { return analysis; @@ -714,8 +708,7 @@ private TreeSchemaAutoCreatorAndVerifier getOrCreateTreeSchemaVerifier() { private LoadTsFileTableSchemaCache getOrCreateTableSchemaCache() { if (tableSchemaCache == null) { tableSchemaCache = - new LoadTsFileTableSchemaCache( - metadata, context, isAutoCreateDatabase, isVerifySchema, getPreferredDataNodeId()); + new LoadTsFileTableSchemaCache(metadata, context, isAutoCreateDatabase, isVerifySchema); } return tableSchemaCache; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileTableSchemaCache.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileTableSchemaCache.java index f538b309ccd92..322a86acfd3b3 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileTableSchemaCache.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileTableSchemaCache.java @@ -100,7 +100,6 @@ public class LoadTsFileTableSchemaCache { private final Metadata metadata; private final MPPQueryContext context; private final boolean shouldVerifyDataType; - private final int preferredDataNodeId; private Map> currentBatchTable2Devices; @@ -122,8 +121,7 @@ public LoadTsFileTableSchemaCache( final Metadata metadata, final MPPQueryContext context, final boolean needToCreateDatabase, - final boolean shouldVerifyDataType, - final int preferredDataNodeId) + final boolean shouldVerifyDataType) throws LoadRuntimeOutOfMemoryException { this.block = LoadTsFileMemoryManager.getInstance() @@ -131,7 +129,6 @@ public LoadTsFileTableSchemaCache( this.metadata = metadata; this.context = context; this.shouldVerifyDataType = shouldVerifyDataType; - this.preferredDataNodeId = preferredDataNodeId; this.currentBatchTable2Devices = new HashMap<>(); this.currentModifications = PatternTreeMapFactory.getModsPatternTreeMap(); this.needToCreateDatabase = needToCreateDatabase; @@ -352,11 +349,8 @@ private void autoCreateTableDatabaseIfAbsent(final String database) throws LoadA AuthorityChecker.getAccessControl() .checkCanCreateDatabase(context.getSession().getUserName(), database, context); - final TDatabaseSchema schema = new TDatabaseSchema(database).setIsTableModel(true); - if (preferredDataNodeId >= 0) { - schema.setPreferredDataNodeId(preferredDataNodeId); - } - final CreateDBTask task = new CreateDBTask(schema, true); + final CreateDBTask task = + new CreateDBTask(new TDatabaseSchema(database).setIsTableModel(true), true); try { final ListenableFuture future = task.execute(ClusterConfigTaskExecutor.getInstance()); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/TreeSchemaAutoCreatorAndVerifier.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/TreeSchemaAutoCreatorAndVerifier.java index ed8fca06dead1..3791109240592 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/TreeSchemaAutoCreatorAndVerifier.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/TreeSchemaAutoCreatorAndVerifier.java @@ -407,10 +407,6 @@ private void autoCreateDatabase() final DatabaseSchemaStatement statement = new DatabaseSchemaStatement(DatabaseSchemaStatement.DatabaseSchemaStatementType.CREATE); statement.setDatabasePath(databasePath); - final int preferredDataNodeId = loadTsFileAnalyzer.getPreferredDataNodeId(); - if (preferredDataNodeId >= 0) { - statement.setPreferredDataNodeId(preferredDataNodeId); - } // do not print exception log because it is not an error statement.setEnablePrintExceptionLog(false); executeSetDatabaseStatement(statement); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTask.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTask.java index 5d9dd31b08bf3..581b1ed2b9980 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTask.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTask.java @@ -76,9 +76,6 @@ public static TDatabaseSchema constructDatabaseSchema( if (databaseSchemaStatement.isSetNeedLastCache()) { databaseSchema.setNeedLastCache(databaseSchemaStatement.isNeedLastCache()); } - if (databaseSchemaStatement.getPreferredDataNodeId() != null) { - databaseSchema.setPreferredDataNodeId(databaseSchemaStatement.getPreferredDataNodeId()); - } databaseSchema.setIsTableModel(false); return databaseSchema; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/DatabaseSchemaStatement.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/DatabaseSchemaStatement.java index 76da63a975938..dc18d56055649 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/DatabaseSchemaStatement.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/DatabaseSchemaStatement.java @@ -41,7 +41,6 @@ public class DatabaseSchemaStatement extends Statement implements IConfigStateme private boolean enablePrintExceptionLog = true; private boolean needLastCache = true; private boolean isNeedLastCacheSet = false; - private Integer preferredDataNodeId = null; // Deprecated private Integer schemaReplicationFactor = null; @@ -134,14 +133,6 @@ public void setNeedLastCache(final boolean needLastCache) { this.isNeedLastCacheSet = true; } - public Integer getPreferredDataNodeId() { - return preferredDataNodeId; - } - - public void setPreferredDataNodeId(final Integer preferredDataNodeId) { - this.preferredDataNodeId = preferredDataNodeId; - } - @Override public R accept(final StatementVisitor visitor, final C context) { switch (subType) { @@ -180,8 +171,6 @@ public String toString() { + maxSchemaRegionGroupNum + ", maxDataRegionGroupNum=" + maxDataRegionGroupNum - + ", preferredDataNodeId=" - + preferredDataNodeId + '}'; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTsFileDataTypeConverter.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTsFileDataTypeConverter.java index 5382afc6f0388..52c4d20e886cd 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTsFileDataTypeConverter.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTsFileDataTypeConverter.java @@ -25,6 +25,7 @@ import org.apache.iotdb.commons.pipe.config.PipeConfig; import org.apache.iotdb.commons.queryengine.common.SqlDialect; import org.apache.iotdb.db.auth.AuthorityChecker; +import org.apache.iotdb.db.conf.IoTDBConfig; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.exception.load.LoadRuntimeOutOfMemoryException; import org.apache.iotdb.db.i18n.StorageEngineMessages; @@ -123,6 +124,11 @@ public static boolean isMemoryPressureException(final Throwable throwable) { return false; } + private static int getPreferredDataNodeIdForLoad() { + final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); + return config.isLoadTsFilePreferLocalNode() ? config.getDataNodeId() : -1; + } + public static TSStatus getMemoryPressureStatus(final Throwable throwable) { Throwable current = throwable; while (current != null) { @@ -176,8 +182,13 @@ public Optional convertForTableModel(final LoadTsFile loadTsFileTableS try { getTabletConversionSemaphore().acquire(); isPermitAcquired = true; - return loadTsFileTableStatement.accept( - tableStatementDataTypeConvertExecutionVisitor, loadTsFileTableStatement.getDatabase()); + ClusterPartitionFetcher.setLoadPreferredDataNodeId(getPreferredDataNodeIdForLoad()); + try { + return loadTsFileTableStatement.accept( + tableStatementDataTypeConvertExecutionVisitor, loadTsFileTableStatement.getDatabase()); + } finally { + ClusterPartitionFetcher.clearLoadPreferredDataNodeId(); + } } catch (final InterruptedException e) { return getInterruptedConversionStatus(e); } catch (Exception e) { @@ -242,7 +253,12 @@ public Optional convertForTreeModel(final LoadTsFileStatement loadTsFi try { getTabletConversionSemaphore().acquire(); isPermitAcquired = true; - return loadTsFileTreeStatement.accept(treeStatementDataTypeConvertExecutionVisitor, null); + ClusterPartitionFetcher.setLoadPreferredDataNodeId(getPreferredDataNodeIdForLoad()); + try { + return loadTsFileTreeStatement.accept(treeStatementDataTypeConvertExecutionVisitor, null); + } finally { + ClusterPartitionFetcher.clearLoadPreferredDataNodeId(); + } } catch (final InterruptedException e) { return getInterruptedConversionStatus(e); } catch (Exception e) { diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java index 8a2d1f786de1e..1e9aae0842826 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzerTest.java @@ -447,7 +447,7 @@ private LoadTsFileTreeSchemaCache getTreeSchemaCache( private LoadTsFileTableSchemaCache createTableSchemaCache(final boolean shouldVerifyDataType) throws LoadRuntimeOutOfMemoryException { return new LoadTsFileTableSchemaCache( - null, new MPPQueryContext(new QueryId("load_test")), false, shouldVerifyDataType, -1); + null, new MPPQueryContext(new QueryId("load_test")), false, shouldVerifyDataType); } private Method getVerifyTableDataTypeMethod() throws NoSuchMethodException { @@ -495,7 +495,7 @@ private static class TrackingLoadTsFileTableSchemaCache extends LoadTsFileTableS private final Set> verifiedDevices = new HashSet<>(); private TrackingLoadTsFileTableSchemaCache() throws LoadRuntimeOutOfMemoryException { - super(null, new MPPQueryContext(new QueryId("load_test")), false, true, -1); + super(null, new MPPQueryContext(new QueryId("load_test")), false, true); } @Override diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTaskTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTaskTest.java index 4ed9f534bb2bb..41703a4583493 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTaskTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTaskTest.java @@ -54,17 +54,4 @@ public void testConstructDatabaseSchemaSetsNeedLastCacheWhenPresent() throws Exc assertTrue(databaseSchema.isSetNeedLastCache()); assertEquals(false, databaseSchema.isNeedLastCache()); } - - @Test - public void testConstructDatabaseSchemaSetsPreferredDataNodeWhenPresent() throws Exception { - final DatabaseSchemaStatement statement = - new DatabaseSchemaStatement(DatabaseSchemaStatement.DatabaseSchemaStatementType.ALTER); - statement.setDatabasePath(new PartialPath("root.sg")); - statement.setPreferredDataNodeId(7); - - final TDatabaseSchema databaseSchema = DatabaseSchemaTask.constructDatabaseSchema(statement); - - assertTrue(databaseSchema.isSetPreferredDataNodeId()); - assertEquals(7, databaseSchema.getPreferredDataNodeId()); - } } diff --git a/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift b/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift index 9ff194dfa9c83..ca05dba4138dc 100644 --- a/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift +++ b/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift @@ -241,7 +241,6 @@ struct TDatabaseSchema { 10: optional i64 timePartitionOrigin 11: optional bool isTableModel 12: optional bool needLastCache - 13: optional i32 preferredDataNodeId } // Schema