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

JackieTien97 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new 5d67dae28c8 Fix init of DeviceEntryBatchSizeInBytes config when there 
is no properties file (#18660)
5d67dae28c8 is described below

commit 5d67dae28c814fa62bea6e446ec1fc6b87241c8c
Author: Weihao Li <[email protected]>
AuthorDate: Thu Sep 17 17:03:11 2026 +0800

    Fix init of DeviceEntryBatchSizeInBytes config when there is no properties 
file (#18660)
---
 .../apache/iotdb/db/conf/DataNodeMemoryConfig.java | 42 ++++++++++++++++++-
 .../java/org/apache/iotdb/db/conf/IoTDBConfig.java | 14 -------
 .../org/apache/iotdb/db/conf/IoTDBDescriptor.java  | 49 ++++------------------
 .../metadata/fetcher/TableDeviceSchemaFetcher.java |  3 +-
 .../distribute/TableDistributedPlanGenerator.java  | 10 ++---
 .../org/apache/iotdb/db/conf/PropertiesTest.java   | 15 ++++---
 6 files changed, 67 insertions(+), 66 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..12ac2293d3d 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
@@ -37,6 +37,7 @@ public class DataNodeMemoryConfig {
   private static final int LEGACY_QUERY_MEMORY_COMPONENT_COUNT = 8;
   private static final int QUERY_MEMORY_COMPONENT_COUNT = 9;
   private static final int SUBSCRIPTION_MEMORY_INDEX = 8;
+  private static final int DEVICE_ENTRY_RPC_FRAME_RESERVED_BYTES = 1024;
   private static final int[] DEFAULT_QUERY_MEMORY_PROPORTIONS = {
     1, 100, 200, 50, 200, 200, 200, 50, 250
   };
@@ -133,6 +134,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;
 
@@ -175,7 +179,7 @@ public class DataNodeMemoryConfig {
         getMemoryAllocateProportion(properties, false));
   }
 
-  public void init(TrimProperties properties) {
+  public void init(TrimProperties properties, int thriftMaxFrameSize, Logger 
logger) {
     // on heap memory
     String memoryAllocateProportion = getMemoryAllocateProportion(properties, 
true);
     // Get global memory manager here
@@ -252,6 +256,8 @@ public class DataNodeMemoryConfig {
     initSchemaMemoryAllocate(schemaEngineMemoryManager, properties);
     initStorageEngineAllocate(storageEngineMemoryManager, properties);
     initQueryEngineMemoryAllocate(queryEngineMemoryManager, properties);
+    tableQueryDeviceEntryBatchSizeInBytes = 0;
+    loadTableQueryDeviceEntryBatchSize(properties, thriftMaxFrameSize, logger);
 
     String offHeapMemoryStr = System.getProperty("OFF_HEAP_MEMORY");
     offHeapMemoryManager =
@@ -626,6 +632,32 @@ public class DataNodeMemoryConfig {
             properties.getProperty("query_thread_count", 
Integer.toString(getQueryThreadCount()))));
   }
 
+  public void loadTableQueryDeviceEntryBatchSize(
+      TrimProperties properties, int thriftMaxFrameSize, Logger logger) {
+    long defaultTableQueryDeviceEntryBatchSizeInBytes =
+        operatorsMemoryManager.getTotalMemorySizeInBytes() / queryThreadCount 
/ 4;
+    long deviceEntryBatchSize =
+        Long.parseLong(
+            properties.getProperty(
+                "table_query_device_entry_batch_size_in_bytes",
+                Long.toString(tableQueryDeviceEntryBatchSizeInBytes)));
+    if (deviceEntryBatchSize <= 0) {
+      deviceEntryBatchSize = defaultTableQueryDeviceEntryBatchSizeInBytes;
+    }
+    long maxBatchSize = Math.max(1, thriftMaxFrameSize - 
DEVICE_ENTRY_RPC_FRAME_RESERVED_BYTES);
+    long effectiveBatchSize = Math.min(deviceEntryBatchSize, maxBatchSize);
+    if (deviceEntryBatchSize > maxBatchSize) {
+      logger.warn(
+          String.format(
+              DataNodeMiscMessages
+                  
.LOG_TABLE_QUERY_DEVICE_ENTRY_BATCH_SIZE_IN_BYTES_ARG_EXCEEDS_DN_THRIFT_MAX_FRAME_SIZE_ARG_USING_ARG_AS_THE_EFFECTIVE_VALUE_2AE1BEDA,
+              deviceEntryBatchSize,
+              thriftMaxFrameSize,
+              effectiveBatchSize));
+    }
+    tableQueryDeviceEntryBatchSizeInBytes = effectiveBatchSize;
+  }
+
   public double getRejectProportion() {
     return rejectProportion;
   }
@@ -774,6 +806,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..1eff0fcb123 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
@@ -119,8 +119,6 @@ public class IoTDBDescriptor {
 
   private static final double MIN_DIR_USE_PROPORTION = 0.5;
 
-  private static final long DEVICE_ENTRY_RPC_FRAME_RESERVED_BYTES = 1024;
-
   private static final String[] DEFAULT_WAL_THRESHOLD_NAME = {
     "iot_consensus_throttle_threshold_in_byte", 
"wal_throttle_threshold_in_byte"
   };
@@ -170,7 +168,7 @@ public class IoTDBDescriptor {
     }
     // If no configuration source initialized the memory config, initialize it 
with defaults.
     if (!hasLoadedProperties && !hasProperties) {
-      memoryConfig.init(new TrimProperties());
+      memoryConfig.init(new TrimProperties(), conf.getThriftMaxFrameSize(), 
LOGGER);
     }
   }
 
@@ -337,7 +335,11 @@ public class IoTDBDescriptor {
                 "write_memory_variation_report_proportion",
                 
Double.toString(conf.getWriteMemoryVariationReportProportion()))));
 
-    memoryConfig.init(properties);
+    conf.setThriftMaxFrameSize(
+        Integer.parseInt(
+            properties.getProperty(
+                "dn_thrift_max_frame_size", 
String.valueOf(conf.getThriftMaxFrameSize()))));
+    memoryConfig.init(properties, conf.getThriftMaxFrameSize(), LOGGER);
 
     String systemDir = properties.getProperty("dn_system_dir");
     if (systemDir == null) {
@@ -831,13 +833,6 @@ public class IoTDBDescriptor {
             properties.getProperty(
                 "primitive_array_size", 
String.valueOf(conf.getPrimitiveArraySize())))));
 
