diff --git a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java index f13bca9a4a8a6..c56efafbcd9d8 100644 --- a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java +++ b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java @@ -2585,4 +2585,6 @@ private DataNodePipeMessages() {} "Topic metadata for %s is unavailable during consensus subscription setup"; public static final String EXCEPTION_TOPIC_CONFIG_FOR_ARG_IS_UNAVAILABLE_DURING_CONSENSUS_SUBSCRIPTION_SETUP_B94404EE = "Topic config for %s is unavailable during consensus subscription setup"; + public static final String LOG_FAILED_TO_RELEASE_TSFILE_PARSER_MEMORY_FOR_PIPE_ARG_CREATION_TIME_ARG_IN_DATAREGION_ARG_BECAUSE_NO_RESERVATION_EXISTS_BB8321C0 = + "Failed to release TsFile parser memory for Pipe {} (creation time {}) in DataRegion {} because no reservation exists."; } diff --git a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java index 17dd4c6d836f7..3c4cc12959622 100644 --- a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java +++ b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java @@ -2413,4 +2413,6 @@ private DataNodePipeMessages() {} "共识订阅设置期间 topic %s 的元数据不可用"; public static final String EXCEPTION_TOPIC_CONFIG_FOR_ARG_IS_UNAVAILABLE_DURING_CONSENSUS_SUBSCRIPTION_SETUP_B94404EE = "共识订阅设置期间 topic %s 的配置不可用"; + public static final String LOG_FAILED_TO_RELEASE_TSFILE_PARSER_MEMORY_FOR_PIPE_ARG_CREATION_TIME_ARG_IN_DATAREGION_ARG_BECAUSE_NO_RESERVATION_EXISTS_BB8321C0 = + "无法释放 Pipe {}(创建时间 {})在 DataRegion {} 中的 TsFile 解析器内存,因为不存在对应的预留。"; } 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 ecd3d73a3e30d..09584dfe61907 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 @@ -42,6 +42,7 @@ import org.apache.iotdb.db.conf.rest.IoTDBRestServiceDescriptor; import org.apache.iotdb.db.consensus.DataRegionConsensusImpl; import org.apache.iotdb.db.i18n.DataNodeMiscMessages; +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; @@ -2766,6 +2767,7 @@ private void loadLoadTsFileHotModifiedProp(TrimProperties properties) throws IOE 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 988f57ce5cb4f..ba2a9d275c39d 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 @@ -47,6 +47,7 @@ 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; @@ -69,7 +70,6 @@ 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; @@ -81,6 +81,7 @@ public class PipeTsFileInsertionEvent extends PipeInsertionEvent private static final Logger LOGGER = LoggerFactory.getLogger(PipeTsFileInsertionEvent.class); private final TsFileResource resource; + private final String dataRegionId; private File tsFile; private long extractTime = 0; @@ -99,6 +100,8 @@ public class PipeTsFileInsertionEvent extends PipeInsertionEvent private final AtomicBoolean isClosed; private final AtomicReference eventParser; private final AtomicBoolean isTsFileParserMemoryReserved = new AtomicBoolean(false); + private final TsFileParserMemoryReservation tsFileParserMemoryReservationKey = + new TsFileParserMemoryReservation(); private final AtomicReference> tabletInsertionEventIterator = new AtomicReference<>(); private final AtomicReference pendingTabletInsertionEvent = @@ -229,6 +232,7 @@ private PipeTsFileInsertionEvent( databaseNameFromDataRegion); 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 @@ -294,6 +298,17 @@ private PipeTsFileInsertionEvent( }); } + 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. @@ -944,32 +959,12 @@ private void waitForResourceEnough4Parsing(final long timeoutMs) throws Interrup 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( - DataNodePipeMessages - .PIPE_EXCEPTION_TIMEOUTEXCEPTION_WAITED_S_SECONDS_FOR_MEMORY_TO_PARSE_TSFILE_0E4EF8FD, - 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( DataNodePipeMessages.WAIT_FOR_MEMORY_ENOUGH_FOR_PARSING_FOR, resource != null ? resource.getTsFilePath() : "tsfile", @@ -982,35 +977,28 @@ private void waitForResourceEnough4Parsing(final long timeoutMs) throws Interrup waitTimeSeconds); } - if (tryReserveTsFileParserMemory(memoryManager)) { - LOGGER.info( - DataNodePipeMessages.WAIT_FOR_MEMORY_ENOUGH_FOR_PARSING_FOR, - resource != null ? resource.getTsFilePath() : "tsfile", - waitTimeSeconds); - return; + if (waitTimeInMs > timeoutMs) { + // should contain 'TimeoutException' in exception message + throw new PipeRuntimeOutOfMemoryCriticalException( + String.format( + DataNodePipeMessages + .PIPE_EXCEPTION_TIMEOUTEXCEPTION_WAITED_S_SECONDS_FOR_MEMORY_TO_PARSE_TSFILE_0E4EF8FD, + 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( + DataNodePipeMessages.WAIT_FOR_MEMORY_ENOUGH_FOR_PARSING_FOR, + resource != null ? resource.getTsFilePath() : "tsfile", + waitTimeSeconds); } private boolean tryReserveTsFileParserMemory(final PipeMemoryManager memoryManager) { @@ -1019,7 +1007,8 @@ private boolean tryReserveTsFileParserMemory(final PipeMemoryManager memoryManag return true; } - if (!memoryManager.tryReserveTsFileParserMemory()) { + if (!memoryManager.tryReserveTsFileParserMemory( + pipeName, creationTime, dataRegionId, tsFileParserMemoryReservationKey)) { return false; } @@ -1031,11 +1020,20 @@ private boolean tryReserveTsFileParserMemory(final PipeMemoryManager memoryManag 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 IoTConsensusV2 */ public boolean isGeneratedByIoTConsensusV2() { return isGeneratedByIoTConsensusV2; @@ -1104,6 +1102,7 @@ public long count(final boolean skipReportOnCommit) throws Exception { /** Release the resource of {@link TsFileInsertionEventParser}. */ @Override public void close() { + cancelTsFileParserMemoryReservationIfPending(); tabletInsertionEventIterator.set(null); releaseParsedTabletEvent(pendingTabletInsertionEvent.getAndSet(null)); parsedTabletInsertionEventCount.set(0); @@ -1153,12 +1152,14 @@ public PipeEventResource eventResourceBuilder() { this.referenceCount, this.pipeName, this.creationTime, + this.dataRegionId, this.tsFile, this.isWithMod, this.modFile, this.sharedModFile, this.eventParser, - this.isTsFileParserMemoryReserved); + this.isTsFileParserMemoryReserved, + this.tsFileParserMemoryReservationKey); } private static class PipeTsFileInsertionEventResource extends PipeEventResource { @@ -1170,33 +1171,42 @@ private static class PipeTsFileInsertionEventResource extends PipeEventResource private final AtomicReference eventParser; 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 File sharedModFile, final AtomicReference eventParser, - 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.sharedModFile = sharedModFile; this.eventParser = eventParser; this.isTsFileParserMemoryReserved = isTsFileParserMemoryReserved; + this.tsFileParserMemoryReservationKey = tsFileParserMemoryReservationKey; } @Override protected void finalizeResource() { try { + PipeDataNodeResourceManager.memory() + .cancelTsFileParserMemoryReservation( + pipeName, creationTime, dataRegionId, tsFileParserMemoryReservationKey); final String pipeTsFileResourcePipeName = PipeTsFileResourceManager.getPipeTsFileResourcePipeName(pipeName, creationTime); // decrease reference count @@ -1217,7 +1227,8 @@ protected void finalizeResource() { }); 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 acd85b983780b..46d957424f51c 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 @@ -31,10 +31,15 @@ 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; @@ -58,6 +63,15 @@ public class PipeMemoryManager { private volatile long reservedTsFileParserCount; + private final Map reservedTsFileParserCountByPipe = new HashMap<>(); + private final Map reservedTsFileParserCountByPipeRegion = + new HashMap<>(); + private final Map> + waitingTsFileParserRequestsByPipeRegion = new HashMap<>(); + private final Map> + waitingTsFileParserRegionOrderByPipe = new HashMap<>(); + private final ArrayDeque waitingTsFileParserPipeOrder = new ArrayDeque<>(); + // Only non-zero memory blocks will be added to this set. private final Set allocatedBlocks = new HashSet<>(); private final Set shrinkableBlocks = new HashSet<>(); @@ -128,9 +142,12 @@ private boolean isEnough4TabletParsingWithReservedParserMemory(final long extraM < 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(); @@ -149,32 +166,236 @@ < allowedMaxMemorySizeInBytesOfTabletsAndTsFiles() && (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 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( + DataNodePipeMessages + .LOG_FAILED_TO_RELEASE_TSFILE_PARSER_MEMORY_FOR_PIPE_ARG_CREATION_TIME_ARG_IN_DATAREGION_ARG_BECAUSE_NO_RESERVATION_EXISTS_BB8321C0, + 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 requestsOfPipeRegion = + waitingTsFileParserRequestsByPipeRegion.computeIfAbsent( + pipeRegionIdentity, + key -> { + final ArrayDeque 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(); } - public synchronized void waitForTsFileParserMemory(final long timeoutInMs) - throws InterruptedException { - this.wait(Math.max(1, timeoutInMs)); + 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 requestsOfPipeRegion = + waitingTsFileParserRequestsByPipeRegion.get(nextPipeRegion); + if (requestsOfPipeRegion != null && !requestsOfPipeRegion.isEmpty()) { + requestsOfPipeRegion.iterator().next().signal(); + } + } + + 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 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 requestsOfPipeRegion = + waitingTsFileParserRequestsByPipeRegion.get(pipeRegionIdentity); + if (requestsOfPipeRegion == null || !requestsOfPipeRegion.remove(reservationKey)) { + return; + } + + final PipeIdentity pipeIdentity = pipeRegionIdentity.pipeIdentity; + final ArrayDeque 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( @@ -194,7 +415,7 @@ public boolean shouldReleaseTsFileParserOnOutOfMemory( return elapsedTimeInMs >= maxRetryTimeInMs; } - if (!isHardEnough4TabletParsingWithReservedParserMemory()) { + if (!isHardEnough4TabletParsingWithReservedParserMemory(0)) { return true; } @@ -442,6 +663,7 @@ public synchronized void resize( allocatedBlocks.remove(block); } + notifyNextTsFileParserMemoryReservationInternal(); this.notifyAll(); return; } @@ -732,6 +954,7 @@ public synchronized void release(PipeMemoryBlock block) { } block.markAsReleased(); + notifyNextTsFileParserMemoryReservationInternal(); this.notifyAll(); } @@ -749,6 +972,7 @@ public synchronized boolean release(PipeMemoryBlock block, long sizeInBytes) { } block.setMemoryUsageInBytes(block.getMemoryUsageInBytes() - sizeInBytes); + notifyNextTsFileParserMemoryReservationInternal(); this.notifyAll(); return true; @@ -785,4 +1009,61 @@ public long getTotalFloatingMemorySizeInBytes() { public long getTotalMemorySizeInBytes() { return memoryBlock.getTotalMemorySizeInBytes(); } + + 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 24cb6a1b8c522..aad4851f6bc12 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.ConfigurationFileUtils; import org.apache.iotdb.commons.conf.TrimProperties; import org.apache.iotdb.commons.utils.RegionMigrationFileRemoveRateLimiter; @@ -39,6 +41,33 @@ 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 testHotReloadRegionMigrationFileRemoveSpeedLimit() throws Exception { 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 d3d7ae9cf8e9f..74bea8a5ea5e1 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.pipe.datastructure.pattern.PrefixTreePattern; 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; @@ -42,54 +43,28 @@ 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> 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( @@ -109,7 +84,7 @@ public void testParserAdmissionIsWokenWhenMemoryIsReleased() throws Exception { false, false, null, - null, + "testPipe", 0, null, new PrefixTreePattern("root"), @@ -130,8 +105,8 @@ public void testParserAdmissionIsWokenWhenMemoryIsReleased() throws Exception { 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( @@ -147,13 +122,13 @@ public void testParserAdmissionIsWokenWhenMemoryIsReleased() throws Exception { 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 0000000000000..72bffbc2c610c --- /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 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 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 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> 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 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 dd3ba2d295ecb..02dcb1b3b5de7 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 @@ -2036,6 +2036,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/i18n/en/org/apache/iotdb/commons/i18n/PipeMessages.java b/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/PipeMessages.java index e69dc15c1ddbd..742239bc1ca75 100644 --- a/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/PipeMessages.java +++ b/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/PipeMessages.java @@ -48,6 +48,10 @@ public final class PipeMessages { public static final String CONFIG_IS_PIPE_ENABLE_MEMORY_CHECK = "IsPipeEnableMemoryCheck: {}"; public static final String CONFIG_PIPE_TSFILE_PARSER_MEMORY = "PipeTsFileParserMemory: {}"; + public static final String CONFIG_PIPE_TSFILE_PARSER_IN_FLIGHT_MAX_NUM = + "PipeTsFileParserInFlightMaxNum: {}"; + public static final String CONFIG_PIPE_TSFILE_PARSER_IN_FLIGHT_MAX_NUM_PER_PIPE_REGION = + "PipeTsFileParserInFlightMaxNumPerPipeRegion: {}"; public static final String CONFIG_SINK_BATCH_MEMORY_INSERT_NODE = "SinkBatchMemoryInsertNode: {}"; public static final String CONFIG_SINK_BATCH_MEMORY_TSFILE = "SinkBatchMemoryTsFile: {}"; diff --git a/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/PipeMessages.java b/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/PipeMessages.java index ca3243cdc01d7..0a8c0badb3544 100644 --- a/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/PipeMessages.java +++ b/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/PipeMessages.java @@ -47,6 +47,10 @@ public final class PipeMessages { public static final String CONFIG_IS_PIPE_ENABLE_MEMORY_CHECK = "IsPipeEnableMemoryCheck: {}"; public static final String CONFIG_PIPE_TSFILE_PARSER_MEMORY = "PipeTsFileParserMemory: {}"; + public static final String CONFIG_PIPE_TSFILE_PARSER_IN_FLIGHT_MAX_NUM = + "PipeTsFileParserInFlightMaxNum: {}"; + public static final String CONFIG_PIPE_TSFILE_PARSER_IN_FLIGHT_MAX_NUM_PER_PIPE_REGION = + "PipeTsFileParserInFlightMaxNumPerPipeRegion: {}"; public static final String CONFIG_SINK_BATCH_MEMORY_INSERT_NODE = "SinkBatchMemoryInsertNode: {}"; public static final String CONFIG_SINK_BATCH_MEMORY_TSFILE = "SinkBatchMemoryTsFile: {}"; 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 14675b8b19e7a..73f5bf52fc33e 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 @@ -248,6 +248,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; @@ -1041,6 +1048,39 @@ public void setPipeTsFileParserMemory(long pipeTsFileParserMemory) { logger.info(ConfigMessages.CONFIG_SET_TO, "pipeTsFileParserMemory", 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(ConfigMessages.CONFIG_SET_TO, "pipeTsFileParserInFlightMaxNum", 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( + ConfigMessages.CONFIG_SET_TO, + "pipeTsFileParserInFlightMaxNumPerPipeRegion", + 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 6e85519a85f68..561d2923c07e7 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 @@ -84,6 +84,14 @@ public long getTsFileParserMemory() { 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(); } @@ -506,6 +514,12 @@ public void printAllConfigs() { LOGGER.info(PipeMessages.CONFIG_IS_PIPE_ENABLE_MEMORY_CHECK, isPipeEnableMemoryCheck()); LOGGER.info(PipeMessages.CONFIG_PIPE_TSFILE_PARSER_MEMORY, getTsFileParserMemory()); + LOGGER.info( + PipeMessages.CONFIG_PIPE_TSFILE_PARSER_IN_FLIGHT_MAX_NUM, + getPipeTsFileParserInFlightMaxNum()); + LOGGER.info( + PipeMessages.CONFIG_PIPE_TSFILE_PARSER_IN_FLIGHT_MAX_NUM_PER_PIPE_REGION, + getPipeTsFileParserInFlightMaxNumPerPipeRegion()); LOGGER.info(PipeMessages.CONFIG_SINK_BATCH_MEMORY_INSERT_NODE, getSinkBatchMemoryInsertNode()); LOGGER.info(PipeMessages.CONFIG_SINK_BATCH_MEMORY_TSFILE, getSinkBatchMemoryTsFile()); LOGGER.info(PipeMessages.CONFIG_SEND_TSFILE_READ_BUFFER, 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 7909ef89e80d9..3a1da5f18ba46 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 static void loadPipeInternalConfig(CommonConfig config, TrimProperties pr 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(