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);

Reply via email to