-    conf.setThriftMaxFrameSize(
-        Integer.parseInt(
-            properties.getProperty(
-                "dn_thrift_max_frame_size", 
String.valueOf(conf.getThriftMaxFrameSize()))));
-
-    loadTableQueryDeviceEntryBatchSize(properties);
-
     conf.setThriftDefaultBufferSize(
         Integer.parseInt(
             properties.getProperty(
@@ -2279,7 +2274,8 @@ public class IoTDBDescriptor {
                   ConfigurationFileUtils.getConfigurationDefaultValue(
                       "enable_topk_runtime_filter"))));
 
-      loadTableQueryDeviceEntryBatchSize(properties);
+      memoryConfig.loadTableQueryDeviceEntryBatchSize(
+          properties, conf.getThriftMaxFrameSize(), LOGGER);
 
       // update wal config
       long prevDeleteWalFilesPeriodInMs = conf.getDeleteWalFilesPeriodInMs();
@@ -2457,38 +2453,11 @@ 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()));
   }
 
-  private void loadTableQueryDeviceEntryBatchSize(TrimProperties properties) {
-    long deviceEntryBatchSize =
-        Long.parseLong(
-            properties.getProperty(
-                "table_query_device_entry_batch_size_in_bytes",
-                
Long.toString(conf.getTableQueryDeviceEntryBatchSizeInBytes())));
-    if (deviceEntryBatchSize <= 0) {
-      deviceEntryBatchSize =
-          memoryConfig.getOperatorsMemoryManager().getTotalMemorySizeInBytes()
-              / memoryConfig.getQueryThreadCount()
-              / 4;
-    }
-    long maxBatchSize =
-        Math.max(1, conf.getThriftMaxFrameSize() - 
DEVICE_ENTRY_RPC_FRAME_RESERVED_BYTES);
-    long effectiveBatchSize = Math.min(deviceEntryBatchSize, maxBatchSize);
-    if (deviceEntryBatchSize > maxBatchSize) {
-      LOGGER.warn(
-          String.format(
-              DataNodeMiscMessages
-                  
.LOG_TABLE_QUERY_DEVICE_ENTRY_BATCH_SIZE_IN_BYTES_ARG_EXCEEDS_DN_THRIFT_MAX_FRAME_SIZE_ARG_USING_ARG_AS_THE_EFFECTIVE_VALUE_2AE1BEDA,
-              deviceEntryBatchSize,
-              conf.getThriftMaxFrameSize(),
-              effectiveBatchSize));
-    }
-    conf.setTableQueryDeviceEntryBatchSizeInBytes(effectiveBatchSize);
-  }
-
   private void loadQuerySampleThroughput(TrimProperties properties) throws 
IOException {
     String querySamplingRateLimitNumber =
         properties.getProperty(
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();

Reply via email to