From ea2ab0ef013757257feebe79c23c2cb1ea54c369 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Fri, 31 Jul 2026 17:40:26 +0800 Subject: [PATCH] Pipe: Fix concurrent TsFile reference increases --- .../tsfile/PipeTsFileResourceManager.java | 37 ++++++----- .../PipeTsFileResourceManagerTest.java | 62 +++++++++++++++++++ 2 files changed, 83 insertions(+), 16 deletions(-) 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 5bcdaec14a11f..17b15cb26ea08 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 @@ -122,23 +122,28 @@ private File increaseFileReference( segmentLock.lock(hardlinkOrCopiedFile); try { - resultFile = - isTsFile - ? FileUtils.createHardLink(source, hardlinkOrCopiedFile) - : FileUtils.copyFile(source, hardlinkOrCopiedFile); - - // If the file is not a hardlink or copied file, and there is no related hardlink or copied - // 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)); + final PipeTsFileResource existingResource = + getResourceMap(pipeName).get(hardlinkOrCopiedFile.getPath()); + if (existingResource != null) { + existingResource.increaseReferenceCount(); + resultFile = existingResource.getFile(); } else { - hardlinkOrCopiedFileToTsFilePublicResourceMap.put( - resultFile.getPath(), new PipeTsFilePublicResource(resultFile)); + resultFile = + isTsFile + ? FileUtils.createHardLink(source, hardlinkOrCopiedFile) + : FileUtils.copyFile(source, hardlinkOrCopiedFile); + + // Create the hardlink or copy and its reference-counted resource only when none exists. + if (Objects.nonNull(pipeName)) { + pipeNameToPipeTsFileDirPathMap.putIfAbsent( + pipeName, hardlinkOrCopiedFile.getParentFile().getPath()); + hardlinkOrCopiedFileToPipeTsFileResourceMap + .computeIfAbsent(pipeName, k -> new ConcurrentHashMap<>()) + .put(resultFile.getPath(), new PipeTsFileResource(resultFile)); + } else { + hardlinkOrCopiedFileToTsFilePublicResourceMap.put( + resultFile.getPath(), new PipeTsFilePublicResource(resultFile)); + } } } finally { segmentLock.unlock(hardlinkOrCopiedFile); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/PipeTsFileResourceManagerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/PipeTsFileResourceManagerTest.java index 0c69684f25fbd..68cd29223f61f 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/PipeTsFileResourceManagerTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/PipeTsFileResourceManagerTest.java @@ -47,6 +47,13 @@ import java.io.File; import java.io.IOException; import java.nio.file.Files; +import java.util.ArrayList; +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; import static org.junit.Assert.fail; @@ -238,4 +245,59 @@ public void testDecreaseTsFile() throws IOException { Assert.assertFalse(Files.exists(originFile.toPath())); Assert.assertFalse(Files.exists(originModFile.toPath())); } + + @Test + public void testConcurrentIncreaseTsFile() throws Exception { + final int concurrency = 64; + final File originTsFile = new File(TS_FILE_NAME); + final CountDownLatch readyLatch = new CountDownLatch(concurrency); + final CountDownLatch startLatch = new CountDownLatch(1); + final ExecutorService executor = Executors.newFixedThreadPool(concurrency); + final List> futures = new ArrayList<>(concurrency); + + try { + for (int i = 0; i < concurrency; i++) { + futures.add( + executor.submit( + () -> { + readyLatch.countDown(); + startLatch.await(); + return pipeTsFileResourceManager.increaseFileReference( + originTsFile, true, PIPE_NAME); + })); + } + + Assert.assertTrue(readyLatch.await(30, TimeUnit.SECONDS)); + startLatch.countDown(); + + File pipeTsFile = null; + for (final Future future : futures) { + final File referencedFile = future.get(30, TimeUnit.SECONDS); + if (pipeTsFile == null) { + pipeTsFile = referencedFile; + } else { + Assert.assertEquals(pipeTsFile, referencedFile); + } + } + + Assert.assertNotNull(pipeTsFile); + Assert.assertEquals( + concurrency, pipeTsFileResourceManager.getFileReferenceCount(pipeTsFile, PIPE_NAME)); + Assert.assertEquals( + concurrency, pipeTsFileResourceManager.getFileReferenceCount(pipeTsFile, null)); + Assert.assertTrue(Files.exists(pipeTsFile.toPath())); + + for (int i = 0; i < concurrency; i++) { + pipeTsFileResourceManager.decreaseFileReference(pipeTsFile, PIPE_NAME); + } + Assert.assertEquals( + 0, pipeTsFileResourceManager.getFileReferenceCount(pipeTsFile, PIPE_NAME)); + Assert.assertEquals(0, pipeTsFileResourceManager.getFileReferenceCount(pipeTsFile, null)); + Assert.assertFalse(Files.exists(pipeTsFile.toPath())); + } finally { + startLatch.countDown(); + executor.shutdownNow(); + Assert.assertTrue(executor.awaitTermination(30, TimeUnit.SECONDS)); + } + } }