This is an automated email from the ASF dual-hosted git repository. spricoder pushed a commit to branch feature/memory_auto in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 1e6e12630e3e4c2f7c63dc0e55d3b9d861520b71 Merge: dfb98462bd2 1a4ac0ea0ad Author: spricoder <[email protected]> AuthorDate: Thu Feb 27 11:25:25 2025 +0800 merge feature/memory_collect .github/workflows/vulnerability-check.yml | 7 +- .../org/apache/iotdb/ainode/it/AINodeBasicIT.java | 2 + .../org/apache/iotdb/jdbc/IoTDBJDBCResultSet.java | 25 +- .../apache/iotdb/jdbc/IoTDBPreparedStatement.java | 31 +- .../IoTDBRelationalDatabaseMetadata.java | 52 +++- .../base/AbstractSubscriptionConsumer.java | 22 +- .../iotdb/confignode/manager/ConfigManager.java | 9 + .../apache/iotdb/confignode/manager/IManager.java | 3 + .../iotdb/confignode/manager/ModelManager.java | 5 + .../manager/load/balancer/RouteBalancer.java | 152 ++++++--- .../iotdb/confignode/manager/node/NodeManager.java | 8 + .../iotdb/confignode/persistence/ModelInfo.java | 2 +- .../procedure/impl/StateMachineProcedure.java | 9 +- .../impl/pipe/AbstractOperatePipeProcedureV2.java | 2 + .../iotdb/confignode/service/ConfigNode.java | 44 +-- .../consensus/config/PipeConsensusConfig.java | 21 +- .../consensus/pipe/PipeConsensusServerImpl.java | 21 +- ...xManager.java => ReplicateProgressManager.java} | 8 +- .../pipe/metric/PipeConsensusSyncLagManager.java | 71 ++++- .../java/org/apache/iotdb/db/conf/IoTDBConfig.java | 286 +---------------- .../org/apache/iotdb/db/conf/IoTDBDescriptor.java | 105 ++++--- .../db/consensus/DataRegionConsensusImpl.java | 8 +- .../client/IoTDBDataNodeAsyncClientManager.java | 4 + .../evolvable/batch/PipeTabletEventBatch.java | 37 ++- .../evolvable/batch/PipeTabletEventPlainBatch.java | 41 +-- .../batch/PipeTabletEventTsFileBatch.java | 11 +- .../pipeconsensus/PipeConsensusAsyncConnector.java | 16 +- .../pipeconsensus/PipeConsensusSyncConnector.java | 9 +- .../PipeConsensusTabletInsertionEventHandler.java | 8 +- .../PipeConsensusTsFileInsertionEventHandler.java | 8 +- .../PipeConsensusTransferBatchReqBuilder.java | 5 +- .../async/IoTDBDataRegionAsyncConnector.java | 1 + ....java => ReplicateProgressDataNodeManager.java} | 36 ++- .../deletion/DeletionResourceManager.java | 4 +- .../deletion/persist/PageCacheDeletionBuffer.java | 5 +- .../common/tsfile/PipeTsFileInsertionEvent.java | 17 +- .../event/realtime/PipeRealtimeEventFactory.java | 71 ++++- ...oricalDataRegionTsFileAndDeletionExtractor.java | 17 ++ .../realtime/assigner/PipeDataRegionAssigner.java | 2 + .../listener/PipeInsertionDataNodeListener.java | 16 +- .../pipeconsensus/PipeConsensusProcessor.java | 43 ++- .../pipeconsensus/PipeConsensusReceiver.java | 338 +++++++++------------ .../db/pipe/resource/memory/PipeMemoryManager.java | 5 +- .../execution/memory/LocalMemoryManager.java | 3 +- .../analyze/cache/partition/PartitionCache.java | 4 +- .../cache/schema/DataNodeDevicePathCache.java | 4 +- .../plan/planner/LocalExecutionPlanner.java | 6 +- .../fetcher/cache/TableDeviceSchemaCache.java | 4 +- .../metric/SchemaEngineCachedMetric.java | 7 +- .../rescon/MemSchemaEngineStatistics.java | 5 +- .../ReleaseFlushStrategySizeBasedImpl.java | 7 +- .../java/org/apache/iotdb/db/service/DataNode.java | 64 +--- .../iotdb/db/service/ExternalRPCService.java | 2 +- ...viceMBean.java => ExternalRPCServiceMBean.java} | 2 +- .../metrics/memory/ConsensusMemoryMetrics.java | 6 +- .../metrics/memory/GlobalMemoryMetrics.java | 6 +- .../metrics/memory/OffHeapMemoryMetrics.java | 6 +- .../metrics/memory/QueryEngineMemoryMetrics.java | 30 +- .../metrics/memory/SchemaEngineMemoryMetrics.java | 16 +- .../metrics/memory/StorageEngineMemoryMetrics.java | 28 +- .../metrics/memory/StreamEngineMemoryMetrics.java | 6 +- .../iotdb/db/storageengine/StorageEngine.java | 19 +- .../db/storageengine/buffer/BloomFilterCache.java | 4 +- .../iotdb/db/storageengine/buffer/ChunkCache.java | 6 +- .../buffer/TimeSeriesMetadataCache.java | 4 +- .../storageengine/dataregion/DataRegionInfo.java | 6 +- .../dataregion/memtable/AbstractMemTable.java | 61 +++- .../memtable/AlignedWritableMemChunk.java | 13 + .../memtable/AlignedWritableMemChunkGroup.java | 15 +- .../dataregion/memtable/IMemTable.java | 2 + .../memtable/IWritableMemChunkGroup.java | 2 + .../dataregion/memtable/TsFileProcessor.java | 41 ++- .../dataregion/memtable/WritableMemChunk.java | 8 + .../dataregion/memtable/WritableMemChunkGroup.java | 20 ++ .../dataregion/tsfile/TsFileResource.java | 8 +- .../dataregion/wal/utils/WALInsertNodeCache.java | 4 +- .../rescon/memory/PrimitiveArrayManager.java | 5 +- .../db/storageengine/rescon/memory/SystemInfo.java | 18 +- .../rescon/memory/TimePartitionManager.java | 5 +- .../rescon/memory/TsFileResourceManager.java | 6 +- .../db/utils/datastructure/AlignedTVList.java | 44 ++- .../iotdb/db/utils/datastructure/BinaryTVList.java | 38 ++- .../db/utils/datastructure/BooleanTVList.java | 38 ++- .../iotdb/db/utils/datastructure/DoubleTVList.java | 38 ++- .../iotdb/db/utils/datastructure/FloatTVList.java | 38 ++- .../iotdb/db/utils/datastructure/IntTVList.java | 38 ++- .../iotdb/db/utils/datastructure/LongTVList.java | 38 ++- .../iotdb/db/utils/datastructure/TVList.java | 70 ++++- .../fetcher/cache/TableDeviceSchemaCacheTest.java | 10 +- .../dataregion/memtable/TsFileProcessorTest.java | 36 +-- .../rescon/memory/ResourceManagerTest.java | 5 +- .../rescon/memory/TimePartitionManagerTest.java | 7 +- .../datastructure/PrimitiveArrayManagerTest.java | 5 +- .../async/AsyncPipeDataTransferServiceClient.java | 14 + .../commons/memory/AtomicLongMemoryBlock.java | 36 +-- .../apache/iotdb/commons/memory/IMemoryBlock.java | 18 +- .../apache/iotdb/commons/memory/MemoryConfig.java | 315 +++++++++++++++++++ .../iotdb/commons/memory/MemoryException.java | 7 +- .../apache/iotdb/commons/memory/MemoryManager.java | 33 +- .../iotdb/commons/pipe/event/EnrichedEvent.java | 17 +- .../iotdb/commons/memory/MemoryManagerTest.java | 2 +- .../thrift-commons/src/main/thrift/common.thrift | 1 + .../src/main/thrift/pipeconsensus.thrift | 5 +- 103 files changed, 1727 insertions(+), 1166 deletions(-) diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java index 5c7f5ad7242,d60920c0b31..339acc919c5 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java @@@ -113,7 -115,9 +115,9 @@@ public class DataRegionConsensusImpl private static ConsensusConfig buildConsensusConfig() { IMemoryBlock memoryBlock = - CONF.getConsensusMemoryManager().forceAllocate("Consensus", MemoryBlockType.DYNAMIC); + MEMORY_CONFIG + .getConsensusMemoryManager() - .forceAllocate("Consensus", MemoryBlockType.FUNCTION); ++ .forceAllocate("Consensus", MemoryBlockType.DYNAMIC); return ConsensusConfig.newBuilder() .setThisNodeId(CONF.getDataNodeId()) .setThisNode(new TEndPoint(CONF.getInternalAddress(), CONF.getDataRegionConsensusPort())) diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManager.java index bede51da24d,e72c8ff3375..80784009b08 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManager.java @@@ -53,10 -53,9 +53,9 @@@ public class PipeMemoryManager // TODO @spricoder: consider combine memory block and used MemorySizeInBytes private IMemoryBlock memoryBlock = - IoTDBDescriptor.getInstance() - .getConfig() + MemoryConfig.getInstance() .getPipeMemoryManager() - .forceAllocate("Stream", MemoryBlockType.FUNCTION); + .forceAllocate("Stream", MemoryBlockType.DYNAMIC); private static final double EXCEED_PROTECT_THRESHOLD = 0.95; diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/cache/partition/PartitionCache.java index 143d92e62fa,17f0a82dda5..e1de5c71acb --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/cache/partition/PartitionCache.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/cache/partition/PartitionCache.java @@@ -119,9 -121,9 +121,9 @@@ public class PartitionCache public PartitionCache() { this.memoryBlock = - config + memoryConfig .getPartitionCacheMemoryManager() - .forceAllocate("PartitionCache", MemoryBlockType.FUNCTION); + .forceAllocate("PartitionCache", MemoryBlockType.STATIC); this.memoryBlock.allocate(this.memoryBlock.getTotalMemorySizeInBytes()); // TODO @spricoder: PartitionCache need to be controlled according to memory this.schemaPartitionCache = diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/cache/schema/DataNodeDevicePathCache.java index fca17a11430,bd3c49bf998..ce3d1bc1737 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/cache/schema/DataNodeDevicePathCache.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/cache/schema/DataNodeDevicePathCache.java @@@ -43,9 -45,9 +45,9 @@@ public class DataNodeDevicePathCache private DataNodeDevicePathCache() { devicePathCacheMemoryBlock = - config + memoryConfig .getDevicePathCacheMemoryManager() - .forceAllocate("DevicePathCache", MemoryBlockType.PERFORMANCE); + .forceAllocate("DevicePathCache", MemoryBlockType.STATIC); // TODO @spricoder: later we can find a way to get the byte size of cache devicePathCacheMemoryBlock.allocate(devicePathCacheMemoryBlock.getTotalMemorySizeInBytes()); devicePathCache = diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/LocalExecutionPlanner.java index 93f45bebb5e,b380f34ce89..1f799e7401b --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/LocalExecutionPlanner.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/LocalExecutionPlanner.java @@@ -67,9 -68,12 +68,12 @@@ public class LocalExecutionPlanner static { IoTDBConfig CONFIG = IoTDBDescriptor.getInstance().getConfig(); + MemoryConfig MEMORY_CONFIG = MemoryConfig.getInstance(); OPERATORS_MEMORY_BLOCK = - CONFIG.getOperatorsMemoryManager().forceAllocate("Operators", MemoryBlockType.DYNAMIC); + MEMORY_CONFIG + .getOperatorsMemoryManager() - .forceAllocate("Operators", MemoryBlockType.FUNCTION); ++ .forceAllocate("Operators", MemoryBlockType.DYNAMIC); MIN_REST_MEMORY_FOR_QUERY_AFTER_LOAD = (long) ((OPERATORS_MEMORY_BLOCK.getTotalMemorySizeInBytes()) diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/fetcher/cache/TableDeviceSchemaCache.java index 604208b1edc,06233626a28..1380c06e883 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/fetcher/cache/TableDeviceSchemaCache.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/fetcher/cache/TableDeviceSchemaCache.java @@@ -106,9 -108,9 +108,9 @@@ public class TableDeviceSchemaCache private TableDeviceSchemaCache() { memoryBlock = - config + memoryConfig .getSchemaCacheMemoryManager() - .forceAllocate("TableDeviceSchemaCache", MemoryBlockType.FUNCTION); + .forceAllocate("TableDeviceSchemaCache", MemoryBlockType.STATIC); dualKeyCache = new DualKeyCacheBuilder<TableId, IDeviceID, TableDeviceCacheEntry>() .cacheEvictionPolicy( diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/rescon/MemSchemaEngineStatistics.java index 51d6ad337a7,32d3c9c2f75..e6c44c52e43 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/rescon/MemSchemaEngineStatistics.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/rescon/MemSchemaEngineStatistics.java @@@ -56,10 -56,9 +56,9 @@@ public class MemSchemaEngineStatistics public MemSchemaEngineStatistics() { memoryBlock = - IoTDBDescriptor.getInstance() - .getConfig() + MemoryConfig.getInstance() .getSchemaRegionMemoryManager() - .forceAllocate("SchemaRegion", MemoryBlockType.FUNCTION); + .forceAllocate("SchemaRegion", MemoryBlockType.DYNAMIC); } @Override diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/buffer/BloomFilterCache.java index be2a48c19aa,be23fe45db5..8dbf8691310 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/buffer/BloomFilterCache.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/buffer/BloomFilterCache.java @@@ -58,9 -60,9 +60,9 @@@ public class BloomFilterCache static { CACHE_MEMORY_BLOCK = - CONFIG + MEMORY_CONFIG .getBloomFilterCacheMemoryManager() - .forceAllocate("BloomFilterCache", MemoryBlockType.PERFORMANCE); + .forceAllocate("BloomFilterCache", MemoryBlockType.STATIC); // TODO @spricoder: find a way to get the size of the BloomFilterCache CACHE_MEMORY_BLOCK.allocate(CACHE_MEMORY_BLOCK.getTotalMemorySizeInBytes()); } diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/buffer/ChunkCache.java index d964a635b8e,301480239bd..30d17a41ce0 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/buffer/ChunkCache.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/buffer/ChunkCache.java @@@ -73,7 -75,9 +75,9 @@@ public class ChunkCache static { CACHE_MEMORY_BLOCK = - CONFIG.getChunkCacheMemoryManager().forceAllocate("ChunkCache", MemoryBlockType.STATIC); + MEMORY_CONFIG + .getChunkCacheMemoryManager() - .forceAllocate("ChunkCache", MemoryBlockType.PERFORMANCE); ++ .forceAllocate("ChunkCache", MemoryBlockType.STATIC); // TODO @spricoder: find a way to get the size of the ChunkCache CACHE_MEMORY_BLOCK.allocate(CACHE_MEMORY_BLOCK.getTotalMemorySizeInBytes()); } diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/buffer/TimeSeriesMetadataCache.java index 3ba8e8a4af3,b7d85968f31..3e09d144a1f --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/buffer/TimeSeriesMetadataCache.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/buffer/TimeSeriesMetadataCache.java @@@ -85,9 -87,9 +87,9 @@@ public class TimeSeriesMetadataCache static { CACHE_MEMORY_BLOCK = - config + memoryConfig .getTimeSeriesMetaDataCacheMemoryManager() - .forceAllocate("TimeSeriesMetadataCache", MemoryBlockType.PERFORMANCE); + .forceAllocate("TimeSeriesMetadataCache", MemoryBlockType.STATIC); // TODO @spricoder find a better way to get the size of cache CACHE_MEMORY_BLOCK.allocate(CACHE_MEMORY_BLOCK.getTotalMemorySizeInBytes()); } diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/PrimitiveArrayManager.java index 1b9b523dd12,ad73a1bdc9c..95095655312 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/PrimitiveArrayManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/PrimitiveArrayManager.java @@@ -53,9 -56,9 +56,9 @@@ public class PrimitiveArrayManager /** memory block for all arrays */ private static final IMemoryBlock POOLED_ARRAYS_MEMORY_BLOCK = - CONFIG + MEMORY_CONFIG .getBufferedArraysMemoryManager() - .forceAllocate("BufferedArrays", MemoryBlockType.FUNCTION); + .forceAllocate("BufferedArrays", MemoryBlockType.DYNAMIC); /** threshold total size of arrays for all data types */ private static final double POOLED_ARRAYS_MEMORY_THRESHOLD = diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/SystemInfo.java index 7819e9b5a53,ba8d704d9fb..46ef78ac47a --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/SystemInfo.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/SystemInfo.java @@@ -78,13 -80,17 +80,17 @@@ public class SystemInfo private SystemInfo() { compactionMemoryBlock = - config.getCompactionMemoryManager().forceAllocate("Compaction", MemoryBlockType.DYNAMIC); + memoryConfig + .getCompactionMemoryManager() - .forceAllocate("Compaction", MemoryBlockType.FUNCTION); ++ .forceAllocate("Compaction", MemoryBlockType.DYNAMIC); walBufferQueueMemoryBlock = - config.getWalBufferQueueManager().forceAllocate("WalBufferQueue", MemoryBlockType.DYNAMIC); + memoryConfig + .getWalBufferQueueManager() - .forceAllocate("WalBufferQueue", MemoryBlockType.FUNCTION); ++ .forceAllocate("WalBufferQueue", MemoryBlockType.DYNAMIC); directBufferMemoryBlock = - config + memoryConfig .getDirectBufferMemoryManager() - .forceAllocate("DirectBuffer", MemoryBlockType.FUNCTION); + .forceAllocate("DirectBuffer", MemoryBlockType.DYNAMIC); loadWriteMemory(); } diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/TimePartitionManager.java index 85e96683413,2d8b2206c06..53f44867ff2 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/TimePartitionManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/TimePartitionManager.java @@@ -48,10 -48,9 +48,9 @@@ public class TimePartitionManager private TimePartitionManager() { timePartitionInfoMap = new HashMap<>(); timePartitionInfoMemoryBlock = - IoTDBDescriptor.getInstance() - .getConfig() + MemoryConfig.getInstance() .getTimePartitionInfoMemoryManager() - .forceAllocate("TimePartitionInfoMemoryBlock", MemoryBlockType.FUNCTION); + .forceAllocate("TimePartitionInfoMemoryBlock", MemoryBlockType.DYNAMIC); } public void registerTimePartitionInfo(TimePartitionInfo timePartitionInfo) { diff --cc iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/TsFileResourceManager.java index 2cecc1ed3a6,50fbc3a62c2..9d929967ab7 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/TsFileResourceManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/TsFileResourceManager.java @@@ -49,7 -51,9 +51,9 @@@ public class TsFileResourceManager private TsFileResourceManager() { memoryBlock = - CONFIG.getTimeIndexMemoryManager().forceAllocate("TimeIndex", MemoryBlockType.DYNAMIC); + MEMORY_CONFIG + .getTimeIndexMemoryManager() - .forceAllocate("TimeIndex", MemoryBlockType.FUNCTION); ++ .forceAllocate("TimeIndex", MemoryBlockType.DYNAMIC); } @TestOnly diff --cc iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/memory/MemoryManager.java index be4210a6226,360418303fe..405ddff7523 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/memory/MemoryManager.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/memory/MemoryManager.java @@@ -73,11 -69,9 +73,10 @@@ public class MemoryManager this.enable = false; } - private MemoryManager( - String name, MemoryManager parentMemoryManager, long totalMemorySizeInBytes) { + MemoryManager(String name, MemoryManager parentMemoryManager, long totalMemorySizeInBytes) { this.name = name; this.parentMemoryManager = parentMemoryManager; + this.totalAllocatedMemorySizeInBytes = totalMemorySizeInBytes; this.totalMemorySizeInBytes = totalMemorySizeInBytes; this.enable = false; }
