From 5d134803f7588e40d10b806ddd1211408016e1f9 Mon Sep 17 00:00:00 2001 From: Zhenyu Luo Date: Fri, 10 Jul 2026 14:15:25 +0800 Subject: [PATCH] Optimize pipe TsFile cleanup on drop (#18163) * Optimize pipe TsFile cleanup on drop * Use ThreadName for pipe TsFile cleanup executor * Make pipe TsFile cleanup thread lazy * Use IoTDB thread pool factory for pipe cleanup * Clarify pipe TsFile cleanup deletion guard * Reuse periodical cleaner for Pipe TsFile cleanup * Clarify stale Pipe dir collection * Revert "Clarify stale Pipe dir collection" This reverts commit 25625fb1cf2a76fb54360800b990d27449b12feb. * Avoid duplicate stale Pipe dir collection (cherry picked from commit 350c6c7bcd25b883eca53897b294b57637850119) --- .../agent/task/PipeDataNodeTaskAgent.java | 22 +++++ .../tsfile/PipeTsFileInsertionEvent.java | 33 +++++-- ...HardlinkOrCopiedFileDirStartupCleaner.java | 86 +++++++++++++------ .../resource/tsfile/PipeTsFileResource.java | 8 +- .../tsfile/PipeTsFileResourceManager.java | 44 +++++++++- 5 files changed, 159 insertions(+), 34 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java index 17504ab65573c..31aaf36422f68 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java @@ -56,6 +56,7 @@ import org.apache.iotdb.db.pipe.metric.overview.PipeTsFileToTabletsMetrics; import org.apache.iotdb.db.pipe.resource.PipeDataNodeResourceManager; import org.apache.iotdb.db.pipe.resource.memory.PipeMemoryManager; +import org.apache.iotdb.db.pipe.resource.tsfile.PipeTsFileResourceManager; import org.apache.iotdb.db.pipe.source.dataregion.DataRegionListeningFilter; import org.apache.iotdb.db.pipe.source.dataregion.realtime.listener.PipeInsertionDataNodeListener; import org.apache.iotdb.db.pipe.source.schemaregion.SchemaRegionListeningFilter; @@ -304,13 +305,20 @@ protected void freezeRate(final String pipeName, final long creationTime) { @Override protected boolean dropPipe(final String pipeName, final long creationTime) { + final String pipeTsFileResourcePipeName = + PipeTsFileResourceManager.getPipeTsFileResourcePipeName(pipeName, creationTime); + PipeDataNodeResourceManager.tsfile().markPipeTsFileDirUnderDeletion(pipeTsFileResourcePipeName); + if (!super.dropPipe(pipeName, creationTime)) { + PipeDataNodeResourceManager.tsfile() + .unmarkPipeTsFileDirUnderDeletion(pipeTsFileResourcePipeName); return false; } final String taskId = pipeName + "_" + creationTime; PipeTsFileToTabletsMetrics.getInstance().deregister(taskId); PipeDataNodeSinglePipeMetrics.getInstance().deregister(taskId); + PipeDataNodeResourceManager.tsfile().cleanPipeTsFileDir(pipeTsFileResourcePipeName); return true; } @@ -319,6 +327,15 @@ protected boolean dropPipe(final String pipeName, final long creationTime) { protected boolean dropPipe(final String pipeName) { // Get the pipe meta first because it is removed after super#dropPipe(pipeName) final PipeMeta pipeMeta = pipeMetaKeeper.getPipeMeta(pipeName); + final String pipeTsFileResourcePipeName = + Objects.isNull(pipeMeta) + ? null + : PipeTsFileResourceManager.getPipeTsFileResourcePipeName( + pipeName, pipeMeta.getStaticMeta().getCreationTime()); + if (Objects.nonNull(pipeTsFileResourcePipeName)) { + PipeDataNodeResourceManager.tsfile() + .markPipeTsFileDirUnderDeletion(pipeTsFileResourcePipeName); + } // Record whether there are pipe tasks before dropping the pipe final boolean hasPipeTasks; @@ -331,6 +348,10 @@ protected boolean dropPipe(final String pipeName) { } if (!super.dropPipe(pipeName)) { + if (Objects.nonNull(pipeTsFileResourcePipeName)) { + PipeDataNodeResourceManager.tsfile() + .unmarkPipeTsFileDirUnderDeletion(pipeTsFileResourcePipeName); + } return false; } @@ -339,6 +360,7 @@ protected boolean dropPipe(final String pipeName) { final String taskId = pipeName + "_" + creationTime; PipeTsFileToTabletsMetrics.getInstance().deregister(taskId); PipeDataNodeSinglePipeMetrics.getInstance().deregister(taskId); + PipeDataNodeResourceManager.tsfile().cleanPipeTsFileDir(pipeTsFileResourcePipeName); // When the pipe contains no pipe tasks, there is no corresponding prefetching queue for the // subscribed pipe, so the subscription needs to be manually marked as completed. if (!hasPipeTasks && PipeStaticMeta.isSubscriptionPipe(pipeName)) { 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 ece6e25a151ae..c3518d0d628c8 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 @@ -321,11 +321,16 @@ public long getExtractTime() { @Override public boolean internallyIncreaseResourceReferenceCount(final String holderMessage) { extractTime = System.nanoTime(); + final String pipeTsFileResourcePipeName = + PipeTsFileResourceManager.getPipeTsFileResourcePipeName(pipeName, creationTime); try { - tsFile = PipeDataNodeResourceManager.tsfile().increaseFileReference(tsFile, true, pipeName); + tsFile = + PipeDataNodeResourceManager.tsfile() + .increaseFileReference(tsFile, true, pipeTsFileResourcePipeName); if (isWithMod) { modFile = - PipeDataNodeResourceManager.tsfile().increaseFileReference(modFile, false, pipeName); + PipeDataNodeResourceManager.tsfile() + .increaseFileReference(modFile, false, pipeTsFileResourcePipeName); } return true; } catch (final Exception e) { @@ -345,10 +350,14 @@ public boolean internallyIncreaseResourceReferenceCount(final String holderMessa @Override public boolean internallyDecreaseResourceReferenceCount(final String holderMessage) { + final String pipeTsFileResourcePipeName = + PipeTsFileResourceManager.getPipeTsFileResourcePipeName(pipeName, creationTime); try { - PipeDataNodeResourceManager.tsfile().decreaseFileReference(tsFile, pipeName); + PipeDataNodeResourceManager.tsfile() + .decreaseFileReference(tsFile, pipeTsFileResourcePipeName); if (isWithMod) { - PipeDataNodeResourceManager.tsfile().decreaseFileReference(modFile, pipeName); + PipeDataNodeResourceManager.tsfile() + .decreaseFileReference(modFile, pipeTsFileResourcePipeName); } close(); return true; @@ -461,7 +470,9 @@ public boolean mayEventPathsOverlappedWithPattern() { PipeDataNodeResourceManager.tsfile() .getDeviceIsAlignedMapFromCache( PipeTsFileResourceManager.getHardlinkOrCopiedFileInPipeDir( - resource.getTsFile(), pipeName), + resource.getTsFile(), + PipeTsFileResourceManager.getPipeTsFileResourcePipeName( + pipeName, creationTime)), false); final Set deviceSet = Objects.nonNull(deviceIsAlignedMap) ? deviceIsAlignedMap.keySet() : resource.getDevices(); @@ -877,6 +888,7 @@ public PipeEventResource eventResourceBuilder() { this.isReleased, this.referenceCount, this.pipeName, + this.creationTime, this.tsFile, this.isWithMod, this.modFile, @@ -891,12 +903,14 @@ private static class PipeTsFileInsertionEventResource extends PipeEventResource private final File modFile; private final AtomicReference dataContainer; private final String pipeName; + private final long creationTime; private final AtomicBoolean isTsFileParserMemoryReserved; private PipeTsFileInsertionEventResource( final AtomicBoolean isReleased, final AtomicInteger referenceCount, final String pipeName, + final long creationTime, final File tsFile, final boolean isWithMod, final File modFile, @@ -904,6 +918,7 @@ private PipeTsFileInsertionEventResource( final AtomicBoolean isTsFileParserMemoryReserved) { super(isReleased, referenceCount); this.pipeName = pipeName; + this.creationTime = creationTime; this.tsFile = tsFile; this.isWithMod = isWithMod; this.modFile = modFile; @@ -914,10 +929,14 @@ private PipeTsFileInsertionEventResource( @Override protected void finalizeResource() { try { + final String pipeTsFileResourcePipeName = + PipeTsFileResourceManager.getPipeTsFileResourcePipeName(pipeName, creationTime); // decrease reference count - PipeDataNodeResourceManager.tsfile().decreaseFileReference(tsFile, pipeName); + PipeDataNodeResourceManager.tsfile() + .decreaseFileReference(tsFile, pipeTsFileResourcePipeName); if (isWithMod) { - PipeDataNodeResourceManager.tsfile().decreaseFileReference(modFile, pipeName); + PipeDataNodeResourceManager.tsfile() + .decreaseFileReference(modFile, pipeTsFileResourcePipeName); } // close data container diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner.java index 009daea6c243b..859734efe262a 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner.java @@ -36,6 +36,7 @@ import java.nio.file.attribute.BasicFileAttributes; import java.util.ArrayList; import java.util.List; +import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; @@ -48,6 +49,9 @@ public class PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner { "PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner#cleanTsFileDir()"; private static final long DELETE_MAX_PATH_COUNT_PER_ROUND = 100_000L; private static final long DELETE_MAX_TIME_PER_ROUND_MS = 1_000L; + private static final PeriodicalStalePipeDirCleaner PERIODICAL_STALE_PIPE_DIR_CLEANER = + new PeriodicalStalePipeDirCleaner(); + private static final AtomicBoolean PERIODICAL_CLEANUP_JOB_REGISTERED = new AtomicBoolean(false); /** * Delete the data directory and all of its subdirectories that contain the @@ -70,7 +74,8 @@ private static void cleanTsFileDir(final PeriodicalJobRegistrar periodicalJobReg moveAsideAndCollect(pipeHardLinkDir, pipeHardlinkBaseDirName, stalePipeDirs); } } - registerPeriodicalCleanupJob(periodicalJobRegistrar, stalePipeDirs); + PERIODICAL_STALE_PIPE_DIR_CLEANER.addStalePipeDirs(stalePipeDirs); + registerPeriodicalCleanupJob(periodicalJobRegistrar); } private static void collectInterruptedStalePipeDirs( @@ -81,7 +86,7 @@ private static void collectInterruptedStalePipeDirs( localDataDir.listFiles( file -> file.isDirectory() - && file.getName().startsWith(pipeHardlinkBaseDirName + STALE_PIPE_DIR_SUFFIX)); + && file.getName().contains(pipeHardlinkBaseDirName + STALE_PIPE_DIR_SUFFIX)); if (stalePipeDirFiles == null) { return; } @@ -137,17 +142,48 @@ private static File moveAside(final File pipeHardLinkDir, final String pipeHardl } private static void registerPeriodicalCleanupJob( - final PeriodicalJobRegistrar periodicalJobRegistrar, final List stalePipeDirs) { - if (stalePipeDirs.isEmpty()) { + final PeriodicalJobRegistrar periodicalJobRegistrar) { + if (!PERIODICAL_CLEANUP_JOB_REGISTERED.compareAndSet(false, true)) { return; } periodicalJobRegistrar.register( PERIODICAL_CLEANUP_JOB_ID, - new PeriodicalStalePipeDirCleaner(stalePipeDirs)::cleanOneRound, + PERIODICAL_STALE_PIPE_DIR_CLEANER::cleanOneRound, PipeConfig.getInstance().getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds()); } + public static void submitStalePipeDirForPeriodicalCleanup(final File stalePipeDir) { + if (stalePipeDir == null) { + return; + } + + final File stalePipeDirToDelete; + if (stalePipeDir.isDirectory()) { + try { + stalePipeDirToDelete = moveAside(stalePipeDir, stalePipeDir.getName()); + LOGGER.info( + "Pipe hardlink dir found, moved it from {} to {} for throttled periodical deletion.", + stalePipeDir, + stalePipeDirToDelete); + } catch (final IOException e) { + LOGGER.warn( + "Failed to move pipe hardlink dir {} for periodical deletion, skip registering the " + + "original dir to avoid deleting files of a recreated pipe.", + stalePipeDir, + e); + return; + } + } else { + stalePipeDirToDelete = stalePipeDir; + } + + LOGGER.info( + "Stale pipe hardlink dir found, registering it for throttled periodical deletion: {}", + stalePipeDirToDelete); + PERIODICAL_STALE_PIPE_DIR_CLEANER.addStalePipeDir(stalePipeDirToDelete); + } + private static CleanupRoundResult deleteQuietlyWithThrottle(final File stalePipeDir) { if (!stalePipeDir.exists()) { return CleanupRoundResult.finished(); @@ -230,33 +266,36 @@ public interface PeriodicalJobRegistrar { private static class PeriodicalStalePipeDirCleaner { - private final List stalePipeDirs; - private int currentDirIndex; - private boolean finished; + private final ConcurrentLinkedQueue stalePipeDirs = new ConcurrentLinkedQueue<>(); + private File currentStalePipeDir; - private PeriodicalStalePipeDirCleaner(final List stalePipeDirs) { - this.stalePipeDirs = stalePipeDirs; - currentDirIndex = 0; - finished = false; + private void addStalePipeDir(final File stalePipeDir) { + stalePipeDirs.offer(stalePipeDir); } - private void cleanOneRound() { - if (finished) { - return; - } + private void addStalePipeDirs(final List stalePipeDirs) { + stalePipeDirs.forEach(this::addStalePipeDir); + } + private void cleanOneRound() { long deletedPathCount = 0; - while (currentDirIndex < stalePipeDirs.size()) { - final File stalePipeDir = stalePipeDirs.get(currentDirIndex); - final CleanupRoundResult result = deleteQuietlyWithThrottle(stalePipeDir); + while (true) { + if (currentStalePipeDir == null) { + currentStalePipeDir = stalePipeDirs.poll(); + } + if (currentStalePipeDir == null) { + return; + } + + final CleanupRoundResult result = deleteQuietlyWithThrottle(currentStalePipeDir); deletedPathCount += result.deletedPathCount; if (result.finished) { LOGGER.info( "Finished deleting stale pipe hardlink dir {} by periodical job, result: {}", - stalePipeDir, + currentStalePipeDir, result.success); - ++currentDirIndex; + currentStalePipeDir = null; continue; } @@ -265,14 +304,11 @@ private void cleanOneRound() { "Periodically deleted {} paths from stale pipe hardlink dirs, current dir: {}, " + "current round result: {}", deletedPathCount, - stalePipeDir, + currentStalePipeDir, result.success); } return; } - - finished = true; - LOGGER.info("Finished deleting all stale pipe hardlink dirs by periodical job."); } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResource.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResource.java index 8b37f87709447..6329ff9b849a5 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResource.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResource.java @@ -68,9 +68,15 @@ public void increaseReferenceCount() { } public boolean decreaseReferenceCount() { + return decreaseReferenceCount(true); + } + + public boolean decreaseReferenceCount(final boolean deleteFileWhenNoReference) { final int finalReferenceCount = referenceCount.addAndGet(-1); if (finalReferenceCount == 0) { - close(); + if (deleteFileWhenNoReference) { + close(); + } return true; } if (finalReferenceCount < 0) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResourceManager.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResourceManager.java index 2bc3ebf12dc59..3afacdca8d9bc 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResourceManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResourceManager.java @@ -23,6 +23,7 @@ import org.apache.iotdb.commons.pipe.config.PipeConfig; import org.apache.iotdb.commons.utils.FileUtils; import org.apache.iotdb.commons.utils.TestOnly; +import org.apache.iotdb.db.pipe.resource.PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner; import org.apache.iotdb.db.storageengine.dataregion.modification.ModificationFile; import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource; @@ -39,6 +40,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Set; import java.util.concurrent.ConcurrentHashMap; public class PipeTsFileResourceManager { @@ -54,8 +56,16 @@ public class PipeTsFileResourceManager { // PipeName -> TsFilePath -> PipeTsFileResource private final Map> hardlinkOrCopiedFileToPipeTsFileResourceMap = new ConcurrentHashMap<>(); + private final Map pipeNameToPipeTsFileDirPathMap = new ConcurrentHashMap<>(); + private final Set pipeTsFileResourcePipeNameSetUnderDeletion = + ConcurrentHashMap.newKeySet(); private final PipeTsFileResourceSegmentLock segmentLock = new PipeTsFileResourceSegmentLock(); + public static String getPipeTsFileResourcePipeName( + final @Nullable String pipeName, final long creationTime) { + return Objects.isNull(pipeName) ? null : pipeName + "_" + creationTime; + } + public File increaseFileReference( final File file, final boolean isTsFile, final @Nullable String pipeName) throws IOException { return increaseFileReference(file, isTsFile, pipeName, null); @@ -121,6 +131,8 @@ private File increaseFileReference( // file in pipe dir, create a hardlink or copy it to pipe dir, maintain a reference count for // the hardlink or copied file, and return the hardlink or copied file. if (Objects.nonNull(pipeName)) { + pipeNameToPipeTsFileDirPathMap.putIfAbsent( + pipeName, hardlinkOrCopiedFile.getParentFile().getPath()); hardlinkOrCopiedFileToPipeTsFileResourceMap .computeIfAbsent(pipeName, k -> new ConcurrentHashMap<>()) .put(resultFile.getPath(), new PipeTsFileResource(resultFile)); @@ -220,7 +232,8 @@ public void decreaseFileReference( try { final String filePath = hardlinkOrCopiedFile.getPath(); final PipeTsFileResource resource = getResourceMap(pipeName).get(filePath); - if (resource != null && resource.decreaseReferenceCount()) { + if (resource != null + && resource.decreaseReferenceCount(shouldDeleteFileWhenNoReference(pipeName))) { getResourceMap(pipeName).remove(filePath); } } finally { @@ -241,6 +254,35 @@ private void decreasePublicReferenceIfExists(final File file, final @Nullable St decreaseFileReference(new File(getCommonFilePath(file)), null); } + private boolean shouldDeleteFileWhenNoReference( + final @Nullable String pipeTsFileResourcePipeName) { + return Objects.isNull(pipeTsFileResourcePipeName) + || !pipeTsFileResourcePipeNameSetUnderDeletion.contains(pipeTsFileResourcePipeName); + } + + public void markPipeTsFileDirUnderDeletion(final @Nonnull String pipeTsFileResourcePipeName) { + pipeTsFileResourcePipeNameSetUnderDeletion.add(pipeTsFileResourcePipeName); + } + + public void unmarkPipeTsFileDirUnderDeletion(final @Nonnull String pipeTsFileResourcePipeName) { + pipeTsFileResourcePipeNameSetUnderDeletion.remove(pipeTsFileResourcePipeName); + } + + public void cleanPipeTsFileDir(final @Nonnull String pipeTsFileResourcePipeName) { + final String pipeTsFileDirPath = + pipeNameToPipeTsFileDirPathMap.remove(pipeTsFileResourcePipeName); + hardlinkOrCopiedFileToPipeTsFileResourceMap.remove(pipeTsFileResourcePipeName); + + if (Objects.isNull(pipeTsFileDirPath)) { + pipeTsFileResourcePipeNameSetUnderDeletion.remove(pipeTsFileResourcePipeName); + return; + } + + PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner.submitStalePipeDirForPeriodicalCleanup( + new File(pipeTsFileDirPath)); + pipeTsFileResourcePipeNameSetUnderDeletion.remove(pipeTsFileResourcePipeName); + } + // Warning: Shall not be called by the assigner private String getCommonFilePath(final @Nonnull File file) { // If the parent or grandparent is null then this is testing scenario