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

jt2594838 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 01d1c15aaf7 Fix pipe tablet memory self-lock during batching (#18266)
01d1c15aaf7 is described below

commit 01d1c15aaf76903a358e863fed361693a1e7e56b
Author: Caideyipi <[email protected]>
AuthorDate: Wed Jul 22 16:16:58 2026 +0800

    Fix pipe tablet memory self-lock during batching (#18266)
---
 .../db/pipe/resource/memory/PipeMemoryManager.java |  33 ++++--
 .../memory/PipeMemoryManagerResizeTest.java        | 116 +++++++++++++++++++++
 2 files changed, 142 insertions(+), 7 deletions(-)

diff --git 
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
index 98b820925c7..acd85b98378 100644
--- 
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
@@ -48,11 +48,7 @@ public class PipeMemoryManager {
       PipeConfig.getInstance().getPipeMemoryManagementEnabled();
 
   // TODO @spricoder: consider combine memory block and used MemorySizeInBytes
-  private final IMemoryBlock memoryBlock =
-      IoTDBDescriptor.getInstance()
-          .getMemoryConfig()
-          .getPipeMemoryManager()
-          .exactAllocate("Stream", MemoryBlockType.DYNAMIC);
+  private final IMemoryBlock memoryBlock;
 
   private static final double EXCEED_PROTECT_THRESHOLD = 0.95;
 
@@ -68,6 +64,11 @@ public class PipeMemoryManager {
   private final Set<PipeMemoryBlock> expandableBlocks = new HashSet<>();
 
   public PipeMemoryManager() {
+    this(
+        IoTDBDescriptor.getInstance()
+            .getMemoryConfig()
+            .getPipeMemoryManager()
+            .exactAllocate("Stream", MemoryBlockType.DYNAMIC));
     PipeDataNodeAgent.runtime()
         .registerPeriodicalJob(
             "PipeMemoryManager#tryExpandAll()",
@@ -75,6 +76,10 @@ public class PipeMemoryManager {
             PipeConfig.getInstance().getPipeMemoryExpanderIntervalSeconds());
   }
 
+  PipeMemoryManager(final IMemoryBlock memoryBlock) {
+    this.memoryBlock = memoryBlock;
+  }
+
   // NOTE: Here we unify the memory threshold judgment for tablet and tsfile 
memory block, because
   // introducing too many heuristic rules not conducive to flexible dynamic 
adjustment of memory
   // configuration:
@@ -210,6 +215,16 @@ public class PipeMemoryManager {
         && (double) usedMemorySizeInBytesOfTsFiles < 
allowedMaxMemorySizeInBytesOfTsTiles();
   }
 
+  private boolean isHardEnoughForResizing(final PipeMemoryBlock block) {
+    if (block instanceof PipeTabletMemoryBlock) {
+      return isHardEnough4TabletParsing();
+    }
+    if (block instanceof PipeTsFileMemoryBlock) {
+      return isHardEnough4TsFileSlicing();
+    }
+    return true;
+  }
+
   public synchronized PipeMemoryBlock forceAllocate(long sizeInBytes)
       throws PipeRuntimeOutOfMemoryCriticalException {
     if (!PIPE_MEMORY_MANAGEMENT_ENABLED) {
@@ -434,8 +449,12 @@ public class PipeMemoryManager {
     long sizeInBytes = targetSize - oldSize;
     final int memoryAllocateMaxRetries = 
PIPE_CONFIG.getPipeMemoryAllocateMaxRetries();
     for (int i = 1; i <= memoryAllocateMaxRetries; i++) {
-      if (getTotalNonFloatingMemorySizeInBytes() - 
memoryBlock.getUsedMemoryInBytes()
-          >= sizeInBytes) {
+      // Dynamically resized data-structure blocks must obey the same 
admission thresholds as
+      // blocks allocated with a non-zero initial size. Otherwise they can 
exhaust the pool and
+      // prevent downstream consumers from allocating the memory needed to 
release them.
+      if (isHardEnoughForResizing(block)
+          && getTotalNonFloatingMemorySizeInBytes() - 
memoryBlock.getUsedMemoryInBytes()
+              >= sizeInBytes) {
         memoryBlock.forceAllocateWithoutLimitation(sizeInBytes);
         if (oldSize == 0) {
           // If the memory block is not registered, we need to register it 
first.
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManagerResizeTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManagerResizeTest.java
new file mode 100644
index 00000000000..6c320e973dd
--- /dev/null
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManagerResizeTest.java
@@ -0,0 +1,116 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.pipe.resource.memory;
+
+import org.apache.iotdb.commons.conf.CommonConfig;
+import org.apache.iotdb.commons.conf.CommonDescriptor;
+import 
org.apache.iotdb.commons.exception.pipe.PipeRuntimeOutOfMemoryCriticalException;
+import org.apache.iotdb.commons.memory.AtomicLongMemoryBlock;
+import org.apache.iotdb.commons.memory.MemoryBlockType;
+
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+
+public class PipeMemoryManagerResizeTest {
+
+  private static final long TOTAL_MEMORY_SIZE_IN_BYTES = 2000;
+  private static final long TABLET_MEMORY_SIZE_IN_BYTES = 451;
+  private static final long SINK_MEMORY_SIZE_IN_BYTES = 100;
+
+  private final CommonConfig config = 
CommonDescriptor.getInstance().getConfig();
+
+  private boolean originalMemoryManagementEnabled;
+  private int originalAllocateMaxRetries;
+  private long originalAllocateRetryIntervalInMs;
+  private double originalFloatingMemoryProportion;
+  private double originalTabletRejectThreshold;
+  private double originalTsFileRejectThreshold;
+
+  @Before
+  public void setUp() {
+    originalMemoryManagementEnabled = config.getPipeMemoryManagementEnabled();
+    originalAllocateMaxRetries = config.getPipeMemoryAllocateMaxRetries();
+    originalAllocateRetryIntervalInMs = 
config.getPipeMemoryAllocateRetryIntervalInMs();
+    originalFloatingMemoryProportion = 
config.getPipeTotalFloatingMemoryProportion();
+    originalTabletRejectThreshold =
+        
config.getPipeDataStructureTabletMemoryBlockAllocationRejectThreshold();
+    originalTsFileRejectThreshold =
+        
config.getPipeDataStructureTsFileMemoryBlockAllocationRejectThreshold();
+
+    config.setPipeMemoryManagementEnabled(true);
+    config.setPipeMemoryAllocateMaxRetries(1);
+    config.setPipeMemoryAllocateRetryIntervalInMs(1);
+    config.setPipeTotalFloatingMemoryProportion(0.5);
+    config.setPipeDataStructureTabletMemoryBlockAllocationRejectThreshold(0.3);
+    config.setPipeDataStructureTsFileMemoryBlockAllocationRejectThreshold(0.3);
+  }
+
+  @After
+  public void tearDown() {
+    config.setPipeMemoryManagementEnabled(originalMemoryManagementEnabled);
+    config.setPipeMemoryAllocateMaxRetries(originalAllocateMaxRetries);
+    
config.setPipeMemoryAllocateRetryIntervalInMs(originalAllocateRetryIntervalInMs);
+    
config.setPipeTotalFloatingMemoryProportion(originalFloatingMemoryProportion);
+    config.setPipeDataStructureTabletMemoryBlockAllocationRejectThreshold(
+        originalTabletRejectThreshold);
+    config.setPipeDataStructureTsFileMemoryBlockAllocationRejectThreshold(
+        originalTsFileRejectThreshold);
+  }
+
+  @Test
+  public void testTabletResizeLeavesMemoryForSinkForwardProgress() {
+    final PipeMemoryManager manager =
+        new PipeMemoryManager(
+            new AtomicLongMemoryBlock(
+                "PipeMemoryManagerResizeTest",
+                null,
+                TOTAL_MEMORY_SIZE_IN_BYTES,
+                MemoryBlockType.DYNAMIC));
+    final PipeTabletMemoryBlock retainedTablet =
+        manager.forceAllocateForTabletWithRetry(TABLET_MEMORY_SIZE_IN_BYTES);
+    final PipeTabletMemoryBlock pendingTablet = 
manager.forceAllocateForTabletWithRetry(0);
+    final PipeMemoryBlock sinkBatch = manager.forceAllocate(0);
+
+    try {
+      Assert.assertThrows(
+          PipeRuntimeOutOfMemoryCriticalException.class,
+          () -> manager.forceResize(pendingTablet, 1));
+      Assert.assertEquals(TABLET_MEMORY_SIZE_IN_BYTES, 
manager.getUsedMemorySizeInBytes());
+      Assert.assertEquals(TABLET_MEMORY_SIZE_IN_BYTES, 
manager.getUsedMemorySizeInBytesOfTablets());
+
+      manager.forceResize(sinkBatch, SINK_MEMORY_SIZE_IN_BYTES);
+      Assert.assertEquals(
+          TABLET_MEMORY_SIZE_IN_BYTES + SINK_MEMORY_SIZE_IN_BYTES,
+          manager.getUsedMemorySizeInBytes());
+
+      manager.release(retainedTablet);
+      manager.forceResize(pendingTablet, 1);
+      Assert.assertEquals(1, manager.getUsedMemorySizeInBytesOfTablets());
+    } finally {
+      manager.release(retainedTablet);
+      manager.release(pendingTablet);
+      manager.release(sinkBatch);
+    }
+
+    Assert.assertEquals(0, manager.getUsedMemorySizeInBytes());
+  }
+}

Reply via email to