This is an automated email from the ASF dual-hosted git repository. Wei-hao-Li pushed a commit to branch fixDeviceEntryConf in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit f2108d51332bb8fa20339e26d2c5d75e6687943d Author: Weihao Li <[email protected]> AuthorDate: Thu Sep 17 10:44:18 2026 +0800 fix conf Signed-off-by: Weihao Li <[email protected]> --- .../org/apache/iotdb/db/conf/DataNodeMemoryConfig.java | 18 ++++++++++++++++++ .../java/org/apache/iotdb/db/conf/IoTDBConfig.java | 14 -------------- .../java/org/apache/iotdb/db/conf/IoTDBDescriptor.java | 15 ++++++++++----- .../metadata/fetcher/TableDeviceSchemaFetcher.java | 3 ++- .../distribute/TableDistributedPlanGenerator.java | 10 +++++----- .../java/org/apache/iotdb/db/conf/PropertiesTest.java | 15 ++++++++++----- 6 files changed, 45 insertions(+), 30 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/DataNodeMemoryConfig.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/DataNodeMemoryConfig.java index a6fe6e0ef99..1ef9480658e 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/DataNodeMemoryConfig.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/DataNodeMemoryConfig.java @@ -133,6 +133,9 @@ public class DataNodeMemoryConfig { /** Memory manager for operators */ private MemoryManager operatorsMemoryManager; + /** Maximum DeviceEntry bytes kept in memory before a table-query spill. */ + private long tableQueryDeviceEntryBatchSizeInBytes; + /** Memory manager for operators */ private MemoryManager dataExchangeMemoryManager; @@ -624,6 +627,13 @@ public class DataNodeMemoryConfig { setQueryThreadCount( Integer.parseInt( properties.getProperty("query_thread_count", Integer.toString(getQueryThreadCount())))); + + tableQueryDeviceEntryBatchSizeInBytes = + Long.parseLong( + properties.getProperty( + "table_query_device_entry_batch_size_in_bytes", + Long.toString( + operatorsMemoryManager.getTotalMemorySizeInBytes() / queryThreadCount / 4))); } public double getRejectProportion() { @@ -774,6 +784,14 @@ public class DataNodeMemoryConfig { return operatorsMemoryManager; } + public long getTableQueryDeviceEntryBatchSizeInBytes() { + return tableQueryDeviceEntryBatchSizeInBytes; + } + + public void setTableQueryDeviceEntryBatchSizeInBytes(long tableQueryDeviceEntryBatchSizeInBytes) { + this.tableQueryDeviceEntryBatchSizeInBytes = tableQueryDeviceEntryBatchSizeInBytes; + } + public MemoryManager getDataExchangeMemoryManager() { return dataExchangeMemoryManager; } 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 80f9f1f1b30..44a9daab3a8 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 @@ -260,12 +260,6 @@ public class IoTDBConfig { private String queryDir = IoTDBConstant.DN_DEFAULT_DATA_DIR + File.separator + IoTDBConstant.QUERY_FOLDER_NAME; - /** - * Maximum DeviceEntry bytes kept in memory before a table-query spill, capped by the effective - * Thrift frame size minus 1 KiB reserved for the RPC response envelope. - */ - private long tableQueryDeviceEntryBatchSizeInBytes; - /** External lib directory, stores user-uploaded JAR files */ private String extDir = IoTDBConstant.EXT_FOLDER_NAME; @@ -1821,14 +1815,6 @@ public class IoTDBConfig { this.queryDir = queryDir; } - public long getTableQueryDeviceEntryBatchSizeInBytes() { - return tableQueryDeviceEntryBatchSizeInBytes; - } - - public void setTableQueryDeviceEntryBatchSizeInBytes(long tableQueryDeviceEntryBatchSizeInBytes) { - this.tableQueryDeviceEntryBatchSizeInBytes = tableQueryDeviceEntryBatchSizeInBytes; - } - public String getRatisDataRegionSnapshotDir() { return ratisDataRegionSnapshotDir; } 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 565f0b7832a..b17536ef1ae 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 @@ -168,9 +168,14 @@ public class IoTDBDescriptor { .getConfig() .setCustomizedProperties(loader.getCustomizedProperties()); } - // If no configuration source initialized the memory config, initialize it with defaults. + // If no configuration source initialized the config, run the normal loading path with defaults. if (!hasLoadedProperties && !hasProperties) { - memoryConfig.init(new TrimProperties()); + try { + loadProperties(new TrimProperties()); + } catch (Exception e) { + LOGGER.error(DataNodeMiscMessages.INCORRECT_FORMAT_CONFIG_FILE, e); + System.exit(-1); + } } } @@ -2457,7 +2462,7 @@ public class IoTDBDescriptor { "mods_cache_size_limit_per_fi_in_bytes", Long.toString(conf.getModsCacheSizeLimitPerFI())); ConfigurationFileUtils.updateAppliedProperties( "table_query_device_entry_batch_size_in_bytes", - Long.toString(conf.getTableQueryDeviceEntryBatchSizeInBytes())); + Long.toString(memoryConfig.getTableQueryDeviceEntryBatchSizeInBytes())); ConfigurationFileUtils.updateAppliedProperties( DEFAULT_WAL_THRESHOLD_NAME[1], Long.toString(conf.getThrottleThreshold())); } @@ -2467,7 +2472,7 @@ public class IoTDBDescriptor { Long.parseLong( properties.getProperty( "table_query_device_entry_batch_size_in_bytes", - Long.toString(conf.getTableQueryDeviceEntryBatchSizeInBytes()))); + Long.toString(memoryConfig.getTableQueryDeviceEntryBatchSizeInBytes()))); if (deviceEntryBatchSize <= 0) { deviceEntryBatchSize = memoryConfig.getOperatorsMemoryManager().getTotalMemorySizeInBytes() @@ -2486,7 +2491,7 @@ public class IoTDBDescriptor { conf.getThriftMaxFrameSize(), effectiveBatchSize)); } - conf.setTableQueryDeviceEntryBatchSizeInBytes(effectiveBatchSize); + memoryConfig.setTableQueryDeviceEntryBatchSizeInBytes(effectiveBatchSize); } private void loadQuerySampleThroughput(TrimProperties properties) throws IOException { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/fetcher/TableDeviceSchemaFetcher.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/fetcher/TableDeviceSchemaFetcher.java index c78ec5de9aa..1c797046f40 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/fetcher/TableDeviceSchemaFetcher.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/fetcher/TableDeviceSchemaFetcher.java @@ -330,7 +330,8 @@ public class TableDeviceSchemaFetcher { private AbstractDeviceEntryMaterializer createDataSetMaterializer( MPPQueryContext queryContext, PlanNodeId planNodeId, boolean distinct) { - long batchSize = CONFIG.getTableQueryDeviceEntryBatchSizeInBytes(); + long batchSize = + IoTDBDescriptor.getInstance().getMemoryConfig().getTableQueryDeviceEntryBatchSizeInBytes(); if (distinct) { return new DeviceEntrySortedMaterializer( queryContext.getQueryId().getId(), diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/distribute/TableDistributedPlanGenerator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/distribute/TableDistributedPlanGenerator.java index 00db7b95e49..ec095730f54 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/distribute/TableDistributedPlanGenerator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/distribute/TableDistributedPlanGenerator.java @@ -988,7 +988,7 @@ public class TableDistributedPlanGenerator Comparator<DeviceEntry> comparator = sortPropertyContext.map(property -> property.comparator).orElse(null); long batchSize = - IoTDBDescriptor.getInstance().getConfig().getTableQueryDeviceEntryBatchSizeInBytes(); + IoTDBDescriptor.getInstance().getMemoryConfig().getTableQueryDeviceEntryBatchSizeInBytes(); Map<TRegionReplicaSet, DeviceTableScanNode> scanNodes = new HashMap<>(); Map<TRegionReplicaSet, AbstractDeviceEntryMaterializer> materializers = new HashMap<>(); Map<TRegionReplicaSet, Integer> regionEntryCounts = new HashMap<>(); @@ -1211,7 +1211,7 @@ public class TableDistributedPlanGenerator Comparator<DeviceEntry> comparator = sortPropertyContext.map(property -> property.comparator).orElse(null); long batchSize = - IoTDBDescriptor.getInstance().getConfig().getTableQueryDeviceEntryBatchSizeInBytes(); + IoTDBDescriptor.getInstance().getMemoryConfig().getTableQueryDeviceEntryBatchSizeInBytes(); Map<TRegionReplicaSet, DeviceTableScanNode> scanNodes = new HashMap<>(); Map<TRegionReplicaSet, AbstractDeviceEntryMaterializer> materializers = new HashMap<>(); Map<TRegionReplicaSet, Integer> regionEntryCounts = new HashMap<>(); @@ -1534,7 +1534,7 @@ public class TableDistributedPlanGenerator Comparator<DeviceEntry> comparator = sortPropertyContext.map(property -> property.comparator).orElse(null); long batchSize = - IoTDBDescriptor.getInstance().getConfig().getTableQueryDeviceEntryBatchSizeInBytes(); + IoTDBDescriptor.getInstance().getMemoryConfig().getTableQueryDeviceEntryBatchSizeInBytes(); Map<TRegionReplicaSet, Pair<TreeAlignedDeviceViewScanNode, TreeNonAlignedDeviceViewScanNode>> scanNodes = new HashMap<>(); Map<DeviceTableScanNode, AbstractDeviceEntryMaterializer> materializers = new HashMap<>(); @@ -2097,7 +2097,7 @@ public class TableDistributedPlanGenerator } long batchSize = - IoTDBDescriptor.getInstance().getConfig().getTableQueryDeviceEntryBatchSizeInBytes(); + IoTDBDescriptor.getInstance().getMemoryConfig().getTableQueryDeviceEntryBatchSizeInBytes(); Map<Integer, List<TRegionReplicaSet>> cachedSeriesSlotWithRegions = new HashMap<>(); Map<TRegionReplicaSet, PlanNodeId> regionPlanNodeIds = new HashMap<>(); Map<TRegionReplicaSet, AbstractDeviceEntryMaterializer> stagingMaterializers = new HashMap<>(); @@ -2473,7 +2473,7 @@ public class TableDistributedPlanGenerator } long batchSize = - IoTDBDescriptor.getInstance().getConfig().getTableQueryDeviceEntryBatchSizeInBytes(); + IoTDBDescriptor.getInstance().getMemoryConfig().getTableQueryDeviceEntryBatchSizeInBytes(); boolean hasCrossRegionDevice = false; Map<Integer, List<TRegionReplicaSet>> cachedSeriesSlotWithRegions = new HashMap<>(); Map<TRegionReplicaSet, Pair<PlanNodeId, PlanNodeId>> regionPlanNodeIds = new HashMap<>(); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/conf/PropertiesTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/conf/PropertiesTest.java index 92e52da7fa1..be8b29a064b 100755 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/conf/PropertiesTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/conf/PropertiesTest.java @@ -74,7 +74,8 @@ public class PropertiesTest { final IoTDBDescriptor descriptor = IoTDBDescriptor.getInstance(); final IoTDBConfig config = descriptor.getConfig(); final int originalFrameSize = config.getThriftMaxFrameSize(); - final long originalBatchSize = config.getTableQueryDeviceEntryBatchSizeInBytes(); + final long originalBatchSize = + descriptor.getMemoryConfig().getTableQueryDeviceEntryBatchSizeInBytes(); try { final TrimProperties properties = new TrimProperties(); @@ -83,7 +84,8 @@ public class PropertiesTest { descriptor.loadProperties(properties); Assert.assertEquals(4096, config.getThriftMaxFrameSize()); - Assert.assertEquals(3072, config.getTableQueryDeviceEntryBatchSizeInBytes()); + Assert.assertEquals( + 3072, descriptor.getMemoryConfig().getTableQueryDeviceEntryBatchSizeInBytes()); Assert.assertEquals( "3072", ConfigurationFileUtils.getAppliedProperties() @@ -102,7 +104,8 @@ public class PropertiesTest { final IoTDBDescriptor descriptor = IoTDBDescriptor.getInstance(); final IoTDBConfig config = descriptor.getConfig(); final int originalFrameSize = config.getThriftMaxFrameSize(); - final long originalBatchSize = config.getTableQueryDeviceEntryBatchSizeInBytes(); + final long originalBatchSize = + descriptor.getMemoryConfig().getTableQueryDeviceEntryBatchSizeInBytes(); try { config.setThriftMaxFrameSize(4096); @@ -110,7 +113,8 @@ public class PropertiesTest { properties.setProperty("table_query_device_entry_batch_size_in_bytes", "4096"); descriptor.loadHotModifiedProps(properties); - Assert.assertEquals(3072, config.getTableQueryDeviceEntryBatchSizeInBytes()); + Assert.assertEquals( + 3072, descriptor.getMemoryConfig().getTableQueryDeviceEntryBatchSizeInBytes()); Assert.assertEquals( "3072", ConfigurationFileUtils.getAppliedProperties() @@ -118,7 +122,8 @@ public class PropertiesTest { properties.setProperty("table_query_device_entry_batch_size_in_bytes", "512"); descriptor.loadHotModifiedProps(properties); - Assert.assertEquals(512, config.getTableQueryDeviceEntryBatchSizeInBytes()); + Assert.assertEquals( + 512, descriptor.getMemoryConfig().getTableQueryDeviceEntryBatchSizeInBytes()); } finally { config.setThriftMaxFrameSize(originalFrameSize); final TrimProperties properties = new TrimProperties();
