From a1faf68cbe9ec44bbb553da17abd6e0b430b77b9 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 1/3] 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)); + } + } } From b872de5060a860de4aa725d3d588e482373b0515 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Wed, 5 Aug 2026 15:15:24 +0800 Subject: [PATCH 2/3] Test: show concise result set diffs --- .../apache/iotdb/db/it/utils/TestUtils.java | 54 ++++++++++++++++++- 1 file changed, 52 insertions(+), 2 deletions(-) diff --git a/integration-test/src/test/java/org/apache/iotdb/db/it/utils/TestUtils.java b/integration-test/src/test/java/org/apache/iotdb/db/it/utils/TestUtils.java index 11cc0654f9065..bf1cb4376bb95 100644 --- a/integration-test/src/test/java/org/apache/iotdb/db/it/utils/TestUtils.java +++ b/integration-test/src/test/java/org/apache/iotdb/db/it/utils/TestUtils.java @@ -73,6 +73,7 @@ public class TestUtils { private static final Logger LOGGER = LoggerFactory.getLogger(TestUtils.class); + private static final int MAX_RESULT_SET_DIFF_ROWS = 20; public static final ZoneId DEFAULT_ZONE_ID = ZoneId.ofOffset("UTC", ZoneOffset.of("Z")); @@ -790,7 +791,11 @@ public static void assertResultSetEqual( System.out.println(builder); } } - assertEquals(expectedResult, actualRetSet); + if (expectedResult instanceof Set) { + assertStringSetEqual((Set) expectedResult, (Set) actualRetSet); + } else { + assertEquals(expectedResult, actualRetSet); + } } catch (final Exception e) { e.printStackTrace(); Assert.fail(String.valueOf(e)); @@ -826,13 +831,58 @@ public static void assertResultSetEqual( } actualRetSet.add(builder.toString()); } - assertEquals(expectedRetSet, actualRetSet); + assertStringSetEqual(expectedRetSet, actualRetSet); } catch (Exception e) { e.printStackTrace(); Assert.fail(String.valueOf(e)); } } + private static void assertStringSetEqual( + final Set expectedResult, final Set actualResult) { + if (expectedResult.equals(actualResult)) { + return; + } + + final List missingRows = new ArrayList<>(expectedResult); + missingRows.removeAll(actualResult); + Collections.sort(missingRows); + + final List unexpectedRows = new ArrayList<>(actualResult); + unexpectedRows.removeAll(expectedResult); + Collections.sort(unexpectedRows); + + final StringBuilder diff = + new StringBuilder("Result set mismatch: expected ") + .append(expectedResult.size()) + .append(" rows but got ") + .append(actualResult.size()) + .append(" rows."); + appendResultSetDiff(diff, "Missing rows", missingRows); + appendResultSetDiff(diff, "Unexpected rows", unexpectedRows); + fail(diff.toString()); + } + + private static void appendResultSetDiff( + final StringBuilder diff, final String title, final List rows) { + diff.append(System.lineSeparator()).append(title).append(" (").append(rows.size()).append("):"); + if (rows.isEmpty()) { + diff.append(" "); + return; + } + + final int displayedRowCount = Math.min(rows.size(), MAX_RESULT_SET_DIFF_ROWS); + for (int i = 0; i < displayedRowCount; i++) { + diff.append(System.lineSeparator()).append(" ").append(rows.get(i)); + } + if (rows.size() > displayedRowCount) { + diff.append(System.lineSeparator()) + .append(" ... and ") + .append(rows.size() - displayedRowCount) + .append(" more"); + } + } + public static void assertSingleResultSetEqual( ResultSet actualResultSet, Map expectedHeaderWithResult) { try { From 25057effe59c18cfc9862d2990b6bffdcd25c672 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Wed, 5 Aug 2026 15:40:01 +0800 Subject: [PATCH 3/3] Pipe: Roll back failed TsFile reference increases --- .../tsfile/PipeTsFileResourceManager.java | 19 ++++++- .../PipeTsFileResourceManagerTest.java | 55 ++++++++++++++----- 2 files changed, 57 insertions(+), 17 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 17b15cb26ea08..5da57496ea5ab 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 @@ -148,7 +148,13 @@ private File increaseFileReference( } finally { segmentLock.unlock(hardlinkOrCopiedFile); } - increasePublicReference(resultFile, pipeName, isTsFile); + try { + increasePublicReference(resultFile, pipeName, isTsFile); + } catch (final IOException e) { + // The private reference must not outlive a failed public reference increase. + decreaseFileReference(resultFile, pipeName, false); + throw e; + } return resultFile; } @@ -233,6 +239,13 @@ private static String getRelativeFilePath(File file) { */ public void decreaseFileReference( final File hardlinkOrCopiedFile, final @Nullable String pipeName) { + decreaseFileReference(hardlinkOrCopiedFile, pipeName, true); + } + + private void decreaseFileReference( + final File hardlinkOrCopiedFile, + final @Nullable String pipeName, + final boolean decreasePublicReference) { segmentLock.lock(hardlinkOrCopiedFile); try { final String filePath = hardlinkOrCopiedFile.getPath(); @@ -247,7 +260,9 @@ public void decreaseFileReference( // Decrease the assigner's file to clear hard-link and memory cache // Note that it does not exist for historical files - decreasePublicReferenceIfExists(hardlinkOrCopiedFile, pipeName); + if (decreasePublicReference) { + decreasePublicReferenceIfExists(hardlinkOrCopiedFile, pipeName); + } } private void decreasePublicReferenceIfExists(final File file, final @Nullable String pipeName) { 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 68cd29223f61f..d4b8c45f0f21f 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 @@ -248,8 +248,34 @@ public void testDecreaseTsFile() throws IOException { @Test public void testConcurrentIncreaseTsFile() throws Exception { + assertConcurrentIncreaseFileReference(new File(TS_FILE_NAME), true); + } + + @Test + public void testConcurrentIncreaseCopiedFile() throws Exception { + assertConcurrentIncreaseFileReference(new File(MODS_FILE_NAME), false); + } + + @Test + public void testIncreaseFileReferenceRollsBackOnPublicReferenceFailure() throws Exception { + final File originModFile = new File(MODS_FILE_NAME); + final File pipeModFile = + PipeTsFileResourceManager.getHardlinkOrCopiedFileInPipeDir(originModFile, PIPE_NAME); + final File publicModFile = + new File(pipeModFile.getParentFile().getParentFile(), pipeModFile.getName()); + Assert.assertTrue(publicModFile.mkdirs()); + + Assert.assertThrows( + IOException.class, + () -> pipeTsFileResourceManager.increaseFileReference(originModFile, false, PIPE_NAME)); + + Assert.assertEquals(0, pipeTsFileResourceManager.getFileReferenceCount(pipeModFile, PIPE_NAME)); + Assert.assertFalse(Files.exists(pipeModFile.toPath())); + } + + private void assertConcurrentIncreaseFileReference(final File originFile, final boolean isTsFile) + 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); @@ -263,37 +289,36 @@ public void testConcurrentIncreaseTsFile() throws Exception { readyLatch.countDown(); startLatch.await(); return pipeTsFileResourceManager.increaseFileReference( - originTsFile, true, PIPE_NAME); + originFile, isTsFile, PIPE_NAME); })); } Assert.assertTrue(readyLatch.await(30, TimeUnit.SECONDS)); startLatch.countDown(); - File pipeTsFile = null; + File pipeFile = null; for (final Future future : futures) { final File referencedFile = future.get(30, TimeUnit.SECONDS); - if (pipeTsFile == null) { - pipeTsFile = referencedFile; + if (pipeFile == null) { + pipeFile = referencedFile; } else { - Assert.assertEquals(pipeTsFile, referencedFile); + Assert.assertEquals(pipeFile, referencedFile); } } - Assert.assertNotNull(pipeTsFile); + Assert.assertNotNull(pipeFile); Assert.assertEquals( - concurrency, pipeTsFileResourceManager.getFileReferenceCount(pipeTsFile, PIPE_NAME)); + concurrency, pipeTsFileResourceManager.getFileReferenceCount(pipeFile, PIPE_NAME)); Assert.assertEquals( - concurrency, pipeTsFileResourceManager.getFileReferenceCount(pipeTsFile, null)); - Assert.assertTrue(Files.exists(pipeTsFile.toPath())); + concurrency, pipeTsFileResourceManager.getFileReferenceCount(pipeFile, null)); + Assert.assertTrue(Files.exists(pipeFile.toPath())); for (int i = 0; i < concurrency; i++) { - pipeTsFileResourceManager.decreaseFileReference(pipeTsFile, PIPE_NAME); + pipeTsFileResourceManager.decreaseFileReference(pipeFile, PIPE_NAME); } - Assert.assertEquals( - 0, pipeTsFileResourceManager.getFileReferenceCount(pipeTsFile, PIPE_NAME)); - Assert.assertEquals(0, pipeTsFileResourceManager.getFileReferenceCount(pipeTsFile, null)); - Assert.assertFalse(Files.exists(pipeTsFile.toPath())); + Assert.assertEquals(0, pipeTsFileResourceManager.getFileReferenceCount(pipeFile, PIPE_NAME)); + Assert.assertEquals(0, pipeTsFileResourceManager.getFileReferenceCount(pipeFile, null)); + Assert.assertFalse(Files.exists(pipeFile.toPath())); } finally { startLatch.countDown(); executor.shutdownNow();