This is an automated email from the ASF dual-hosted git repository.
jt2594838 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 011616824b9 Pipe: Fix TsFile reference races and improve result set
diffs (#18376)
011616824b9 is described below
commit 011616824b9a6bf602fb7361be77062b1d754532
Author: Caideyipi <[email protected]>
AuthorDate: Mon Aug 31 09:35:14 2026 +0800
Pipe: Fix TsFile reference races and improve result set diffs (#18376)
* Pipe: Fix concurrent TsFile reference increases
* Test: show concise result set diffs
* Pipe: Roll back failed TsFile reference increases
* Clarify nullable pipe resource map lookup
---
.../org/apache/iotdb/db/it/utils/TestUtils.java | 54 +++++++++++++-
.../resource/tsfile/PipeTsFileResourceManager.java | 57 +++++++++-----
.../resource/PipeTsFileResourceManagerTest.java | 87 ++++++++++++++++++++++
3 files changed, 178 insertions(+), 20 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 11cc0654f90..bf1cb4376bb 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 @@ import static org.junit.Assert.fail;
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 class TestUtils {
System.out.println(builder);
}
}
- assertEquals(expectedResult, actualRetSet);
+ if (expectedResult instanceof Set) {
+ assertStringSetEqual((Set<String>) expectedResult, (Set<String>)
actualRetSet);
+ } else {
+ assertEquals(expectedResult, actualRetSet);
+ }
} catch (final Exception e) {
e.printStackTrace();
Assert.fail(String.valueOf(e));
@@ -826,13 +831,58 @@ public class TestUtils {
}
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<String> expectedResult, final Set<String> actualResult) {
+ if (expectedResult.equals(actualResult)) {
+ return;
+ }
+
+ final List<String> missingRows = new ArrayList<>(expectedResult);
+ missingRows.removeAll(actualResult);
+ Collections.sort(missingRows);
+
+ final List<String> 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<String> rows) {
+ diff.append(System.lineSeparator()).append(title).append("
(").append(rows.size()).append("):");
+ if (rows.isEmpty()) {
+ diff.append(" <none>");
+ 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<String, String> expectedHeaderWithResult)
{
try {
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 5bcdaec14a1..e325e4f0cf8 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,28 +122,39 @@ public class PipeTsFileResourceManager {
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);
}
- 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;
}
@@ -228,6 +239,13 @@ public class PipeTsFileResourceManager {
*/
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();
@@ -242,7 +260,9 @@ public class PipeTsFileResourceManager {
// 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) {
@@ -382,6 +402,7 @@ public class PipeTsFileResourceManager {
}
}
+ /** Returns the shared public resource map when {@code pipeName} is null. */
public Map<String, ? extends PipeTsFileResource> getResourceMap(final
@Nullable String pipeName) {
return Objects.nonNull(pipeName)
? hardlinkOrCopiedFileToPipeTsFileResourceMap.computeIfAbsent(
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 0c69684f25f..d4b8c45f0f2 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 org.junit.Test;
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,84 @@ public class PipeTsFileResourceManagerTest {
Assert.assertFalse(Files.exists(originFile.toPath()));
Assert.assertFalse(Files.exists(originModFile.toPath()));
}
+
+ @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 CountDownLatch readyLatch = new CountDownLatch(concurrency);
+ final CountDownLatch startLatch = new CountDownLatch(1);
+ final ExecutorService executor = Executors.newFixedThreadPool(concurrency);
+ final List<Future<File>> futures = new ArrayList<>(concurrency);
+
+ try {
+ for (int i = 0; i < concurrency; i++) {
+ futures.add(
+ executor.submit(
+ () -> {
+ readyLatch.countDown();
+ startLatch.await();
+ return pipeTsFileResourceManager.increaseFileReference(
+ originFile, isTsFile, PIPE_NAME);
+ }));
+ }
+
+ Assert.assertTrue(readyLatch.await(30, TimeUnit.SECONDS));
+ startLatch.countDown();
+
+ File pipeFile = null;
+ for (final Future<File> future : futures) {
+ final File referencedFile = future.get(30, TimeUnit.SECONDS);
+ if (pipeFile == null) {
+ pipeFile = referencedFile;
+ } else {
+ Assert.assertEquals(pipeFile, referencedFile);
+ }
+ }
+
+ Assert.assertNotNull(pipeFile);
+ Assert.assertEquals(
+ concurrency,
pipeTsFileResourceManager.getFileReferenceCount(pipeFile, PIPE_NAME));
+ Assert.assertEquals(
+ concurrency,
pipeTsFileResourceManager.getFileReferenceCount(pipeFile, null));
+ Assert.assertTrue(Files.exists(pipeFile.toPath()));
+
+ for (int i = 0; i < concurrency; i++) {
+ pipeTsFileResourceManager.decreaseFileReference(pipeFile, PIPE_NAME);
+ }
+ 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();
+ Assert.assertTrue(executor.awaitTermination(30, TimeUnit.SECONDS));
+ }
+ }
}