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

Caideyipi 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 2cfccada268 [Pipe] Fairly limit concurrent TsFile parsers (#18260) 
(#18295)
2cfccada268 is described below

commit 2cfccada2687aa8dc9271fa86ca56fc9c35e7371
Author: Caideyipi <[email protected]>
AuthorDate: Sat Jul 25 16:26:39 2026 +0800

    [Pipe] Fairly limit concurrent TsFile parsers (#18260) (#18295)
---
 .../org/apache/iotdb/db/conf/IoTDBDescriptor.java  |   2 +
 .../common/tsfile/PipeTsFileInsertionEvent.java    | 124 ++++----
 .../db/pipe/resource/memory/PipeMemoryManager.java | 316 ++++++++++++++++++-
 .../org/apache/iotdb/db/conf/PropertiesTest.java   |  29 ++
 .../PipeTsFileInsertionEventAdmissionTest.java     |  65 ++--
 .../resource/memory/PipeMemoryManagerTest.java     | 347 +++++++++++++++++++++
 .../conf/iotdb-system.properties.template          |  13 +
 .../apache/iotdb/commons/conf/CommonConfig.java    |  37 +++
 .../iotdb/commons/pipe/config/PipeConfig.java      |  12 +
 .../iotdb/commons/pipe/config/PipeDescriptor.java  |  10 +
 10 files changed, 838 insertions(+), 117 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
index b38a7c186ec..5658e10cc60 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
@@ -35,6 +35,7 @@ import org.apache.iotdb.confignode.rpc.thrift.TGlobalConfig;
 import org.apache.iotdb.confignode.rpc.thrift.TRatisConfig;
 import org.apache.iotdb.db.consensus.DataRegionConsensusImpl;
 import org.apache.iotdb.db.exception.query.QueryProcessException;
+import org.apache.iotdb.db.pipe.resource.PipeDataNodeResourceManager;
 import 
org.apache.iotdb.db.queryengine.plan.relational.metadata.fetcher.cache.LastCacheLoadStrategy;
 import org.apache.iotdb.db.service.metrics.IoTDBInternalLocalReporter;
 import org.apache.iotdb.db.storageengine.StorageEngine;
@@ -2647,6 +2648,7 @@ public class IoTDBDescriptor {
 
   private void loadPipeHotModifiedProp(TrimProperties properties) throws 
IOException {
     PipeDescriptor.loadPipeProps(commonDescriptor.getConfig(), properties, 
true);
+    
PipeDataNodeResourceManager.memory().notifyNextTsFileParserMemoryReservation();
     LoggerPeriodicalLogReducer.update();
   }
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
index ece6e25a151..a4859c33579 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
@@ -35,6 +35,7 @@ import 
org.apache.iotdb.db.pipe.event.common.tsfile.container.TsFileInsertionDat
 import org.apache.iotdb.db.pipe.metric.overview.PipeDataNodeSinglePipeMetrics;
 import org.apache.iotdb.db.pipe.resource.PipeDataNodeResourceManager;
 import org.apache.iotdb.db.pipe.resource.memory.PipeMemoryManager;
+import 
org.apache.iotdb.db.pipe.resource.memory.PipeMemoryManager.TsFileParserMemoryReservation;
 import org.apache.iotdb.db.pipe.resource.tsfile.PipeTsFileResourceManager;
 import 
org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner.PipeTsFileEpochProgressIndexKeeper;
 import org.apache.iotdb.db.storageengine.dataregion.memtable.TsFileProcessor;
@@ -56,7 +57,6 @@ import java.util.Iterator;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Set;
-import java.util.concurrent.ThreadLocalRandom;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicInteger;
 import java.util.concurrent.atomic.AtomicLong;
@@ -68,6 +68,7 @@ public class PipeTsFileInsertionEvent extends EnrichedEvent
   private static final Logger LOGGER = 
LoggerFactory.getLogger(PipeTsFileInsertionEvent.class);
 
   protected final TsFileResource resource;
+  private final String dataRegionId;
   protected File tsFile;
   protected long extractTime = 0;
 
@@ -87,6 +88,8 @@ public class PipeTsFileInsertionEvent extends EnrichedEvent
   protected final AtomicBoolean isClosed;
   protected final AtomicReference<TsFileInsertionDataContainer> dataContainer;
   private final AtomicBoolean isTsFileParserMemoryReserved = new 
AtomicBoolean(false);
+  private final TsFileParserMemoryReservation tsFileParserMemoryReservationKey 
=
+      new TsFileParserMemoryReservation();
   private final AtomicReference<Iterator<TabletInsertionEvent>> 
tabletInsertionEventIterator =
       new AtomicReference<>();
   private final AtomicReference<PipeRawTabletInsertionEvent> 
pendingTabletInsertionEvent =
@@ -162,6 +165,7 @@ public class PipeTsFileInsertionEvent extends EnrichedEvent
     super(pipeName, creationTime, pipeTaskMeta, pipePattern, startTime, 
endTime);
 
     this.resource = resource;
+    this.dataRegionId = getDataRegionId(resource);
 
     // For events created at assigner or historical extractor, the tsFile is 
get from the resource
     // For events created for source, the tsFile is inherited from the 
assigner, because the
@@ -226,6 +230,17 @@ public class PipeTsFileInsertionEvent extends EnrichedEvent
         });
   }
 
+  private static String getDataRegionId(final TsFileResource resource) {
+    // TsFileResource#getDataRegionId assumes the storage-engine directory 
structure, while a
+    // synthetic resource may wrap a standalone file.
+    final File resourceTsFile = resource.getTsFile();
+    final File timePartitionDir =
+        Objects.isNull(resourceTsFile) ? null : resourceTsFile.getParentFile();
+    final File dataRegionDir =
+        Objects.isNull(timePartitionDir) ? null : 
timePartitionDir.getParentFile();
+    return Objects.isNull(dataRegionDir) ? "" : dataRegionDir.getName();
+  }
+
   /**
    * @return {@code false} if this file can't be sent by pipe because it is 
empty. {@code true}
    *     otherwise.
@@ -677,31 +692,12 @@ public class PipeTsFileInsertionEvent extends 
EnrichedEvent
     final long startTime = System.currentTimeMillis();
     long lastRecordTime = startTime;
 
-    final long initialMemoryCheckIntervalMs =
-        Math.max(1, 
PipeConfig.getInstance().getPipeCheckMemoryEnoughIntervalMs());
-    final long maxMemoryCheckIntervalMs =
-        getMaxMemoryCheckIntervalMs(
-            initialMemoryCheckIntervalMs,
-            PipeConfig.getInstance().getPipeMemoryAllocateMaxRetries());
-    long memoryCheckIntervalMs = initialMemoryCheckIntervalMs;
-    while (true) {
-      final long elapsedTimeMs = Math.max(0, System.currentTimeMillis() - 
startTime);
-      if (elapsedTimeMs >= timeoutMs) {
-        // should contain 'TimeoutException' in exception message
-        throw new PipeRuntimeOutOfMemoryCriticalException(
-            String.format(
-                "TimeoutException: Waited %s seconds for memory to parse 
TsFile",
-                elapsedTimeMs / 1000.0));
-      }
-
-      memoryManager.waitForTsFileParserMemory(
-          Math.min(
-              getMemoryCheckIntervalWithJitter(memoryCheckIntervalMs), 
timeoutMs - elapsedTimeMs));
-
+    while (!tryReserveTsFileParserMemory(memoryManager)) {
       final long currentTime = System.currentTimeMillis();
-      final double elapsedRecordTimeSeconds = (currentTime - lastRecordTime) / 
1000.0;
-      final double waitTimeSeconds = (currentTime - startTime) / 1000.0;
-      if (elapsedRecordTimeSeconds > 10.0) {
+      final long elapsedRecordTimeInMs = currentTime - lastRecordTime;
+      final long waitTimeInMs = currentTime - startTime;
+      final double waitTimeSeconds = waitTimeInMs / 1000.0;
+      if (elapsedRecordTimeInMs > 10_000) {
         LOGGER.info(
             "Wait for memory enough for parsing {} for {} seconds.",
             resource != null ? resource.getTsFilePath() : "tsfile",
@@ -714,35 +710,26 @@ public class PipeTsFileInsertionEvent extends 
EnrichedEvent
             waitTimeSeconds);
       }
 
-      if (tryReserveTsFileParserMemory(memoryManager)) {
-        LOGGER.info(
-            "Wait for memory enough for parsing {} for {} seconds.",
-            resource != null ? resource.getTsFilePath() : "tsfile",
-            waitTimeSeconds);
-        return;
+      if (waitTimeInMs > timeoutMs) {
+        // should contain 'TimeoutException' in exception message
+        throw new PipeRuntimeOutOfMemoryCriticalException(
+            String.format(
+                "TimeoutException: Waited %s seconds for memory to parse 
TsFile", waitTimeSeconds));
       }
 
-      memoryCheckIntervalMs =
-          getNextMemoryCheckIntervalMs(memoryCheckIntervalMs, 
maxMemoryCheckIntervalMs);
+      tsFileParserMemoryReservationKey.await(
+          Math.max(
+              1,
+              Math.min(
+                  timeoutMs - waitTimeInMs, 10_000 - Math.min(10_000, 
elapsedRecordTimeInMs))));
     }
-  }
-
-  static long getMaxMemoryCheckIntervalMs(final long initialIntervalMs, final 
int maxRetries) {
-    final long multiplier = Math.max(1, maxRetries);
-    return initialIntervalMs > Long.MAX_VALUE / multiplier
-        ? Long.MAX_VALUE
-        : initialIntervalMs * multiplier;
-  }
-
-  static long getNextMemoryCheckIntervalMs(final long currentIntervalMs, final 
long maxIntervalMs) {
-    return currentIntervalMs >= maxIntervalMs - currentIntervalMs
-        ? maxIntervalMs
-        : currentIntervalMs << 1;
-  }
 
-  static long getMemoryCheckIntervalWithJitter(final long intervalMs) {
-    return Math.max(
-        1, (long) (intervalMs * (0.5 + 
ThreadLocalRandom.current().nextDouble() * 0.5)));
+    final long currentTime = System.currentTimeMillis();
+    final double waitTimeSeconds = (currentTime - startTime) / 1000.0;
+    LOGGER.info(
+        "Wait for memory enough for parsing {} for {} seconds.",
+        resource != null ? resource.getTsFilePath() : "tsfile",
+        waitTimeSeconds);
   }
 
   private boolean tryReserveTsFileParserMemory(final PipeMemoryManager 
memoryManager) {
@@ -751,7 +738,8 @@ public class PipeTsFileInsertionEvent extends EnrichedEvent
         return true;
       }
 
-      if (!memoryManager.tryReserveTsFileParserMemory()) {
+      if (!memoryManager.tryReserveTsFileParserMemory(
+          pipeName, creationTime, dataRegionId, 
tsFileParserMemoryReservationKey)) {
         return false;
       }
 
@@ -763,11 +751,20 @@ public class PipeTsFileInsertionEvent extends 
EnrichedEvent
   private void releaseTsFileParserMemoryIfReserved() {
     synchronized (isTsFileParserMemoryReserved) {
       if (isTsFileParserMemoryReserved.compareAndSet(true, false)) {
-        PipeDataNodeResourceManager.memory().releaseTsFileParserMemory();
+        PipeDataNodeResourceManager.memory()
+            .releaseTsFileParserMemory(pipeName, creationTime, dataRegionId);
       }
     }
   }
 
+  private void cancelTsFileParserMemoryReservationIfPending() {
+    if (!isTsFileParserMemoryReserved.get()) {
+      PipeDataNodeResourceManager.memory()
+          .cancelTsFileParserMemoryReservation(
+              pipeName, creationTime, dataRegionId, 
tsFileParserMemoryReservationKey);
+    }
+  }
+
   /** The method is used to prevent circular replication in PipeConsensus */
   public boolean isGeneratedByPipeConsensus() {
     return isGeneratedByPipeConsensus;
@@ -829,6 +826,7 @@ public class PipeTsFileInsertionEvent extends EnrichedEvent
   /** Release the resource of {@link TsFileInsertionDataContainer}. */
   @Override
   public void close() {
+    cancelTsFileParserMemoryReservationIfPending();
     tabletInsertionEventIterator.set(null);
     releaseParsedTabletEvent(pendingTabletInsertionEvent.getAndSet(null));
     parsedTabletInsertionEventCount.set(0);
@@ -877,11 +875,14 @@ public class PipeTsFileInsertionEvent extends 
EnrichedEvent
         this.isReleased,
         this.referenceCount,
         this.pipeName,
+        this.creationTime,
+        this.dataRegionId,
         this.tsFile,
         this.isWithMod,
         this.modFile,
         this.dataContainer,
-        this.isTsFileParserMemoryReserved);
+        this.isTsFileParserMemoryReserved,
+        this.tsFileParserMemoryReservationKey);
   }
 
   private static class PipeTsFileInsertionEventResource extends 
