This is an automated email from the ASF dual-hosted git repository.
jt2594838 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 06ab4f606ac [Pipe] Preserve parser fairness and OOM retry state
(#18306)
06ab4f606ac is described below
commit 06ab4f606ac82a7067f4cb493807242760cc0a65
Author: Caideyipi <[email protected]>
AuthorDate: Mon Jul 27 10:25:09 2026 +0800
[Pipe] Preserve parser fairness and OOM retry state (#18306)
---
.../scan/TsFileInsertionScanDataContainer.java | 5 +++
.../db/pipe/resource/memory/PipeMemoryManager.java | 46 ++++++++++++++++---
.../event/TsFileInsertionDataContainerTest.java | 51 ++++++++++++++++++++++
.../resource/memory/PipeMemoryManagerTest.java | 37 ++++++++++++++++
4 files changed, 134 insertions(+), 5 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/container/scan/TsFileInsertionScanDataContainer.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/container/scan/TsFileInsertionScanDataContainer.java
index 6eee35e87c7..7322a29604f 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/container/scan/TsFileInsertionScanDataContainer.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/container/scan/TsFileInsertionScanDataContainer.java
@@ -20,6 +20,7 @@
package org.apache.iotdb.db.pipe.event.common.tsfile.container.scan;
import org.apache.iotdb.commons.exception.IllegalPathException;
+import
org.apache.iotdb.commons.exception.pipe.PipeRuntimeOutOfMemoryCriticalException;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
import org.apache.iotdb.commons.pipe.config.PipeConfig;
import org.apache.iotdb.commons.pipe.datastructure.pattern.PipePattern;
@@ -352,6 +353,10 @@ public class TsFileInsertionScanDataContainer extends
TsFileInsertionDataContain
}
PipeTabletUtils.compactBitMaps(tablet);
return tablet;
+ } catch (final PipeRuntimeOutOfMemoryCriticalException e) {
+ // Keep the parser state so the caller can yield its parser slot and
retry from the same
+ // unconsumed data after memory is available again.
+ throw e;
} catch (final Exception e) {
close();
throw new PipeException("Failed to get next tablet insertion event.", e);
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 9edecc257ea..d5297f17b09 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
@@ -71,6 +71,7 @@ public class PipeMemoryManager {
private final Map<PipeIdentity, ArrayDeque<PipeRegionIdentity>>
waitingTsFileParserRegionOrderByPipe = new HashMap<>();
private final ArrayDeque<PipeIdentity> waitingTsFileParserPipeOrder = new
ArrayDeque<>();
+ private PipeIdentity lastAdmittedWaitingTsFileParserPipe;
// Only non-zero memory blocks will be added to this set.
private final Set<PipeMemoryBlock> allocatedBlocks = new HashSet<>();
@@ -177,7 +178,8 @@ public class PipeMemoryManager {
final PipeIdentity pipeIdentity = new PipeIdentity(pipeName, creationTime);
final PipeRegionIdentity pipeRegionIdentity =
new PipeRegionIdentity(pipeIdentity, dataRegionId);
- enqueueTsFileParserReservationRequest(pipeRegionIdentity, reservationKey);
+ final boolean wasRequestAlreadyWaiting =
+ enqueueTsFileParserReservationRequest(pipeRegionIdentity,
reservationKey);
final int globalLimit = Math.max(1,
PIPE_CONFIG.getPipeTsFileParserInFlightMaxNum());
final int perPipeRegionLimit =
@@ -212,6 +214,9 @@ public class PipeMemoryManager {
}
removeTsFileParserReservationRequest(pipeRegionIdentity, reservationKey,
true);
+ if (wasRequestAlreadyWaiting) {
+ lastAdmittedWaitingTsFileParserPipe = pipeIdentity;
+ }
reservedTsFileParserCount++;
reservedTsFileParserCountByPipe.merge(pipeIdentity, 1, Integer::sum);
reservedTsFileParserCountByPipeRegion.put(pipeRegionIdentity,
reservedCountOfPipeRegion + 1);
@@ -262,10 +267,11 @@ public class PipeMemoryManager {
reservedTsFileParserCountByPipe.put(pipeIdentity, reservedCountOfPipe -
1);
}
reservedTsFileParserCount--;
+ clearTsFileParserAdmissionCursorIfIdle();
notifyNextTsFileParserMemoryReservationInternal();
}
- private void enqueueTsFileParserReservationRequest(
+ private boolean enqueueTsFileParserReservationRequest(
final PipeRegionIdentity pipeRegionIdentity,
final TsFileParserMemoryReservation reservationKey) {
final LinkedHashSet<TsFileParserMemoryReservation> requestsOfPipeRegion =
@@ -282,7 +288,7 @@ public class PipeMemoryManager {
regionOrder.addLast(key);
return new LinkedHashSet<>();
});
- requestsOfPipeRegion.add(reservationKey);
+ return !requestsOfPipeRegion.add(reservationKey);
}
public synchronized void notifyNextTsFileParserMemoryReservation() {
@@ -322,7 +328,14 @@ public class PipeMemoryManager {
private PipeRegionIdentity getNextEligibleTsFileParserPipeRegion(
final int perPipeRegionLimit, final boolean
requirePipeWithoutReservedParser) {
+ PipeRegionIdentity firstEligiblePipeRegion = null;
+ boolean hasVisitedLastAdmittedPipe = lastAdmittedWaitingTsFileParserPipe
== null;
for (final PipeIdentity pipeIdentity : waitingTsFileParserPipeOrder) {
+ final boolean isLastAdmittedPipe =
pipeIdentity.equals(lastAdmittedWaitingTsFileParserPipe);
+ if (isLastAdmittedPipe) {
+ hasVisitedLastAdmittedPipe = true;
+ }
+
// Under soft memory pressure, reserve the hard-threshold headroom for a
pipe that has no
// parser yet. Otherwise a busy pipe at the queue head can block every
pipe behind it.
if (requirePipeWithoutReservedParser
@@ -335,14 +348,32 @@ public class PipeMemoryManager {
if (regionOrder == null) {
continue;
}
+ PipeRegionIdentity eligiblePipeRegion = null;
for (final PipeRegionIdentity pipeRegionIdentity : regionOrder) {
if
(reservedTsFileParserCountByPipeRegion.getOrDefault(pipeRegionIdentity, 0)
< perPipeRegionLimit) {
- return pipeRegionIdentity;
+ eligiblePipeRegion = pipeRegionIdentity;
+ break;
}
}
+ if (eligiblePipeRegion == null) {
+ continue;
+ }
+
+ if (firstEligiblePipeRegion == null) {
+ firstEligiblePipeRegion = eligiblePipeRegion;
+ }
+ if (hasVisitedLastAdmittedPipe && !isLastAdmittedPipe) {
+ return eligiblePipeRegion;
+ }
+ }
+ return firstEligiblePipeRegion;
+ }
+
+ private void clearTsFileParserAdmissionCursorIfIdle() {
+ if (reservedTsFileParserCount == 0 &&
waitingTsFileParserPipeOrder.isEmpty()) {
+ lastAdmittedWaitingTsFileParserPipe = null;
}
- return null;
}
private void removeTsFileParserReservationRequest(
@@ -365,6 +396,9 @@ public class PipeMemoryManager {
if (regionOrder.isEmpty()) {
waitingTsFileParserRegionOrderByPipe.remove(pipeIdentity);
waitingTsFileParserPipeOrder.remove(pipeIdentity);
+ if (!rotateAfterAdmission) {
+ clearTsFileParserAdmissionCursorIfIdle();
+ }
return;
}
}
@@ -376,6 +410,8 @@ public class PipeMemoryManager {
if (rotateAfterAdmission) {
waitingTsFileParserPipeOrder.remove(pipeIdentity);
waitingTsFileParserPipeOrder.addLast(pipeIdentity);
+ } else {
+ clearTsFileParserAdmissionCursorIfIdle();
}
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/TsFileInsertionDataContainerTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/TsFileInsertionDataContainerTest.java
index 19079f4cf76..ba7d03d74e2 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/TsFileInsertionDataContainerTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/TsFileInsertionDataContainerTest.java
@@ -184,6 +184,47 @@ public class TsFileInsertionDataContainerTest {
}
}
+ @Test
+ public void testScanContainerKeepsIteratorOnOutOfMemory() throws Exception {
+ nonalignedTsFile =
+ TsFileGeneratorUtils.generateNonAlignedTsFile(
+ "nonaligned-retry-tablet-memory.tsfile", 1, 1, 10, 0, 100, 10, 10);
+
+ try (final TsFileInsertionScanDataContainer container =
+ new TsFileInsertionScanDataContainer(
+ nonalignedTsFile,
+ new PrefixPipePattern("root"),
+ Long.MIN_VALUE,
+ Long.MAX_VALUE,
+ null,
+ null,
+ false)) {
+ final AtomicInteger memoryUsageReadCount = new AtomicInteger(0);
+ replaceAllocatedTabletMemory(
+ container,
+ new PipeMemoryBlock(0) {
+ @Override
+ public long getMemoryUsageInBytes() {
+ if (memoryUsageReadCount.incrementAndGet() == 2) {
+ throw new PipeRuntimeOutOfMemoryCriticalException("expected
oom");
+ }
+ return super.getMemoryUsageInBytes();
+ }
+ });
+
+ final Iterator<TabletInsertionEvent> iterator =
+ container.toTabletInsertionEvents().iterator();
+ final PipeRuntimeOutOfMemoryCriticalException exception =
+ Assert.assertThrows(PipeRuntimeOutOfMemoryCriticalException.class,
iterator::next);
+ Assert.assertEquals("expected oom", exception.getMessage());
+
+ Assert.assertTrue(iterator.hasNext());
+ final TabletInsertionEvent event = iterator.next();
+ Assert.assertTrue(event instanceof PipeRawTabletInsertionEvent);
+ ((PipeRawTabletInsertionEvent)
event).clearReferenceCount(getClass().getName());
+ }
+ }
+
@Test
public void
testConsumeTabletInsertionEventsWithRetryPreservesProgressOnOutOfMemory()
throws Exception {
@@ -1309,6 +1350,16 @@ public class TsFileInsertionDataContainerTest {
return (PipeMemoryBlock) field.get(container);
}
+ private void replaceAllocatedTabletMemory(
+ final TsFileInsertionDataContainer container, final PipeMemoryBlock
replacement)
+ throws NoSuchFieldException, IllegalAccessException {
+ final Field field =
+
TsFileInsertionDataContainer.class.getDeclaredField("allocatedMemoryBlockForTablet");
+ field.setAccessible(true);
+ ((PipeMemoryBlock) field.get(container)).close();
+ field.set(container, replacement);
+ }
+
@SuppressWarnings("unchecked")
private AtomicReference<TsFileInsertionDataContainer> getDataContainer(
final PipeTsFileInsertionEvent event) throws NoSuchFieldException,
IllegalAccessException {
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManagerTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManagerTest.java
index 72bffbc2c61..3a50002667e 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManagerTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManagerTest.java
@@ -192,6 +192,43 @@ public class PipeMemoryManagerTest {
Assert.assertTrue(tryAcquire(pipeARegion2));
}
+ @Test
+ public void testPipeFairnessSurvivesTransientSingleRegionQueueGap() {
+ commonConfig.setPipeTsFileParserInFlightMaxNum(1);
+ commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(1);
+
+ final Reservation blocker = new Reservation("blocker", 0);
+ final Reservation multiRegion1 = new Reservation("multi", 1, "1");
+ final Reservation multiRegion2 = new Reservation("multi", 1, "2");
+ final Reservation multiRegion3 = new Reservation("multi", 1, "3");
+ final Reservation singleFirst = new Reservation("single", 2, "1");
+ final Reservation singleSecond = new Reservation("single", 2, "1");
+
+ Assert.assertTrue(tryAcquire(blocker));
+ Assert.assertFalse(tryAcquire(multiRegion1));
+ Assert.assertFalse(tryAcquire(multiRegion2));
+ Assert.assertFalse(tryAcquire(multiRegion3));
+ Assert.assertFalse(tryAcquire(singleFirst));
+
+ release(blocker);
+ Assert.assertTrue(tryAcquire(multiRegion1));
+ release(multiRegion1);
+
+ Assert.assertTrue(tryAcquire(singleFirst));
+ release(singleFirst);
+
+ // The single-region pipe temporarily has no admission request while it
advances to its next
+ // TsFile, so the multi-region pipe can use the otherwise idle parser slot.
+ Assert.assertTrue(tryAcquire(multiRegion2));
+ Assert.assertFalse(tryAcquire(singleSecond));
+ release(multiRegion2);
+
+ // Once the single-region pipe is waiting again, the remembered pipe-level
cursor must prevent
+ // another region of the multi-region pipe from taking a second
consecutive turn.
+ Assert.assertFalse(tryAcquire(multiRegion3));
+ Assert.assertTrue(tryAcquire(singleSecond));
+ }
+
@Test
public void testSoftMemoryHeadroomIsReservedForPipeWithoutParser() {
commonConfig.setPipeTsFileParserInFlightMaxNum(2);