Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}
Expand All @@ -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;
Expand All @@ -331,6 +348,10 @@ protected boolean dropPipe(final String pipeName) {
}

if (!super.dropPipe(pipeName)) {
if (Objects.nonNull(pipeTsFileResourcePipeName)) {
PipeDataNodeResourceManager.tsfile()
.unmarkPipeTsFileDirUnderDeletion(pipeTsFileResourcePipeName);
}
return false;
}

Expand All @@ -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)) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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;
Expand Down Expand Up @@ -461,7 +470,9 @@ public boolean mayEventPathsOverlappedWithPattern() {
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();
Expand Down Expand Up @@ -877,6 +888,7 @@ public PipeEventResource eventResourceBuilder() {
this.isReleased,
this.referenceCount,
this.pipeName,
this.creationTime,
this.tsFile,
this.isWithMod,
this.modFile,
Expand All @@ -891,19 +903,22 @@ private static class PipeTsFileInsertionEventResource extends PipeEventResource
private final File modFile;
private final AtomicReference<TsFileInsertionDataContainer> 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,
final AtomicReference<TsFileInsertionDataContainer> dataContainer,
final AtomicBoolean isTsFileParserMemoryReserved) {
super(isReleased, referenceCount);
this.pipeName = pipeName;
this.creationTime = creationTime;
this.tsFile = tsFile;
this.isWithMod = isWithMod;
this.modFile = modFile;
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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
Expand All @@ -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(
Expand All @@ -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;
}
Expand Down Expand Up @@ -137,17 +142,48 @@ private static File moveAside(final File pipeHardLinkDir, final String pipeHardl
}

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();
Expand Down Expand Up @@ -230,33 +266,36 @@ public interface PeriodicalJobRegistrar {

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;
}

Expand All @@ -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.");
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Loading
Loading