This is an automated email from the ASF dual-hosted git repository.

smengcl pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git


The following commit(s) were added to refs/heads/master by this push:
     new 6384c643fb3 HDDS-14742. Clean up orphan files from a terminated 
compaction. (#10821)
6384c643fb3 is described below

commit 6384c643fb318bb33d86cb0afdac4e8e65d0864e
Author: SaketaChalamchala <[email protected]>
AuthorDate: Tue Aug 11 14:45:27 2026 -0700

    HDDS-14742. Clean up orphan files from a terminated compaction. (#10821)
---
 .../ozone/rocksdiff/RocksDBCheckpointDiffer.java   | 48 ++++++++++++++++++-
 .../rocksdiff/TestRocksDBCheckpointDiffer.java     | 55 ++++++++++++++++++----
 2 files changed, 91 insertions(+), 12 deletions(-)

diff --git 
a/hadoop-hdds/rocksdb-checkpoint-differ/src/main/java/org/apache/ozone/rocksdiff/RocksDBCheckpointDiffer.java
 
b/hadoop-hdds/rocksdb-checkpoint-differ/src/main/java/org/apache/ozone/rocksdiff/RocksDBCheckpointDiffer.java
index 956a0caac7c..470722e3bf5 100644
--- 
a/hadoop-hdds/rocksdb-checkpoint-differ/src/main/java/org/apache/ozone/rocksdiff/RocksDBCheckpointDiffer.java
+++ 
b/hadoop-hdds/rocksdb-checkpoint-differ/src/main/java/org/apache/ozone/rocksdiff/RocksDBCheckpointDiffer.java
@@ -692,14 +692,16 @@ public void loadAllCompactionLogs() {
     synchronized (this) {
       preconditionChecksForLoadAllCompactionLogs();
       addEntriesFromLogFilesToDagAndCompactionLogTable();
-      loadCompactionDagFromDB();
+      Set<String> referencedInputSstFiles = loadCompactionDagFromDB();
+      cleanupOrphanedSstBackupFiles(referencedInputSstFiles);
     }
   }
 
   /**
    * Read a compactionLofTable and create entries in the dags.
    */
-  private void loadCompactionDagFromDB() {
+  private Set<String> loadCompactionDagFromDB() {
+    Set<String> inputSstFilesFromCompactionLog = new HashSet<>();
     try (ManagedRocksIterator managedRocksIterator = new ManagedRocksIterator(
         activeRocksDB.get().newIterator(compactionLogTableCFHandle))) {
       managedRocksIterator.get().seekToFirst();
@@ -707,6 +709,8 @@ private void loadCompactionDagFromDB() {
         byte[] value = managedRocksIterator.get().value();
         CompactionLogEntry compactionLogEntry =
             
CompactionLogEntry.getFromProtobuf(CompactionLogEntryProto.parseFrom(value));
+        compactionLogEntry.getInputFileInfoList()
+            .forEach(inputFile -> 
inputSstFilesFromCompactionLog.add(inputFile.getFileName()));
         
compactionDag.populateCompactionDAG(compactionLogEntry.getInputFileInfoList(),
             compactionLogEntry.getOutputFileInfoList(), 
compactionLogEntry.getDbSequenceNumber());
         // Add the compaction log entry to the prune queue so that the backup 
input sst files can be pruned.
@@ -722,6 +726,7 @@ private void loadCompactionDagFromDB() {
         sstFilePruningMetrics.updateQueueSize(pruneQueue.size());
       }
     }
+    return inputSstFilesFromCompactionLog;
   }
 
   private void preconditionChecksForLoadAllCompactionLogs() {
@@ -1161,6 +1166,45 @@ private synchronized void 
removeKeyFromCompactionLogTable(
     }
   }
 
+  /**
+   * Removes SST files from the backup directory that were hard-linked during
+   * {@code onCompactionBegin} but never recorded in the compaction log because
+   * the compaction did not complete (for example, due to an OM crash/restart).
+   */
+  @VisibleForTesting
+  void cleanupOrphanedSstBackupFiles(Set<String> referencedInputSstFiles) {
+    Path sstBackupDirPath = Paths.get(sstBackupDir);
+    if (!Files.isDirectory(sstBackupDirPath)) {
+      return;
+    }
+
+    Set<String> orphanedFiles = new HashSet<>();
+    try (Stream<Path> pathStream = Files.list(sstBackupDirPath)) {
+      pathStream.filter(path -> 
path.getFileName().toString().endsWith(SST_FILE_EXTENSION))
+          .forEach(path -> {
+            String fileName = 
FilenameUtils.getBaseName(path.getFileName().toString());
+            if (!referencedInputSstFiles.contains(fileName)) {
+              orphanedFiles.add(fileName);
+            }
+          });
+    } catch (IOException e) {
+      LOG.warn("Failed to list SST backup directory " + sstBackupDir, e);
+    }
+
+    if (orphanedFiles.isEmpty()) {
+      return;
+    }
+
+    LOG.info("Removing orphaned SST backup files left by incomplete 
compactions: {}",
+        orphanedFiles);
+    try (UncheckedAutoCloseable ignored = 
getBootstrapStateLock().acquireReadLock()) {
+      removeSstFiles(orphanedFiles);
+    } catch (InterruptedException e) {
+      LOG.warn("Failed to remove orphaned SST backup files", e);
+      Thread.currentThread().interrupt();
+    }
+  }
+
   /**
    * Deletes the SST files from the backup directory if exists.
    */
diff --git 
a/hadoop-hdds/rocksdb-checkpoint-differ/src/test/java/org/apache/ozone/rocksdiff/TestRocksDBCheckpointDiffer.java
 
b/hadoop-hdds/rocksdb-checkpoint-differ/src/test/java/org/apache/ozone/rocksdiff/TestRocksDBCheckpointDiffer.java
index 0fc2df2a596..5c98e190669 100644
--- 
a/hadoop-hdds/rocksdb-checkpoint-differ/src/test/java/org/apache/ozone/rocksdiff/TestRocksDBCheckpointDiffer.java
+++ 
b/hadoop-hdds/rocksdb-checkpoint-differ/src/test/java/org/apache/ozone/rocksdiff/TestRocksDBCheckpointDiffer.java
@@ -1368,8 +1368,7 @@ private static Stream<Arguments> 
sstFilePruningScenarios() {
     List<String> initialFiles1 = Arrays.asList("000015", "000013", "000011",
         "000009");
     List<String> initialFiles2 = Arrays.asList("000015", "000013", "000011",
-        "000009", "000018", "000016", "000017", "000026", "000024", "000022",
-        "000020");
+        "000009", "000018", "000016", "000017");
     List<String> initialFiles3 = Arrays.asList("000015", "000013", "000011",
         "000009", "000018", "000016", "000017", "000026", "000024", "000022",
         "000020", "000027", "000030", "000028", "000031", "000029", "000039",
@@ -1377,21 +1376,19 @@ private static Stream<Arguments> 
sstFilePruningScenarios() {
         "000046", "000041", "000045", "000054", "000052", "000050", "000048",
         "000059", "000055", "000056", "000060", "000057", "000058");
 
-    List<String> expectedFiles1 = Arrays.asList("000015", "000013", "000011",
-        "000009");
     List<String> expectedFiles2 = Arrays.asList("000015", "000013", "000011",
-        "000009", "000026", "000024", "000022", "000020");
+        "000009");
     List<String> expectedFiles3 = Arrays.asList("000013", "000024", "000035",
         "000011", "000022", "000033", "000039", "000015", "000026", "000037",
         "000048", "000009", "000050", "000054", "000020", "000052");
 
     return Stream.of(
         Arguments.of("Case 1 with compaction log file: " +
-                "No compaction.",
+                "No compaction; orphan backup SST files removed on load.",
             "",
             null,
             initialFiles1,
-            expectedFiles1
+            Collections.emptyList()
         ),
         Arguments.of("Case 2 with compaction log file: " +
                 "One level compaction.",
@@ -1415,11 +1412,11 @@ private static Stream<Arguments> 
sstFilePruningScenarios() {
             expectedFiles3
         ),
         Arguments.of("Case 4 with compaction log table: " +
-                "No compaction.",
+                "No compaction; orphan backup SST files removed on load.",
             null,
             Collections.emptyList(),
             initialFiles1,
-            expectedFiles1
+            Collections.emptyList()
         ),
         Arguments.of("Case 5 with compaction log table: " +
                 "One level compaction.",
@@ -1490,7 +1487,7 @@ private static List<CompactionFileInfo> 
toFileInfoList(List<String> files,
   }
 
   /**
-   * End-to-end test for SST file pruning.
+   * End-to-end test for SST file pruning after compaction log load and orphan 
cleanup.
    */
   @ParameterizedTest(name = "{0}")
   @MethodSource("sstFilePruningScenarios")
@@ -1557,6 +1554,44 @@ private void createFileWithContext(String fileName, 
String context)
     }
   }
 
+  @Test
+  public void testCleanupOrphanedSstBackupFiles() throws IOException {
+    CompactionLogEntry compactionLogEntry = new CompactionLogEntry(178, 
System.currentTimeMillis(),
+        Collections.singletonList(
+            new CompactionFileInfo("000078", "/volume/bucket1/key-1", 
"/volume/bucket2/key-5", "keyTable")),
+        Collections.singletonList(
+            new CompactionFileInfo("000081", "/volume/bucket1/key-1", 
"/volume/bucket2/key-10", "keyTable")),
+        null
+    );
+    rocksDBCheckpointDiffer.addToCompactionLogTable(compactionLogEntry);
+
+    createFileWithContext(sstBackUpDir + "/000078" + SST_FILE_EXTENSION, 
"tracked");
+    // 000081 is a logged output that simulates an input hard-linked by a 
subsequent compaction
+    // that terminated before its log entry was written.
+    createFileWithContext(sstBackUpDir + "/000081" + SST_FILE_EXTENSION, 
"logged-output");
+    createFileWithContext(sstBackUpDir + "/000099" + SST_FILE_EXTENSION, 
"orphan");
+
+    rocksDBCheckpointDiffer.loadAllCompactionLogs();
+
+    assertTrue(Files.exists(sstBackUpDir.toPath().resolve("000078" + 
SST_FILE_EXTENSION)));
+    assertFalse(Files.exists(sstBackUpDir.toPath().resolve("000081" + 
SST_FILE_EXTENSION)));
+    assertFalse(Files.exists(sstBackUpDir.toPath().resolve("000099" + 
SST_FILE_EXTENSION)));
+  }
+
+  @Test
+  public void testPruneSstFilesRetainsBackupFilesWhenCompactionDagIsEmpty() 
throws IOException {
+    List<String> backupFiles = Arrays.asList("000015", "000013", "000011", 
"000009");
+    for (String fileName : backupFiles) {
+      createFileWithContext(sstBackUpDir + "/" + fileName + 
SST_FILE_EXTENSION, fileName);
+    }
+
+    rocksDBCheckpointDiffer.pruneSstFiles();
+
+    for (String fileName : backupFiles) {
+      assertTrue(Files.exists(sstBackUpDir.toPath().resolve(fileName + 
SST_FILE_EXTENSION)));
+    }
+  }
+
   /**
    * Test cases for testGetSSTDiffListWithoutDB.
    */


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to