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

Reply via email to