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();