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