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

tanxinyu 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 e1840d0ea7a [IoTConsensus] More accurate statistics on IoTConsensus 
memory management (#14965)
e1840d0ea7a is described below

commit e1840d0ea7ae88e79b536c9b781bdd434c147063
Author: Xiangpeng Hu <[email protected]>
AuthorDate: Thu Feb 27 11:30:18 2025 +0800

    [IoTConsensus] More accurate statistics on IoTConsensus memory management 
(#14965)
    
    * add todo
    
    * get memory size
---
 .../iotdb/consensus/common/request/IConsensusRequest.java    |  5 +++++
 .../consensus/common/request/IndexedConsensusRequest.java    |  8 ++++----
 .../iotdb/consensus/iot/logdispatcher/LogDispatcher.java     | 12 ++++++------
 .../queryengine/plan/planner/plan/node/write/InsertNode.java |  1 +
 4 files changed, 16 insertions(+), 10 deletions(-)

diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IConsensusRequest.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IConsensusRequest.java
index 6cd7f370a6b..cae54c2bcbc 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IConsensusRequest.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IConsensusRequest.java
@@ -36,6 +36,11 @@ public interface IConsensusRequest {
    */
   ByteBuffer serializeToByteBuffer();
 
+  default long getMemorySize() {
+    // return 0 by default
+    return 0;
+  }
+
   default void markAsGeneratedByRemoteConsensusLeader() {
     // do nothing by default
   }
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java
index 58789dbd0ed..1147abc049e 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java
@@ -33,7 +33,7 @@ public class IndexedConsensusRequest implements 
IConsensusRequest {
   private final long syncIndex;
   private final List<IConsensusRequest> requests;
   private final List<ByteBuffer> serializedRequests;
-  private long serializedSize = 0;
+  private long memorySize = 0;
 
   public IndexedConsensusRequest(long searchIndex, List<IConsensusRequest> 
requests) {
     this.searchIndex = searchIndex;
@@ -55,7 +55,7 @@ public class IndexedConsensusRequest implements 
IConsensusRequest {
         r -> {
           ByteBuffer buffer = r.serializeToByteBuffer();
           this.serializedRequests.add(buffer);
-          this.serializedSize += buffer.capacity();
+          this.memorySize += Long.max(buffer.capacity(), r.getMemorySize());
         });
   }
 
@@ -72,8 +72,8 @@ public class IndexedConsensusRequest implements 
IConsensusRequest {
     return serializedRequests;
   }
 
-  public long getSerializedSize() {
-    return serializedSize;
+  public long getMemorySize() {
+    return memorySize;
   }
 
   public long getSearchIndex() {
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
index 6f67bc70c86..2d4e7cd5e01 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
@@ -280,7 +280,7 @@ public class LogDispatcher {
 
     /** try to offer a request into queue with memory control. */
     public boolean offer(IndexedConsensusRequest indexedConsensusRequest) {
-      if 
(!iotConsensusMemoryManager.reserve(indexedConsensusRequest.getSerializedSize(),
 true)) {
+      if 
(!iotConsensusMemoryManager.reserve(indexedConsensusRequest.getMemorySize(), 
true)) {
         return false;
       }
       boolean success;
@@ -288,19 +288,19 @@ public class LogDispatcher {
         success = pendingEntries.offer(indexedConsensusRequest);
       } catch (Throwable t) {
         // If exception occurs during request offer, the reserved memory 
should be released
-        
iotConsensusMemoryManager.free(indexedConsensusRequest.getSerializedSize(), 
true);
+        
iotConsensusMemoryManager.free(indexedConsensusRequest.getMemorySize(), true);
         throw t;
       }
       if (!success) {
         // If offer failed, the reserved memory should be released
-        
iotConsensusMemoryManager.free(indexedConsensusRequest.getSerializedSize(), 
true);
+        
iotConsensusMemoryManager.free(indexedConsensusRequest.getMemorySize(), true);
       }
       return success;
     }
 
     /** try to remove a request from queue with memory control. */
     private void releaseReservedMemory(IndexedConsensusRequest 
indexedConsensusRequest) {
-      
iotConsensusMemoryManager.free(indexedConsensusRequest.getSerializedSize(), 
true);
+      iotConsensusMemoryManager.free(indexedConsensusRequest.getMemorySize(), 
true);
     }
 
     public void stop() {
@@ -322,13 +322,13 @@ public class LogDispatcher {
       }
       long requestSize = 0;
       for (IndexedConsensusRequest indexedConsensusRequest : pendingEntries) {
-        requestSize += indexedConsensusRequest.getSerializedSize();
+        requestSize += indexedConsensusRequest.getMemorySize();
       }
       pendingEntries.clear();
       iotConsensusMemoryManager.free(requestSize, true);
       requestSize = 0;
       for (IndexedConsensusRequest indexedConsensusRequest : bufferedEntries) {
-        requestSize += indexedConsensusRequest.getSerializedSize();
+        requestSize += indexedConsensusRequest.getMemorySize();
       }
       iotConsensusMemoryManager.free(requestSize, true);
       syncStatus.free();
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertNode.java
index 1f1c267075f..da1fc323684 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertNode.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertNode.java
@@ -443,6 +443,7 @@ public abstract class InsertNode extends SearchNode {
         .getPartialPath(ReadWriteIOUtils.readString(stream));
   }
 
+  @Override
   public long getMemorySize() {
     if (memorySize == 0) {
       memorySize = InsertNodeMemoryEstimator.sizeOf(this);

Reply via email to