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]