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