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());
+ }
+}