From 1db7dfc1361db1584f0d60069181808d08134cfd Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Sat, 25 Jul 2026 15:38:22 +0800 Subject: [PATCH] [Pipe] Preserve parser fairness and OOM retry state --- .../TsFileInsertionScanDataContainer.java | 5 ++ .../resource/memory/PipeMemoryManager.java | 46 +++++++++++++++-- .../TsFileInsertionDataContainerTest.java | 51 +++++++++++++++++++ .../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 6eee35e87c715..7322a29604f0e 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 @@ private Tablet getNextTablet() { } 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 9edecc257ea21..d5297f17b0966 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> waitingTsFileParserRegionOrderByPipe = new HashMap<>(); private final ArrayDeque waitingTsFileParserPipeOrder = new ArrayDeque<>(); + private PipeIdentity lastAdmittedWaitingTsFileParserPipe; // Only non-zero memory blocks will be added to this set. private final Set allocatedBlocks = new HashSet<>(); @@ -177,7 +178,8 @@ public synchronized boolean tryReserveTsFileParserMemory( 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 synchronized boolean tryReserveTsFileParserMemory( } 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 synchronized void releaseTsFileParserMemory( reservedTsFileParserCountByPipe.put(pipeIdentity, reservedCountOfPipe - 1); } reservedTsFileParserCount--; + clearTsFileParserAdmissionCursorIfIdle(); notifyNextTsFileParserMemoryReservationInternal(); } - private void enqueueTsFileParserReservationRequest( + private boolean enqueueTsFileParserReservationRequest( final PipeRegionIdentity pipeRegionIdentity, final TsFileParserMemoryReservation reservationKey) { final LinkedHashSet requestsOfPipeRegion = @@ -282,7 +288,7 @@ private void enqueueTsFileParserReservationRequest( regionOrder.addLast(key); return new LinkedHashSet<>(); }); - requestsOfPipeRegion.add(reservationKey); + return !requestsOfPipeRegion.add(reservationKey); } public synchronized void notifyNextTsFileParserMemoryReservation() { @@ -322,7 +328,14 @@ private void notifyNextTsFileParserMemoryReservationInternal() { 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 @@ private PipeRegionIdentity getNextEligibleTsFileParserPipeRegion( 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 @@ private void removeTsFileParserReservationRequest( if (regionOrder.isEmpty()) { waitingTsFileParserRegionOrderByPipe.remove(pipeIdentity); waitingTsFileParserPipeOrder.remove(pipeIdentity); + if (!rotateAfterAdmission) { + clearTsFileParserAdmissionCursorIfIdle(); + } return; } } @@ -376,6 +410,8 @@ private void removeTsFileParserReservationRequest( 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 19079f4cf764e..ba7d03d74e2b8 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 void testScanContainerReleasesTabletMemoryAfterRawTabletGenerated() throw } } + @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 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 @@ private PipeMemoryBlock getAllocatedTabletMemory(final TsFileInsertionDataContai 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 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 72bffbc2c610c..3a50002667e37 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 void testPipeFairnessIsNotWeightedByRegionCount() { 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);