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..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 @@ -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; @@ -39,6 +40,8 @@ import org.apache.iotdb.confignode.manager.partition.PartitionManager; import org.apache.iotdb.confignode.manager.schema.ClusterSchemaManager; +import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -82,6 +85,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 +129,12 @@ 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 : -1) + : -1; // Only considering the specified Database when doing allocation final List databaseAllocatedRegionGroups = getPartitionManager().getAllReplicaSets(database, consensusGroupType); @@ -144,6 +161,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 +193,37 @@ private ProcedureManager getProcedureManager() { return configManager.getProcedureManager(); } + 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..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,6 +337,15 @@ public DataPartition getOrCreateDataPartition( @Override public DataPartition getOrCreateDataPartition( final List dataPartitionQueryParams, final String userName) { + return getOrCreateDataPartition( + dataPartitionQueryParams, userName, getLoadPreferredDataNodeId()); + } + + @Override + public DataPartition getOrCreateDataPartition( + final List dataPartitionQueryParams, + final String userName, + final int preferredDataNodeId) { final Map> splitDataPartitionQueryParams = splitDataPartitionQueryParam( dataPartitionQueryParams, config.isAutoCreateSchemaEnabled(), userName); @@ -330,6 +357,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/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/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-protocol/thrift-confignode/src/main/thrift/confignode.thrift b/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift index 512e2e05f05fc..ca05dba4138dc 100644 --- a/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift +++ b/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift @@ -279,6 +279,7 @@ struct TTimeSlotList { struct TDataPartitionReq { // map> 1: required map> partitionSlotsMap + 2: optional i32 preferredDataNodeId } struct TDataPartitionTableResp {