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 e38636a2cac [To dev/1.3] Optimize pipe TsFile cleanup on drop (#18163) 
(#18303)
e38636a2cac is described below

commit e38636a2caca71de9f1ffe12b47d7d0750e70855
Author: Caideyipi <[email protected]>
AuthorDate: Tue Aug 18 14:22:42 2026 +0800

    [To dev/1.3] Optimize pipe TsFile cleanup on drop (#18163) (#18303)
    
    * 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)
    
    Co-authored-by: Zhenyu Luo <[email protected]>
---
 .../db/pipe/agent/task/PipeDataNodeTaskAgent.java  | 22 ++++++
 .../common/tsfile/PipeTsFileInsertionEvent.java    | 29 ++++++--
 ...aNodeHardlinkOrCopiedFileDirStartupCleaner.java | 86 +++++++++++++++-------
 .../pipe/resource/tsfile/PipeTsFileResource.java   |  8 +-
 .../resource/tsfile/PipeTsFileResourceManager.java | 44 ++++++++++-
 5 files changed, 155 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 d99d5b3cf17..2e4b5e090eb 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
@@ -57,6 +57,7 @@ import 
org.apache.iotdb.db.pipe.metric.overview.PipeDataNodeSinglePipeMetrics;
 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;
@@ -392,13 +393,20 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
 
   @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;
   }
@@ -407,6 +415,15 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
   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;
@@ -419,6 +436,10 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
     }
 
     if (!super.dropPipe(pipeName)) {
+      if (Objects.nonNull(pipeTsFileResourcePipeName)) {
+        PipeDataNodeResourceManager.tsfile()
+            .unmarkPipeTsFileDirUnderDeletion(pipeTsFileResourcePipeName);
+      }
       return false;
     }
 
@@ -427,6 +448,7 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
       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 b6938890cf0..19c845d5015 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
@@ -337,11 +337,16 @@ public class PipeTsFileInsertionEvent extends 
EnrichedEvent
   @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) {
@@ -361,10 +366,14 @@ public class PipeTsFileInsertionEvent extends 
EnrichedEvent
 
   @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;
@@ -488,7 +497,9 @@ public class PipeTsFileInsertionEvent extends EnrichedEvent
           PipeDataNodeResourceManager.tsfile()
               .getDeviceIsAlignedMapFromCache(
                   PipeTsFileResourceManager.getHardlinkOrCopiedFileInPipeDir(
-                      resource.getTsFile(), pipeName),
+                      resource.getTsFile(),
+                      PipeTsFileResourceManager.getPipeTsFileResourcePipeName(
+                          pipeName, creationTime)),
                   false);
       final Set<IDeviceID> deviceSet =
           Objects.nonNull(deviceIsAlignedMap) ? deviceIsAlignedMap.keySet() : 
resource.getDevices();
@@ -940,10 +951,14 @@ public class PipeTsFileInsertionEvent extends 
EnrichedEvent
         PipeDataNodeResourceManager.memory()
             .cancelTsFileParserMemoryReservation(
                 pipeName, creationTime, dataRegionId, 
tsFileParserMemoryReservationKey);
+        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 009daea6c24..859734efe26 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.SimpleFileVisitor;
 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 @@ public class 
PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner {
         moveAsideAndCollect(pipeHardLinkDir, pipeHardlinkBaseDirName, 
stalePipeDirs);
       }
     }
-    registerPeriodicalCleanupJob(periodicalJobRegistrar, stalePipeDirs);
+    PERIODICAL_STALE_PIPE_DIR_CLEANER.addStalePipeDirs(stalePipeDirs);
+    registerPeriodicalCleanupJob(periodicalJobRegistrar);
   }
 
   private static void collectInterruptedStalePipeDirs(
@@ -81,7 +86,7 @@ public class 
PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner {
         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 @@ public class 
PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner {
   }
 
   private static void registerPeriodicalCleanupJob(
-      final PeriodicalJobRegistrar periodicalJobRegistrar, final List<File> 
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 class 
PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner {
 
   private static class PeriodicalStalePipeDirCleaner {
 
-    private final List<File> stalePipeDirs;
-    private int currentDirIndex;
-    private boolean finished;
+    private final ConcurrentLinkedQueue<File> stalePipeDirs = new 
ConcurrentLinkedQueue<>();
+    private File currentStalePipeDir;
 
-    private PeriodicalStalePipeDirCleaner(final List<File> 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<File> 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 @@ public class 
PipeDataNodeHardlinkOrCopiedFileDirStartupCleaner {
               "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 8b37f877094..6329ff9b849 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 class PipeTsFileResource implements AutoCloseable {
   }
 
   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 2bc3ebf12dc..3afacdca8d9 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.conf.IoTDBConstant;
 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.io.IOException;
 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<String, Map<String, PipeTsFileResource>>
       hardlinkOrCopiedFileToPipeTsFileResourceMap = new ConcurrentHashMap<>();
+  private final Map<String, String> pipeNameToPipeTsFileDirPathMap = new 
ConcurrentHashMap<>();
+  private final Set<String> 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 @@ public class PipeTsFileResourceManager {
       // 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 class PipeTsFileResourceManager {
     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 @@ public class PipeTsFileResourceManager {
     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

Reply via email to