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

Reply via email to