PipeEventResource {
@@ -891,29 +892,41 @@ public class PipeTsFileInsertionEvent extends 
EnrichedEvent
     private final File modFile;
     private final AtomicReference<TsFileInsertionDataContainer> dataContainer;
     private final String pipeName;
+    private final long creationTime;
+    private final String dataRegionId;
     private final AtomicBoolean isTsFileParserMemoryReserved;
+    private final TsFileParserMemoryReservation 
tsFileParserMemoryReservationKey;
 
     private PipeTsFileInsertionEventResource(
         final AtomicBoolean isReleased,
         final AtomicInteger referenceCount,
         final String pipeName,
+        final long creationTime,
+        final String dataRegionId,
         final File tsFile,
         final boolean isWithMod,
         final File modFile,
         final AtomicReference<TsFileInsertionDataContainer> dataContainer,
-        final AtomicBoolean isTsFileParserMemoryReserved) {
+        final AtomicBoolean isTsFileParserMemoryReserved,
+        final TsFileParserMemoryReservation tsFileParserMemoryReservationKey) {
       super(isReleased, referenceCount);
       this.pipeName = pipeName;
+      this.creationTime = creationTime;
+      this.dataRegionId = dataRegionId;
       this.tsFile = tsFile;
       this.isWithMod = isWithMod;
       this.modFile = modFile;
       this.dataContainer = dataContainer;
       this.isTsFileParserMemoryReserved = isTsFileParserMemoryReserved;
+      this.tsFileParserMemoryReservationKey = tsFileParserMemoryReservationKey;
     }
 
     @Override
     protected void finalizeResource() {
       try {
+        PipeDataNodeResourceManager.memory()
+            .cancelTsFileParserMemoryReservation(
+                pipeName, creationTime, dataRegionId, 
tsFileParserMemoryReservationKey);
         // decrease reference count
         PipeDataNodeResourceManager.tsfile().decreaseFileReference(tsFile, 
pipeName);
         if (isWithMod) {
@@ -930,7 +943,8 @@ public class PipeTsFileInsertionEvent extends EnrichedEvent
             });
         synchronized (isTsFileParserMemoryReserved) {
           if (isTsFileParserMemoryReserved.compareAndSet(true, false)) {
-            PipeDataNodeResourceManager.memory().releaseTsFileParserMemory();
+            PipeDataNodeResourceManager.memory()
+                .releaseTsFileParserMemory(pipeName, creationTime, 
dataRegionId);
           }
         }
       } catch (final Exception 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 ab4ac53baa8..9edecc257ea 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
@@ -28,10 +28,15 @@ import 
org.apache.iotdb.db.pipe.resource.memory.strategy.ThresholdAllocationStra
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.util.ArrayDeque;
 import java.util.ArrayList;
 import java.util.Collections;
+import java.util.HashMap;
 import java.util.HashSet;
+import java.util.LinkedHashSet;
 import java.util.List;
+import java.util.Map;
+import java.util.Objects;
 import java.util.Set;
 import java.util.function.LongUnaryOperator;
 
@@ -39,6 +44,8 @@ public class PipeMemoryManager {
 
   private static final Logger LOGGER = 
LoggerFactory.getLogger(PipeMemoryManager.class);
 
+  private static final PipeConfig PIPE_CONFIG = PipeConfig.getInstance();
+
   private static final boolean PIPE_MEMORY_MANAGEMENT_ENABLED =
       PipeConfig.getInstance().getPipeMemoryManagementEnabled();
   private static final long TOTAL_MEMORY_SIZE_IN_BYTES =
@@ -56,6 +63,15 @@ public class PipeMemoryManager {
 
   private volatile long reservedTsFileParserCount;
 
+  private final Map<PipeIdentity, Integer> reservedTsFileParserCountByPipe = 
new HashMap<>();
+  private final Map<PipeRegionIdentity, Integer> 
reservedTsFileParserCountByPipeRegion =
+      new HashMap<>();
+  private final Map<PipeRegionIdentity, 
LinkedHashSet<TsFileParserMemoryReservation>>
+      waitingTsFileParserRequestsByPipeRegion = new HashMap<>();
+  private final Map<PipeIdentity, ArrayDeque<PipeRegionIdentity>>
+      waitingTsFileParserRegionOrderByPipe = new HashMap<>();
+  private final ArrayDeque<PipeIdentity> waitingTsFileParserPipeOrder = new 
ArrayDeque<>();
+
   // Only non-zero memory blocks will be added to this set.
   private final Set<PipeMemoryBlock> allocatedBlocks = new HashSet<>();
   private final Set<PipeMemoryBlock> shrinkableBlocks = new HashSet<>();
@@ -125,9 +141,12 @@ public class PipeMemoryManager {
             < EXCEED_PROTECT_THRESHOLD * 
allowedMaxMemorySizeInBytesOfTablets();
   }
 
-  private boolean isHardEnough4TabletParsingWithReservedParserMemory() {
+  private boolean isHardEnough4TabletParsingWithReservedParserMemory(
+      final long extraMemoryInBytes) {
     final double tabletMemoryWithParserMemory =
-        (double) usedMemorySizeInBytesOfTablets + 
getReservedTsFileParserMemorySizeInBytes();
+        (double) usedMemorySizeInBytesOfTablets
+            + getReservedTsFileParserMemorySizeInBytes()
+            + extraMemoryInBytes;
     return tabletMemoryWithParserMemory + (double) 
usedMemorySizeInBytesOfTsFiles
             < allowedMaxMemorySizeInBytesOfTabletsAndTsFiles()
         && tabletMemoryWithParserMemory < 
allowedMaxMemorySizeInBytesOfTablets();
@@ -146,32 +165,235 @@ public class PipeMemoryManager {
         && (double) usedMemorySizeInBytesOfTablets < 
allowedMaxMemorySizeInBytesOfTablets();
   }
 
-  public synchronized boolean tryReserveTsFileParserMemory() {
-    if (!PIPE_MEMORY_MANAGEMENT_ENABLED) {
-      return true;
+  public synchronized boolean tryReserveTsFileParserMemory(
+      final String pipeName,
+      final long creationTime,
+      final String dataRegionId,
+      final TsFileParserMemoryReservation reservationKey) {
+    if (reservationKey == null) {
+      return false;
+    }
+
+    final PipeIdentity pipeIdentity = new PipeIdentity(pipeName, creationTime);
+    final PipeRegionIdentity pipeRegionIdentity =
+        new PipeRegionIdentity(pipeIdentity, dataRegionId);
+    enqueueTsFileParserReservationRequest(pipeRegionIdentity, reservationKey);
+
+    final int globalLimit = Math.max(1, 
PIPE_CONFIG.getPipeTsFileParserInFlightMaxNum());
+    final int perPipeRegionLimit =
+        Math.max(
+            1, Math.min(globalLimit, 
PIPE_CONFIG.getPipeTsFileParserInFlightMaxNumPerPipeRegion()));
+    final int reservedCountOfPipeRegion =
+        reservedTsFileParserCountByPipeRegion.getOrDefault(pipeRegionIdentity, 
0);
+    if (reservedTsFileParserCount >= globalLimit
+        || reservedCountOfPipeRegion >= perPipeRegionLimit) {
+      notifyNextTsFileParserMemoryReservationInternal();
+      return false;
     }
 
     final long parserMemorySizeInBytes = getTsFileParserMemorySizeInBytes();
-    if 
(isEnough4TabletParsingWithReservedParserMemory(parserMemorySizeInBytes)) {
-      reservedTsFileParserCount++;
-      return true;
+    final boolean isSoftMemoryEnough =
+        !PIPE_MEMORY_MANAGEMENT_ENABLED
+            || 
isEnough4TabletParsingWithReservedParserMemory(parserMemorySizeInBytes);
+    if (!isSoftMemoryEnough
+        && 
!isHardEnough4TabletParsingWithReservedParserMemory(parserMemorySizeInBytes)) {
+      return false;
     }
 
-    return false;
+    final PipeRegionIdentity nextPipeRegion =
+        getNextEligibleTsFileParserPipeRegion(perPipeRegionLimit, 
!isSoftMemoryEnough);
+    final LinkedHashSet<TsFileParserMemoryReservation> requestsOfPipeRegion =
+        waitingTsFileParserRequestsByPipeRegion.get(pipeRegionIdentity);
+    if (!pipeRegionIdentity.equals(nextPipeRegion)
+        || requestsOfPipeRegion == null
+        || !reservationKey.equals(requestsOfPipeRegion.iterator().next())) {
+      notifyNextTsFileParserMemoryReservationInternal();
+      return false;
+    }
+
+    removeTsFileParserReservationRequest(pipeRegionIdentity, reservationKey, 
true);
+    reservedTsFileParserCount++;
+    reservedTsFileParserCountByPipe.merge(pipeIdentity, 1, Integer::sum);
+    reservedTsFileParserCountByPipeRegion.put(pipeRegionIdentity, 
reservedCountOfPipeRegion + 1);
+    notifyNextTsFileParserMemoryReservationInternal();
+    return true;
   }
 
-  public synchronized void releaseTsFileParserMemory() {
-    if (!PIPE_MEMORY_MANAGEMENT_ENABLED) {
+  public synchronized void cancelTsFileParserMemoryReservation(
+      final String pipeName,
+      final long creationTime,
+      final String dataRegionId,
+      final TsFileParserMemoryReservation reservationKey) {
+    if (reservationKey == null) {
       return;
     }
+    removeTsFileParserReservationRequest(
+        new PipeRegionIdentity(new PipeIdentity(pipeName, creationTime), 
dataRegionId),
+        reservationKey,
+        false);
+    notifyNextTsFileParserMemoryReservationInternal();
+  }
 
-    reservedTsFileParserCount = Math.max(0, reservedTsFileParserCount - 1);
-    this.notifyAll();
+  public synchronized void releaseTsFileParserMemory(
+      final String pipeName, final long creationTime, final String 
dataRegionId) {
+    final PipeIdentity pipeIdentity = new PipeIdentity(pipeName, creationTime);
+    final PipeRegionIdentity pipeRegionIdentity =
+        new PipeRegionIdentity(pipeIdentity, dataRegionId);
+    final int reservedCountOfPipeRegion =
+        reservedTsFileParserCountByPipeRegion.getOrDefault(pipeRegionIdentity, 
0);
+    if (reservedCountOfPipeRegion <= 0) {
+      LOGGER.warn(
+          "Failed to release TsFile parser memory for pipe {} (creation time 
{}) in DataRegion {} because no reservation exists.",
+          pipeName,
+          creationTime,
+          dataRegionId);
+      return;
+    }
+
+    if (reservedCountOfPipeRegion == 1) {
+      reservedTsFileParserCountByPipeRegion.remove(pipeRegionIdentity);
+    } else {
+      reservedTsFileParserCountByPipeRegion.put(pipeRegionIdentity, 
reservedCountOfPipeRegion - 1);
+    }
+    final int reservedCountOfPipe = 
reservedTsFileParserCountByPipe.getOrDefault(pipeIdentity, 0);
+    if (reservedCountOfPipe == 1) {
+      reservedTsFileParserCountByPipe.remove(pipeIdentity);
+    } else {
+      reservedTsFileParserCountByPipe.put(pipeIdentity, reservedCountOfPipe - 
1);
+    }
+    reservedTsFileParserCount--;
+    notifyNextTsFileParserMemoryReservationInternal();
+  }
+
+  private void enqueueTsFileParserReservationRequest(
+      final PipeRegionIdentity pipeRegionIdentity,
+      final TsFileParserMemoryReservation reservationKey) {
+    final LinkedHashSet<TsFileParserMemoryReservation> requestsOfPipeRegion =
+        waitingTsFileParserRequestsByPipeRegion.computeIfAbsent(
+            pipeRegionIdentity,
+            key -> {
+              final ArrayDeque<PipeRegionIdentity> regionOrder =
+                  waitingTsFileParserRegionOrderByPipe.computeIfAbsent(
+                      key.pipeIdentity,
+                      pipe -> {
+                        waitingTsFileParserPipeOrder.addLast(pipe);
+                        return new ArrayDeque<>();
+                      });
+              regionOrder.addLast(key);
+              return new LinkedHashSet<>();
+            });
+    requestsOfPipeRegion.add(reservationKey);
+  }
+
+  public synchronized void notifyNextTsFileParserMemoryReservation() {
+    notifyNextTsFileParserMemoryReservationInternal();
+  }
+
+  private void notifyNextTsFileParserMemoryReservationInternal() {
+    final int globalLimit = Math.max(1, 
PIPE_CONFIG.getPipeTsFileParserInFlightMaxNum());
+    if (reservedTsFileParserCount >= globalLimit) {
+      return;
+    }
+
+    final long parserMemorySizeInBytes = getTsFileParserMemorySizeInBytes();
+    final boolean isSoftMemoryEnough =
+        !PIPE_MEMORY_MANAGEMENT_ENABLED
+            || 
isEnough4TabletParsingWithReservedParserMemory(parserMemorySizeInBytes);
+    if (!isSoftMemoryEnough
+        && 
!isHardEnough4TabletParsingWithReservedParserMemory(parserMemorySizeInBytes)) {
+      return;
+    }
+
+    final int perPipeRegionLimit =
+        Math.max(
+            1, Math.min(globalLimit, 
PIPE_CONFIG.getPipeTsFileParserInFlightMaxNumPerPipeRegion()));
+    final PipeRegionIdentity nextPipeRegion =
+        getNextEligibleTsFileParserPipeRegion(perPipeRegionLimit, 
!isSoftMemoryEnough);
+    if (nextPipeRegion == null) {
+      return;
+    }
+
+    final LinkedHashSet<TsFileParserMemoryReservation> requestsOfPipeRegion =
+        waitingTsFileParserRequestsByPipeRegion.get(nextPipeRegion);
+    if (requestsOfPipeRegion != null && !requestsOfPipeRegion.isEmpty()) {
+      requestsOfPipeRegion.iterator().next().signal();
+    }
   }
 
-  public synchronized void waitForTsFileParserMemory(final long timeoutInMs)
-      throws InterruptedException {
-    this.wait(Math.max(1, timeoutInMs));
+  private PipeRegionIdentity getNextEligibleTsFileParserPipeRegion(
+      final int perPipeRegionLimit, final boolean 
requirePipeWithoutReservedParser) {
+    for (final PipeIdentity pipeIdentity : waitingTsFileParserPipeOrder) {
+      // 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
+          && reservedTsFileParserCountByPipe.getOrDefault(pipeIdentity, 0) > 
0) {
+        continue;
+      }
+
+      final ArrayDeque<PipeRegionIdentity> regionOrder =
+          waitingTsFileParserRegionOrderByPipe.get(pipeIdentity);
+      if (regionOrder == null) {
+        continue;
+      }
+      for (final PipeRegionIdentity pipeRegionIdentity : regionOrder) {
+        if 
(reservedTsFileParserCountByPipeRegion.getOrDefault(pipeRegionIdentity, 0)
+            < perPipeRegionLimit) {
+          return pipeRegionIdentity;
+        }
+      }
+    }
+    return null;
+  }
+
+  private void removeTsFileParserReservationRequest(
+      final PipeRegionIdentity pipeRegionIdentity,
+      final TsFileParserMemoryReservation reservationKey,
+      final boolean rotateAfterAdmission) {
+    final LinkedHashSet<TsFileParserMemoryReservation> requestsOfPipeRegion =
+        waitingTsFileParserRequestsByPipeRegion.get(pipeRegionIdentity);
+    if (requestsOfPipeRegion == null || 
!requestsOfPipeRegion.remove(reservationKey)) {
+      return;
+    }
+
+    final PipeIdentity pipeIdentity = pipeRegionIdentity.pipeIdentity;
+    final ArrayDeque<PipeRegionIdentity> regionOrder =
+        waitingTsFileParserRegionOrderByPipe.get(pipeIdentity);
+    if (requestsOfPipeRegion.isEmpty()) {
+      waitingTsFileParserRequestsByPipeRegion.remove(pipeRegionIdentity);
+      if (regionOrder != null) {
+        regionOrder.remove(pipeRegionIdentity);
+        if (regionOrder.isEmpty()) {
+          waitingTsFileParserRegionOrderByPipe.remove(pipeIdentity);
+          waitingTsFileParserPipeOrder.remove(pipeIdentity);
+          return;
+        }
+      }
+    } else if (rotateAfterAdmission && regionOrder != null) {
+      regionOrder.remove(pipeRegionIdentity);
+      regionOrder.addLast(pipeRegionIdentity);
+    }
+
+    if (rotateAfterAdmission) {
+      waitingTsFileParserPipeOrder.remove(pipeIdentity);
+      waitingTsFileParserPipeOrder.addLast(pipeIdentity);
+    }
+  }
+
+  public static final class TsFileParserMemoryReservation {
+
+    private boolean isSignaled;
+
+    public synchronized void await(final long timeoutInMs) throws 
InterruptedException {
+      if (!isSignaled) {
+        wait(timeoutInMs);
+      }
+      isSignaled = false;
+    }
+
+    private synchronized void signal() {
+      isSignaled = true;
+      notify();
+    }
   }
 
   public boolean shouldReleaseTsFileParserOnOutOfMemory(
@@ -193,7 +415,7 @@ public class PipeMemoryManager {
       return elapsedTimeInMs >= maxRetryTimeInMs;
     }
 
-    if (!isHardEnough4TabletParsingWithReservedParserMemory()) {
+    if (!isHardEnough4TabletParsingWithReservedParserMemory(0)) {
       return true;
     }
 
@@ -436,6 +658,7 @@ public class PipeMemoryManager {
         allocatedBlocks.remove(block);
       }
 
+      notifyNextTsFileParserMemoryReservationInternal();
       this.notifyAll();
       return;
     }
@@ -730,6 +953,7 @@ public class PipeMemoryManager {
     }
     block.markAsReleased();
 
+    notifyNextTsFileParserMemoryReservationInternal();
     this.notifyAll();
   }
 
@@ -747,6 +971,7 @@ public class PipeMemoryManager {
     }
     block.setMemoryUsageInBytes(block.getMemoryUsageInBytes() - sizeInBytes);
 
+    notifyNextTsFileParserMemoryReservationInternal();
     this.notifyAll();
 
     return true;
@@ -783,4 +1008,61 @@ public class PipeMemoryManager {
   public static long getTotalMemorySizeInBytes() {
     return TOTAL_MEMORY_SIZE_IN_BYTES;
   }
+
+  private static class PipeIdentity {
+
+    private final String pipeName;
+    private final long creationTime;
+
+    private PipeIdentity(final String pipeName, final long creationTime) {
+      this.pipeName = pipeName;
+      this.creationTime = creationTime;
+    }
+
+    @Override
+    public boolean equals(final Object object) {
+      if (this == object) {
+        return true;
+      }
+      if (!(object instanceof PipeIdentity)) {
+        return false;
+      }
+      final PipeIdentity that = (PipeIdentity) object;
+      return creationTime == that.creationTime && Objects.equals(pipeName, 
that.pipeName);
+    }
+
+    @Override
+    public int hashCode() {
+      return Objects.hash(pipeName, creationTime);
+    }
+  }
+
+  private static class PipeRegionIdentity {
+
+    private final PipeIdentity pipeIdentity;
+    private final String dataRegionId;
+
+    private PipeRegionIdentity(final PipeIdentity pipeIdentity, final String 
dataRegionId) {
+      this.pipeIdentity = pipeIdentity;
+      this.dataRegionId = dataRegionId;
+    }
+
+    @Override
+    public boolean equals(final Object object) {
+      if (this == object) {
+        return true;
+      }
+      if (!(object instanceof PipeRegionIdentity)) {
+        return false;
+      }
+      final PipeRegionIdentity that = (PipeRegionIdentity) object;
+      return Objects.equals(pipeIdentity, that.pipeIdentity)
+          && Objects.equals(dataRegionId, that.dataRegionId);
+    }
+
+    @Override
+    public int hashCode() {
+      return Objects.hash(pipeIdentity, dataRegionId);
+    }
+  }
 }
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/conf/PropertiesTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/conf/PropertiesTest.java
index d547963d322..458871ec4f8 100755
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/conf/PropertiesTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/conf/PropertiesTest.java
@@ -19,6 +19,8 @@
 
 package org.apache.iotdb.db.conf;
 
+import org.apache.iotdb.commons.conf.CommonConfig;
+import org.apache.iotdb.commons.conf.CommonDescriptor;
 import org.apache.iotdb.commons.conf.TrimProperties;
 
 import com.tngtech.archunit.core.domain.JavaClasses;
@@ -35,6 +37,33 @@ import java.util.Properties;
 import static com.tngtech.archunit.lang.syntax.ArchRuleDefinition.noClasses;
 
 public class PropertiesTest {
+  @Test
+  public void testHotReloadTsFileParserInFlightLimits() throws Exception {
+    final IoTDBDescriptor descriptor = IoTDBDescriptor.getInstance();
+    final CommonConfig commonConfig = 
CommonDescriptor.getInstance().getConfig();
+    final int originalGlobalLimit = 
commonConfig.getPipeTsFileParserInFlightMaxNum();
+    final int originalPerPipeRegionLimit =
+        commonConfig.getPipeTsFileParserInFlightMaxNumPerPipeRegion();
+
+    try {
+      final TrimProperties properties = new TrimProperties();
+      properties.setProperty("pipe_tsfile_parser_in_flight_max_num", "3");
+      
properties.setProperty("pipe_tsfile_parser_in_flight_max_num_per_pipe_region", 
"2");
+      descriptor.loadHotModifiedProps(properties);
+
+      Assert.assertEquals(3, commonConfig.getPipeTsFileParserInFlightMaxNum());
+      Assert.assertEquals(2, 
commonConfig.getPipeTsFileParserInFlightMaxNumPerPipeRegion());
+    } finally {
+      final TrimProperties properties = new TrimProperties();
+      properties.setProperty(
+          "pipe_tsfile_parser_in_flight_max_num", 
Integer.toString(originalGlobalLimit));
+      properties.setProperty(
+          "pipe_tsfile_parser_in_flight_max_num_per_pipe_region",
+          Integer.toString(originalPerPipeRegionLimit));
+      descriptor.loadHotModifiedProps(properties);
+    }
+  }
+
   @Test
   public void PropertiesWithSpace() {
     IoTDBDescriptor descriptor = IoTDBDescriptor.getInstance();
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEventAdmissionTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEventAdmissionTest.java
index cfeff847699..59d613bb026 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEventAdmissionTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEventAdmissionTest.java
@@ -24,6 +24,7 @@ import org.apache.iotdb.commons.conf.CommonDescriptor;
 import org.apache.iotdb.commons.pipe.datastructure.pattern.PrefixPipePattern;
 import org.apache.iotdb.db.pipe.resource.PipeDataNodeResourceManager;
 import org.apache.iotdb.db.pipe.resource.memory.PipeMemoryManager;
+import 
org.apache.iotdb.db.pipe.resource.memory.PipeMemoryManager.TsFileParserMemoryReservation;
 import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
 import 
org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResourceStatus;
 import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;
@@ -43,54 +44,28 @@ import java.util.concurrent.TimeoutException;
 
 public class PipeTsFileInsertionEventAdmissionTest {
 
-  @Test
-  public void testParserAdmissionBackoffIsBoundedAndJittered() {
-    Assert.assertEquals(100, 
PipeTsFileInsertionEvent.getMaxMemoryCheckIntervalMs(10, 10));
-    Assert.assertEquals(10, 
PipeTsFileInsertionEvent.getMaxMemoryCheckIntervalMs(10, 0));
-    Assert.assertEquals(
-        Long.MAX_VALUE, 
PipeTsFileInsertionEvent.getMaxMemoryCheckIntervalMs(Long.MAX_VALUE, 10));
-
-    Assert.assertEquals(20, 
PipeTsFileInsertionEvent.getNextMemoryCheckIntervalMs(10, 100));
-    Assert.assertEquals(100, 
PipeTsFileInsertionEvent.getNextMemoryCheckIntervalMs(80, 100));
-    Assert.assertEquals(100, 
PipeTsFileInsertionEvent.getNextMemoryCheckIntervalMs(100, 100));
-
-    for (int i = 0; i < 100; i++) {
-      final long intervalWithJitter =
-          PipeTsFileInsertionEvent.getMemoryCheckIntervalWithJitter(100);
-      Assert.assertTrue(intervalWithJitter >= 50);
-      Assert.assertTrue(intervalWithJitter <= 100);
-    }
-  }
-
   @Test(timeout = 10000)
   public void testParserAdmissionIsWokenWhenMemoryIsReleased() throws 
Exception {
     final CommonConfig commonConfig = 
CommonDescriptor.getInstance().getConfig();
     final PipeMemoryManager memoryManager = 
PipeDataNodeResourceManager.memory();
     final long originalParserMemoryInBytes = 
commonConfig.getPipeTsFileParserMemory();
-    final long originalMemoryCheckIntervalMs = 
commonConfig.getPipeCheckMemoryEnoughIntervalMs();
-    final int originalMemoryAllocateMaxRetries = 
commonConfig.getPipeMemoryAllocateMaxRetries();
+    final int originalGlobalLimit = 
commonConfig.getPipeTsFileParserInFlightMaxNum();
+    final int originalPerPipeRegionLimit =
+        commonConfig.getPipeTsFileParserInFlightMaxNumPerPipeRegion();
+    final TsFileParserMemoryReservation blockerReservation = new 
TsFileParserMemoryReservation();
 
     File tsFile = null;
     PipeTsFileInsertionEvent event = null;
     ExecutorService executor = null;
     Future<Iterable<TabletInsertionEvent>> parsingFuture = null;
-    int blockerReservationCount = 0;
+    boolean isBlockerReserved = false;
     try {
-      commonConfig.setPipeTsFileParserMemory(
-          Math.max(1, memoryManager.getTotalNonFloatingMemorySizeInBytes() / 
8));
-      commonConfig.setPipeCheckMemoryEnoughIntervalMs(10000);
-      commonConfig.setPipeMemoryAllocateMaxRetries(10);
-
-      boolean parserMemoryExhausted = false;
-      for (int i = 0; i < 100; i++) {
-        if (!memoryManager.tryReserveTsFileParserMemory()) {
-          parserMemoryExhausted = true;
-          break;
-        }
-        blockerReservationCount++;
-      }
-      Assert.assertTrue(blockerReservationCount > 0);
-      Assert.assertTrue(parserMemoryExhausted);
+      commonConfig.setPipeTsFileParserMemory(1);
+      commonConfig.setPipeTsFileParserInFlightMaxNum(1);
+      commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(1);
+      isBlockerReserved =
+          memoryManager.tryReserveTsFileParserMemory("blocker", 0, "0", 
blockerReservation);
+      Assert.assertTrue(isBlockerReserved);
 
       tsFile =
           TsFileGeneratorUtils.generateNonAlignedTsFile(
@@ -107,7 +82,7 @@ public class PipeTsFileInsertionEventAdmissionTest {
               false,
               false,
               false,
-              null,
+              "testPipe",
               0,
               null,
               new PrefixPipePattern("root"),
@@ -123,8 +98,8 @@ public class PipeTsFileInsertionEventAdmissionTest {
           TimeoutException.class, () -> blockedParsingFuture.get(200, 
TimeUnit.MILLISECONDS));
 
       final long releaseTimeInNanos = System.nanoTime();
-      memoryManager.releaseTsFileParserMemory();
-      blockerReservationCount--;
+      memoryManager.releaseTsFileParserMemory("blocker", 0, "0");
+      isBlockerReserved = false;
 
       Assert.assertNotNull(parsingFuture.get(3, TimeUnit.SECONDS));
       Assert.assertTrue(
@@ -140,13 +115,13 @@ public class PipeTsFileInsertionEventAdmissionTest {
       if (event != null) {
         event.close();
       }
-      while (blockerReservationCount > 0) {
-        memoryManager.releaseTsFileParserMemory();
-        blockerReservationCount--;
+      memoryManager.cancelTsFileParserMemoryReservation("blocker", 0, "0", 
blockerReservation);
+      if (isBlockerReserved) {
+        memoryManager.releaseTsFileParserMemory("blocker", 0, "0");
       }
       commonConfig.setPipeTsFileParserMemory(originalParserMemoryInBytes);
-      
commonConfig.setPipeCheckMemoryEnoughIntervalMs(originalMemoryCheckIntervalMs);
-      
commonConfig.setPipeMemoryAllocateMaxRetries(originalMemoryAllocateMaxRetries);
+      commonConfig.setPipeTsFileParserInFlightMaxNum(originalGlobalLimit);
+      
commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(originalPerPipeRegionLimit);
       if (tsFile != null) {
         tsFile.delete();
       }
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
new file mode 100644
index 00000000000..72bffbc2c61
--- /dev/null
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManagerTest.java
@@ -0,0 +1,347 @@
+/*
+ * 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.db.pipe.resource.PipeDataNodeResourceManager;
+import 
org.apache.iotdb.db.pipe.resource.memory.PipeMemoryManager.TsFileParserMemoryReservation;
+
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+
+public class PipeMemoryManagerTest {
+
+  private final PipeMemoryManager memoryManager = 
PipeDataNodeResourceManager.memory();
+  private final CommonConfig commonConfig = 
CommonDescriptor.getInstance().getConfig();
+  private final List<Reservation> reservations = new ArrayList<>();
+  private int originalGlobalLimit;
+  private int originalPerPipeRegionLimit;
+  private long originalParserMemoryInBytes;
+
+  @Before
+  public void setUp() {
+    originalGlobalLimit = commonConfig.getPipeTsFileParserInFlightMaxNum();
+    originalPerPipeRegionLimit = 
commonConfig.getPipeTsFileParserInFlightMaxNumPerPipeRegion();
+    originalParserMemoryInBytes = commonConfig.getPipeTsFileParserMemory();
+    commonConfig.setPipeTsFileParserMemory(1);
+  }
+
+  @After
+  public void tearDown() {
+    for (final Reservation reservation : reservations) {
+      memoryManager.cancelTsFileParserMemoryReservation(
+          reservation.pipeName,
+          reservation.creationTime,
+          reservation.dataRegionId,
+          reservation.key);
+      if (reservation.acquired) {
+        memoryManager.releaseTsFileParserMemory(
+            reservation.pipeName, reservation.creationTime, 
reservation.dataRegionId);
+      }
+    }
+    commonConfig.setPipeTsFileParserInFlightMaxNum(originalGlobalLimit);
+    
commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(originalPerPipeRegionLimit);
+    commonConfig.setPipeTsFileParserMemory(originalParserMemoryInBytes);
+  }
+
+  @Test
+  public void testWaitingPipesAreAdmittedInRoundRobinOrder() {
+    commonConfig.setPipeTsFileParserInFlightMaxNum(1);
+    commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(1);
+
+    final Reservation pipeAActive = new Reservation("pipeA", 1);
+    final Reservation pipeAFirstWaiting = new Reservation("pipeA", 1);
+    final Reservation pipeASecondWaiting = new Reservation("pipeA", 1);
+    final Reservation pipeBWaiting = new Reservation("pipeB", 2);
+
+    Assert.assertTrue(tryAcquire(pipeAActive));
+    Assert.assertFalse(tryAcquire(pipeAFirstWaiting));
+    Assert.assertFalse(tryAcquire(pipeBWaiting));
+    Assert.assertFalse(tryAcquire(pipeASecondWaiting));
+
+    release(pipeAActive);
+    Assert.assertTrue(tryAcquire(pipeAFirstWaiting));
+    release(pipeAFirstWaiting);
+
+    // Pipe A still has another waiting TsFile, but it was rotated behind pipe 
B after admission.
+    Assert.assertFalse(tryAcquire(pipeASecondWaiting));
+    Assert.assertTrue(tryAcquire(pipeBWaiting));
+    release(pipeBWaiting);
+
+    Assert.assertTrue(tryAcquire(pipeASecondWaiting));
+  }
+
+  @Test
+  public void testGlobalAndPerPipeRegionLimitsAreBothEnforced() {
+    commonConfig.setPipeTsFileParserInFlightMaxNum(2);
+    commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(1);
+
+    final Reservation pipeAFirst = new Reservation("pipeA", 1);
+    final Reservation pipeASecond = new Reservation("pipeA", 1);
+    final Reservation pipeB = new Reservation("pipeB", 2);
+    final Reservation pipeC = new Reservation("pipeC", 3);
+
+    Assert.assertTrue(tryAcquire(pipeAFirst));
+    Assert.assertFalse(tryAcquire(pipeASecond));
+    Assert.assertTrue(tryAcquire(pipeB));
+    Assert.assertFalse(tryAcquire(pipeC));
+
+    release(pipeAFirst);
+    Assert.assertTrue(tryAcquire(pipeASecond));
+    Assert.assertFalse(tryAcquire(pipeC));
+
+    release(pipeB);
+    Assert.assertTrue(tryAcquire(pipeC));
+  }
+
+  @Test
+  public void testDifferentRegionsOfSamePipeCanRunConcurrently() {
+    commonConfig.setPipeTsFileParserInFlightMaxNum(2);
+    commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(1);
+
+    final Reservation pipeARegion1First = new Reservation("pipeA", 1, "1");
+    final Reservation pipeARegion1Second = new Reservation("pipeA", 1, "1");
+    final Reservation pipeARegion2 = new Reservation("pipeA", 1, "2");
+
+    Assert.assertTrue(tryAcquire(pipeARegion1First));
+    Assert.assertFalse(tryAcquire(pipeARegion1Second));
+    Assert.assertTrue(tryAcquire(pipeARegion2));
+  }
+
+  @Test
+  public void testWaitingRegionsWithinPipeAreAdmittedInRoundRobinOrder() {
+    commonConfig.setPipeTsFileParserInFlightMaxNum(1);
+    commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(1);
+
+    final Reservation blocker = new Reservation("blocker", 0);
+    final Reservation pipeARegion1First = new Reservation("pipeA", 1, "1");
+    final Reservation pipeARegion1Second = new Reservation("pipeA", 1, "1");
+    final Reservation pipeARegion2 = new Reservation("pipeA", 1, "2");
+
+    Assert.assertTrue(tryAcquire(blocker));
+    Assert.assertFalse(tryAcquire(pipeARegion1First));
+    Assert.assertFalse(tryAcquire(pipeARegion1Second));
+    Assert.assertFalse(tryAcquire(pipeARegion2));
+
+    release(blocker);
+    Assert.assertTrue(tryAcquire(pipeARegion1First));
+    release(pipeARegion1First);
+
+    Assert.assertFalse(tryAcquire(pipeARegion1Second));
+    Assert.assertTrue(tryAcquire(pipeARegion2));
+    release(pipeARegion2);
+
+    Assert.assertTrue(tryAcquire(pipeARegion1Second));
+  }
+
+  @Test
+  public void testPipeFairnessIsNotWeightedByRegionCount() {
+    commonConfig.setPipeTsFileParserInFlightMaxNum(1);
+    commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(1);
+
+    final Reservation blocker = new Reservation("blocker", 0);
+    final Reservation pipeARegion1 = new Reservation("pipeA", 1, "1");
+    final Reservation pipeARegion2 = new Reservation("pipeA", 1, "2");
+    final Reservation pipeARegion3 = new Reservation("pipeA", 1, "3");
+    final Reservation pipeBRegion1 = new Reservation("pipeB", 2, "1");
+
+    Assert.assertTrue(tryAcquire(blocker));
+    Assert.assertFalse(tryAcquire(pipeARegion1));
+    Assert.assertFalse(tryAcquire(pipeARegion2));
+    Assert.assertFalse(tryAcquire(pipeARegion3));
+    Assert.assertFalse(tryAcquire(pipeBRegion1));
+
+    release(blocker);
+    Assert.assertTrue(tryAcquire(pipeARegion1));
+    release(pipeARegion1);
+
+    Assert.assertFalse(tryAcquire(pipeARegion2));
+    Assert.assertTrue(tryAcquire(pipeBRegion1));
+    release(pipeBRegion1);
+
+    Assert.assertTrue(tryAcquire(pipeARegion2));
+  }
+
+  @Test
+  public void testSoftMemoryHeadroomIsReservedForPipeWithoutParser() {
+    commonConfig.setPipeTsFileParserInFlightMaxNum(2);
+    commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(2);
+
+    final double tabletMemoryLimit =
+        
(commonConfig.getPipeDataStructureTabletMemoryBlockAllocationRejectThreshold()
+                + 
commonConfig.getPipeDataStructureTsFileMemoryBlockAllocationRejectThreshold() / 
2)
+            * memoryManager.getTotalNonFloatingMemorySizeInBytes();
+    final double tabletAndTsFileMemoryLimit =
+        
(commonConfig.getPipeDataStructureTabletMemoryBlockAllocationRejectThreshold()
+                + 
commonConfig.getPipeDataStructureTsFileMemoryBlockAllocationRejectThreshold())
+            * memoryManager.getTotalNonFloatingMemorySizeInBytes();
+    commonConfig.setPipeTsFileParserMemory(
+        Math.max(1, (long) (Math.min(tabletMemoryLimit, 
tabletAndTsFileMemoryLimit) * 0.49)));
+
+    final Reservation pipeAActive = new Reservation("pipeA", 1, "1");
+    final Reservation pipeAWaiting = new Reservation("pipeA", 1, "2");
+    final Reservation pipeBWaiting = new Reservation("pipeB", 2, "1");
+
+    Assert.assertTrue(tryAcquire(pipeAActive));
+    Assert.assertFalse(tryAcquire(pipeAWaiting));
+
+    // The second parser would fit only below the hard threshold. Pipe A 
already has a parser, so
+    // the headroom must go to pipe B even though pipe A is ahead in the 
waiting queue.
+    Assert.assertTrue(tryAcquire(pipeBWaiting));
+  }
+
+  @Test
+  public void testConcurrentTsFilesFromMultiplePipesAreNotStarved() throws 
Exception {
+    commonConfig.setPipeTsFileParserInFlightMaxNum(1);
+    commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(1);
+
+    final Reservation blocker = new Reservation("blocker", 0);
+    Assert.assertTrue(tryAcquire(blocker));
+
+    // Each reservation represents a distinct TsFile event. Pipe A 
deliberately has more waiting
+    // TsFiles so the test can detect whether it monopolizes the single parser 
slot.
+    final List<Reservation> waitingTsFiles = new ArrayList<>();
+    for (int i = 0; i < 5; i++) {
+      waitingTsFiles.add(new Reservation("pipeA", 1));
+    }
+    for (int i = 0; i < 2; i++) {
+      waitingTsFiles.add(new Reservation("pipeB", 2));
+      waitingTsFiles.add(new Reservation("pipeC", 3));
+    }
+    reservations.addAll(waitingTsFiles);
+
+    final List<String> acquisitionOrder = Collections.synchronizedList(new 
ArrayList<>());
+    final CountDownLatch ready = new CountDownLatch(waitingTsFiles.size());
+    final CountDownLatch start = new CountDownLatch(1);
+    final CountDownLatch enqueued = new CountDownLatch(waitingTsFiles.size());
+    final ExecutorService executor = 
Executors.newFixedThreadPool(waitingTsFiles.size());
+    final List<Future<Boolean>> futures = new ArrayList<>();
+
+    try {
+      for (final Reservation reservation : waitingTsFiles) {
+        futures.add(
+            executor.submit(
+                () -> {
+                  ready.countDown();
+                  start.await();
+
+                  boolean acquired = tryAcquireWithoutTracking(reservation);
+                  enqueued.countDown();
+                  final long deadline = System.nanoTime() + 
TimeUnit.SECONDS.toNanos(10);
+                  while (!acquired && System.nanoTime() < deadline) {
+                    reservation.key.await(
+                        Math.max(
+                            1,
+                            TimeUnit.NANOSECONDS.toMillis(
+                                Math.max(1, deadline - System.nanoTime()))));
+                    acquired = tryAcquireWithoutTracking(reservation);
+                  }
+                  if (!acquired) {
+                    return false;
+                  }
+
+                  acquisitionOrder.add(reservation.pipeName);
+                  Thread.sleep(5);
+                  release(reservation);
+                  return true;
+                }));
+      }
+
+      Assert.assertTrue(ready.await(5, TimeUnit.SECONDS));
+      start.countDown();
+      Assert.assertTrue(enqueued.await(5, TimeUnit.SECONDS));
+      release(blocker);
+
+      for (final Future<Boolean> future : futures) {
+        Assert.assertTrue(future.get(15, TimeUnit.SECONDS));
+      }
+      Assert.assertEquals(waitingTsFiles.size(), acquisitionOrder.size());
+      Assert.assertEquals(3, new HashSet<>(acquisitionOrder.subList(0, 
3)).size());
+    } finally {
+      release(blocker);
+      executor.shutdownNow();
+      Assert.assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+    }
+  }
+
+  private boolean tryAcquire(final Reservation reservation) {
+    if (!reservations.contains(reservation)) {
+      reservations.add(reservation);
+    }
+    reservation.acquired =
+        memoryManager.tryReserveTsFileParserMemory(
+            reservation.pipeName,
+            reservation.creationTime,
+            reservation.dataRegionId,
+            reservation.key);
+    return reservation.acquired;
+  }
+
+  private boolean tryAcquireWithoutTracking(final Reservation reservation) {
+    reservation.acquired =
+        memoryManager.tryReserveTsFileParserMemory(
+            reservation.pipeName,
+            reservation.creationTime,
+            reservation.dataRegionId,
+            reservation.key);
+    return reservation.acquired;
+  }
+
+  private void release(final Reservation reservation) {
+    if (!reservation.acquired) {
+      return;
+    }
+    memoryManager.releaseTsFileParserMemory(
+        reservation.pipeName, reservation.creationTime, 
reservation.dataRegionId);
+    reservation.acquired = false;
+  }
+
+  private static class Reservation {
+
+    private final String pipeName;
+    private final long creationTime;
+    private final String dataRegionId;
+    private final TsFileParserMemoryReservation key = new 
TsFileParserMemoryReservation();
+    private volatile boolean acquired;
+
+    private Reservation(final String pipeName, final long creationTime) {
+      this(pipeName, creationTime, "0");
+    }
+
+    private Reservation(final String pipeName, final long creationTime, final 
String dataRegionId) {
+      this.pipeName = pipeName;
+      this.creationTime = creationTime;
+      this.dataRegionId = dataRegionId;
+    }
+  }
+}
diff --git 
a/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-system.properties.template
 
b/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-system.properties.template
index bda26c3cb93..2bc03179bd1 100644
--- 
a/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-system.properties.template
+++ 
b/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-system.properties.template
@@ -1811,6 +1811,19 @@ pipe_lib_dir=ext/pipe
 # Datatype: int
 pipe_subtask_executor_max_thread_num=0
 
+# The maximum number of TsFile parsers that can run concurrently across all 
Pipes on this DataNode.
+# When <= 0, use max(1, CPU core number / 2).
+# effectiveMode: hot_reload
+# Datatype: int
+pipe_tsfile_parser_in_flight_max_num=0
+
+# The maximum number of TsFile parsers that can run concurrently for one 
DataRegion of one Pipe.
+# Different DataRegions of the same Pipe have independent limits.
+# When <= 0, use 1.
+# effectiveMode: hot_reload
+# Datatype: int
+pipe_tsfile_parser_in_flight_max_num_per_pipe_region=1
+
 # The connection timeout (in milliseconds) for the thrift client.
 # effectiveMode: restart
 # Datatype: int
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
index e178ae230c8..6bc02fbef2c 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
@@ -237,6 +237,13 @@ public class CommonConfig {
   // Note: Pipes that do not decompose pattern/time do not need this part of 
memory
   private long pipeTsFileParserMemory = 17 * MB;
 
+  // Limit concurrently active TsFile parsers globally and for each region 
task of a pipe. The
+  // per-pipe-region limit also serves as an approximate parser memory quota 
because every admitted
+  // parser reserves pipeTsFileParserMemory bytes.
+  private int pipeTsFileParserInFlightMaxNum =
+      Math.max(1, Runtime.getRuntime().availableProcessors() / 2);
+  private int pipeTsFileParserInFlightMaxNumPerPipeRegion = 1;
+
   // Memory for Sink batch sending (InsertNode/TsFile, choose one)
   // 1. InsertNode: 15MB, used for batch sending data to the downstream system
   private long pipeSinkBatchMemoryInsertNode = 15 * MB;
@@ -895,6 +902,36 @@ public class CommonConfig {
     logger.info("pipeTsFileParserMemory is set to {}.", 
pipeTsFileParserMemory);
   }
 
+  public int getPipeTsFileParserInFlightMaxNum() {
+    return pipeTsFileParserInFlightMaxNum;
+  }
+
+  public void setPipeTsFileParserInFlightMaxNum(final int 
pipeTsFileParserInFlightMaxNum) {
+    final int validatedValue =
+        pipeTsFileParserInFlightMaxNum > 0
+            ? pipeTsFileParserInFlightMaxNum
+            : Math.max(1, Runtime.getRuntime().availableProcessors() / 2);
+    if (this.pipeTsFileParserInFlightMaxNum == validatedValue) {
+      return;
+    }
+    this.pipeTsFileParserInFlightMaxNum = validatedValue;
+    logger.info("pipeTsFileParserInFlightMaxNum is set to {}.", 
validatedValue);
+  }
+
+  public int getPipeTsFileParserInFlightMaxNumPerPipeRegion() {
+    return pipeTsFileParserInFlightMaxNumPerPipeRegion;
+  }
+
+  public void setPipeTsFileParserInFlightMaxNumPerPipeRegion(
+      final int pipeTsFileParserInFlightMaxNumPerPipeRegion) {
+    final int validatedValue = Math.max(1, 
pipeTsFileParserInFlightMaxNumPerPipeRegion);
+    if (this.pipeTsFileParserInFlightMaxNumPerPipeRegion == validatedValue) {
+      return;
+    }
+    this.pipeTsFileParserInFlightMaxNumPerPipeRegion = validatedValue;
+    logger.info("pipeTsFileParserInFlightMaxNumPerPipeRegion is set to {}.", 
validatedValue);
+  }
+
   public long getPipeSinkBatchMemoryInsertNode() {
     return pipeSinkBatchMemoryInsertNode;
   }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeConfig.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeConfig.java
index a0456f7a5c4..fe81318f5dd 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeConfig.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeConfig.java
@@ -83,6 +83,14 @@ public class PipeConfig {
     return COMMON_CONFIG.getPipeTsFileParserMemory();
   }
 
+  public int getPipeTsFileParserInFlightMaxNum() {
+    return COMMON_CONFIG.getPipeTsFileParserInFlightMaxNum();
+  }
+
+  public int getPipeTsFileParserInFlightMaxNumPerPipeRegion() {
+    return COMMON_CONFIG.getPipeTsFileParserInFlightMaxNumPerPipeRegion();
+  }
+
   public long getSinkBatchMemoryInsertNode() {
     return COMMON_CONFIG.getPipeSinkBatchMemoryInsertNode();
   }
@@ -502,6 +510,10 @@ public class PipeConfig {
 
     LOGGER.info("IsPipeEnableMemoryCheck: {}", isPipeEnableMemoryCheck());
     LOGGER.info("PipeTsFileParserMemory: {}", getTsFileParserMemory());
+    LOGGER.info("PipeTsFileParserInFlightMaxNum: {}", 
getPipeTsFileParserInFlightMaxNum());
+    LOGGER.info(
+        "PipeTsFileParserInFlightMaxNumPerPipeRegion: {}",
+        getPipeTsFileParserInFlightMaxNumPerPipeRegion());
     LOGGER.info("SinkBatchMemoryInsertNode: {}", 
getSinkBatchMemoryInsertNode());
     LOGGER.info("SinkBatchMemoryTsFile: {}", getSinkBatchMemoryTsFile());
     LOGGER.info("SendTsFileReadBuffer: {}", getSendTsFileReadBuffer());
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeDescriptor.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeDescriptor.java
index b9115769e11..2db65d60d5e 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeDescriptor.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeDescriptor.java
@@ -238,6 +238,16 @@ public class PipeDescriptor {
         Long.parseLong(
             properties.getProperty(
                 "pipe_tsfile_parser_memory", 
String.valueOf(config.getPipeTsFileParserMemory()))));
+    config.setPipeTsFileParserInFlightMaxNum(
+        Integer.parseInt(
+            properties.getProperty(
+                "pipe_tsfile_parser_in_flight_max_num",
+                String.valueOf(config.getPipeTsFileParserInFlightMaxNum()))));
+    config.setPipeTsFileParserInFlightMaxNumPerPipeRegion(
+        Integer.parseInt(
+            properties.getProperty(
+                "pipe_tsfile_parser_in_flight_max_num_per_pipe_region",
+                
String.valueOf(config.getPipeTsFileParserInFlightMaxNumPerPipeRegion()))));
     config.setPipeSinkBatchMemoryInsertNode(
         Long.parseLong(
             properties.getProperty(

Reply via email to