This is an automated email from the ASF dual-hosted git repository.
rong pushed a commit to branch dev/1.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/dev/1.3 by this push:
new f069e60892c Load: fix memory leak when failed in 2nd phase (#15503)
(#15513)
f069e60892c is described below
commit f069e60892ce4441ba513425d32b5f66bd881530
Author: Zikun Ma <[email protected]>
AuthorDate: Thu May 15 20:01:56 2025 +0800
Load: fix memory leak when failed in 2nd phase (#15503) (#15513)
(cherry picked from commit 06ccdd163d0c7579d4ed2015686de6ada91c7681)
---
.../receiver/protocol/thrift/IoTDBDataNodeReceiver.java | 6 ++++--
.../plan/scheduler/load/LoadTsFileScheduler.java | 14 +++++++++-----
2 files changed, 13 insertions(+), 7 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
index d7161ca4669..975e8beb77d 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
@@ -680,11 +680,13 @@ public class IoTDBDataNodeReceiver extends
IoTDBFileReceiver {
} catch (final PipeRuntimeOutOfMemoryCriticalException e) {
final String message =
String.format(
- "Temporarily out of memory when executing statement %s,
Requested memory: %s, used memory: %s, total memory: %s",
+ "Temporarily out of memory when executing statement %s,
Requested memory: %s, "
+ + "used memory: %s, free memory: %s, total non-floating
memory: %s",
statement,
estimatedMemory,
PipeDataNodeResourceManager.memory().getUsedMemorySizeInBytes(),
- PipeDataNodeResourceManager.memory().getFreeMemorySizeInBytes());
+ PipeDataNodeResourceManager.memory().getFreeMemorySizeInBytes(),
+
PipeDataNodeResourceManager.memory().getTotalNonFloatingMemorySizeInBytes());
if (LOGGER.isDebugEnabled()) {
LOGGER.debug("Receiver id = {}: {}", receiverId.get(), message, e);
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java
index a86f22ac2cd..f45e9805200 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java
@@ -414,7 +414,9 @@ public class LoadTsFileScheduler implements IScheduler {
result.getFailureStatus().getMessage());
TSStatus status = result.getFailureStatus();
status.setMessage(
- String.format("Load %s error in 2nd phase. Because ", tsFile) +
status.getMessage());
+ String.format(
+ "Load %s error in second phase. Because %s, first phase is %s",
+ tsFile, status.getMessage(), isFirstPhaseSuccess ? "success" :
"failed"));
stateMachine.transitionToFailed(status);
return false;
}
@@ -740,19 +742,21 @@ public class LoadTsFileScheduler implements IScheduler {
private boolean sendAllTsFileData() throws LoadFileException {
routeChunkData();
+ boolean isAllSuccess = true;
for (Map.Entry<TConsensusGroupId, Pair<TRegionReplicaSet,
LoadTsFilePieceNode>> entry :
regionId2ReplicaSetAndNode.entrySet()) {
block.reduceMemoryUsage(entry.getValue().getRight().getDataSize());
- if (!scheduler.dispatchOnePieceNode(
- entry.getValue().getRight(), entry.getValue().getLeft())) {
+ if (isAllSuccess
+ && !scheduler.dispatchOnePieceNode(
+ entry.getValue().getRight(), entry.getValue().getLeft())) {
LOGGER.warn(
"Dispatch piece node {} of TsFile {} error.",
entry.getValue(),
singleTsFileNode.getTsFileResource().getTsFile());
- return false;
+ isAllSuccess = false;
}
}
- return true;
+ return isAllSuccess;
}
private void clear() {