SaketaChalamchala commented on code in PR #8214:
URL: https://github.com/apache/ozone/pull/8214#discussion_r2080756018
##########
hadoop-hdds/rocksdb-checkpoint-differ/src/main/java/org/apache/ozone/rocksdiff/RocksDBCheckpointDiffer.java:
##########
@@ -1200,6 +1249,109 @@ public void pruneSstFiles() {
}
}
+ /**
+ * Defines the task that removes OMKeyInfo from SST files from backup
directory to
+ * save disk space.
+ */
+ public void pruneSstFileValues() {
+ if (!shouldRun()) {
+ return;
+ }
+
+ Path sstBackupDirPath = Paths.get(sstBackupDir);
+ Path prunedSSTFilePath = sstBackupDirPath.resolve(PRUNED_SST_FILE_TEMP);
+ try (ManagedOptions managedOptions = new ManagedOptions();
+ ManagedEnvOptions envOptions = new ManagedEnvOptions();
+ ManagedSstFileWriter sstFileWriter = new
ManagedSstFileWriter(envOptions, managedOptions)) {
+ byte[] compactionLogEntryKey;
+ int batchCounter = 0;
+ while ((compactionLogEntryKey = pruneQueue.peek()) != null &&
++batchCounter <= pruneSSTFileBatchSize) {
+ CompactionLogEntry compactionLogEntry;
+ synchronized (this) {
+ try {
+ compactionLogEntry =
CompactionLogEntry.getCodec().fromPersistedFormat(
+ activeRocksDB.get().get(compactionLogTableCFHandle,
compactionLogEntryKey));
+ } catch (RocksDBException ex) {
+ throw new RocksDatabaseException("Failed to get compaction log
entry.", ex);
+ }
+ }
+ boolean shouldUpdateTable = false;
+ List<CompactionFileInfo> fileInfoList =
compactionLogEntry.getInputFileInfoList();
+ List<CompactionFileInfo> updatedFileInfoList = new ArrayList<>();
+ for (CompactionFileInfo fileInfo : fileInfoList) {
+ if (fileInfo.isPruned()) {
+ updatedFileInfoList.add(fileInfo);
+ continue;
+ }
+ Path sstFilePath = sstBackupDirPath.resolve(fileInfo.getFileName() +
ROCKSDB_SST_SUFFIX);
+ if (Files.notExists(sstFilePath)) {
+ LOG.debug("Skipping pruning SST file {} as it does not exist in
backup directory.", sstFilePath);
+ updatedFileInfoList.add(fileInfo);
+ continue;
+ }
+
+ // Write the file.sst => pruned.sst.tmp
+ Files.deleteIfExists(prunedSSTFilePath);
+ try (ManagedRawSSTFileReader<Pair<byte[], Integer>> sstFileReader =
new ManagedRawSSTFileReader<>(
+ managedOptions, sstFilePath.toFile().getAbsolutePath(),
SST_READ_AHEAD_SIZE);
+ ManagedRawSSTFileIterator<Pair<byte[], Integer>> itr =
sstFileReader.newIterator(
+ keyValue -> Pair.of(keyValue.getKey(), keyValue.getType()),
null, null)) {
+ sstFileWriter.open(prunedSSTFilePath.toFile().getAbsolutePath());
+ while (itr.hasNext()) {
+ Pair<byte[], Integer> keyValue = itr.next();
+ if (keyValue.getValue() == 0) {
+ sstFileWriter.delete(keyValue.getKey());
+ } else {
+ sstFileWriter.put(keyValue.getKey(), new byte[0]);
+ }
+ }
+ } catch (RocksDBException ex) {
+ throw new RocksDatabaseException("Failed to write pruned entries
for " + sstFilePath, ex);
+ } finally {
+ try {
+ sstFileWriter.finish();
+ } catch (RocksDBException ex) {
+ throw new RocksDatabaseException("Failed to finish writing to "
+ prunedSSTFilePath, ex);
+ }
+ }
+
+ // Move file.sst.tmp to file.sst and replace existing file atomically
+ try (BootstrapStateHandler.Lock lock =
getBootstrapStateLock().lock()) {
+ Files.move(prunedSSTFilePath, sstFilePath,
+ StandardCopyOption.ATOMIC_MOVE,
StandardCopyOption.REPLACE_EXISTING);
+ }
+ shouldUpdateTable = true;
+ fileInfo.setPruned();
+ updatedFileInfoList.add(fileInfo);
+ LOG.debug("Completed pruning OMKeyInfo from {}", sstFilePath);
+ }
+
+ // Update Compaction Log table. Track keys that need updating.
+ if (shouldUpdateTable) {
Review Comment:
We are processing all the files in a compaction log entry, one file at a
time. Sometimes, the compaction entry might exist but `pruneSstFiles` might
delete the sst files if they are non-leaf nodes. In which case, we will skip
pruning the deleted file.
If all the files in the compaction log are deleted, then we don't have to
update the compaction entry, If some files were deleted and some files are
still present on the disk we prune what files are available, if all files are
present on the disk we prune everything.
In the latter 2 cases, we update the compaction log table.
`shouldUpdateTable` indicates if at least of the source files has been pruned.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]