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 b4a4227cf4f HDDS-14876. Make snapshot diff progress more granular 
(#10901)
b4a4227cf4f is described below

commit b4a4227cf4fa8701493bf3ce5babdb3346981d56
Author: rhalm <[email protected]>
AuthorDate: Tue Aug 11 23:45:57 2026 +0200

    HDDS-14876. Make snapshot diff progress more granular (#10901)
---
 .../ozone/shell/snapshot/SnapshotDiffHandler.java  |  7 +++
 .../hadoop/ozone/om/helpers/SnapshotDiffJob.java   |  3 +-
 ...OzoneManagerProtocolClientSideTranslatorPB.java | 11 +++-
 .../ozone/snapshot/SnapshotDiffResponse.java       | 28 +++++++--
 .../ozone/snapshot/TestSnapshotDiffResponse.java   | 40 +++++++++++--
 .../src/main/proto/OmClientProtocol.proto          |  6 +-
 .../interface-client/src/main/resources/proto.lock | 18 ++++++
 .../ozone/om/snapshot/SnapshotDiffManager.java     | 66 +++++++++++++++++-----
 .../diff/delta/CompositeDeltaDiffComputer.java     |  3 +
 .../protocolPB/OzoneManagerRequestHandler.java     |  6 ++
 .../ozone/om/snapshot/TestSnapshotDiffManager.java | 31 +++++++++-
 .../om/snapshot/TestSnapshotDiffManagerMXBean.java |  4 +-
 .../protocolPB/TestOzoneManagerRequestHandler.java | 65 +++++++++++++++++++++
 13 files changed, 253 insertions(+), 35 deletions(-)

diff --git 
a/hadoop-ozone/cli-shell/src/main/java/org/apache/hadoop/ozone/shell/snapshot/SnapshotDiffHandler.java
 
b/hadoop-ozone/cli-shell/src/main/java/org/apache/hadoop/ozone/shell/snapshot/SnapshotDiffHandler.java
index a4fdbdfb579..4e0943e5dd1 100644
--- 
a/hadoop-ozone/cli-shell/src/main/java/org/apache/hadoop/ozone/shell/snapshot/SnapshotDiffHandler.java
+++ 
b/hadoop-ozone/cli-shell/src/main/java/org/apache/hadoop/ozone/shell/snapshot/SnapshotDiffHandler.java
@@ -181,6 +181,13 @@ private ObjectNode getJsonObject(SnapshotDiffResponse 
diffResponse) {
     if (StringUtils.isNotEmpty(diffResponse.getReason())) {
       diffResponseNode.put("reason", diffResponse.getReason());
     }
+    if (diffResponse.getSubStatus() != null) {
+      SnapshotDiffResponse.SubStatus sub = diffResponse.getSubStatus();
+      diffResponseNode.put("subStatus", sub.name());
+      if (sub.hasProgress()) {
+        diffResponseNode.put("progressPercent", 
diffResponse.getProgressPercent());
+      }
+    }
     return diffResponseNode;
   }
 
diff --git 
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/SnapshotDiffJob.java
 
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/SnapshotDiffJob.java
index baf304e898c..387bf8c6b21 100644
--- 
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/SnapshotDiffJob.java
+++ 
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/SnapshotDiffJob.java
@@ -229,8 +229,7 @@ public String toString() {
     }
     if (status.equals(JobStatus.IN_PROGRESS) && subStatus != null) {
       sb.append(", subStatus: ").append(subStatus);
-      if (subStatus.equals(SubStatus.OBJECT_ID_MAP_GEN_FSO) ||
-          subStatus.equals(SubStatus.OBJECT_ID_MAP_GEN_OBS)) {
+      if (subStatus.hasProgress()) {
         sb.append(String.format(", keysProcessedPct: %.2f", keysProcessedPct));
       }
     }
diff --git 
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java
 
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java
index 3a2007d774a..ed1d5c35785 100644
--- 
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java
+++ 
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java
@@ -1470,12 +1470,19 @@ private SnapshotDiffResponse 
snapshotDiffInternal(String volumeName,
     OzoneManagerProtocolProtos.SnapshotDiffResponse diffResponse =
         omResponse.getSnapshotDiffResponse();
 
-    return new SnapshotDiffResponse(SnapshotDiffReportOzone.fromProtobuf(
-        diffResponse.getSnapshotDiffReport()),
+    SnapshotDiffResponse result = new SnapshotDiffResponse(
+        
SnapshotDiffReportOzone.fromProtobuf(diffResponse.getSnapshotDiffReport()),
         JobStatus.fromProtobuf(diffResponse.getJobStatus()),
         diffResponse.getWaitTimeInMs(),
         diffResponse.getReason(),
         reportOnly);
+    if (diffResponse.hasSubStatus()) {
+      
result.setSubStatus(SnapshotDiffResponse.SubStatus.fromProtoBuf(diffResponse.getSubStatus()));
+      if (diffResponse.hasProgressPercent()) {
+        result.setProgressPercent(diffResponse.getProgressPercent());
+      }
+    }
+    return result;
   }
 
   /**
diff --git 
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/snapshot/SnapshotDiffResponse.java
 
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/snapshot/SnapshotDiffResponse.java
index de16a3aec40..816bfb0ce9f 100644
--- 
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/snapshot/SnapshotDiffResponse.java
+++ 
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/snapshot/SnapshotDiffResponse.java
@@ -86,10 +86,18 @@ public String getReason() {
     return reason;
   }
 
+  public SubStatus getSubStatus() {
+    return subStatus;
+  }
+
   public void setSubStatus(SubStatus subStatus) {
     this.subStatus = subStatus;
   }
 
+  public double getProgressPercent() {
+    return progressPercent;
+  }
+
   public void setProgressPercent(double progressPercent) {
     this.progressPercent = progressPercent;
   }
@@ -143,11 +151,10 @@ public String toString() {
       str.append(".\n");
       if (subStatus != null) {
         str.append("SubStatus : ")
-            .append(subStatus);
-        if (subStatus.equals(SubStatus.OBJECT_ID_MAP_GEN_OBS) ||
-            subStatus.equals(SubStatus.OBJECT_ID_MAP_GEN_FSO)) {
-          str.append("Keys Processed Estimated Percentage : ")
-              .append(progressPercent);
+            .append(subStatus)
+            .append('\n');
+        if (subStatus.hasProgress()) {
+          str.append(String.format("Keys Processed Estimated Percentage : 
%.1f%n", progressPercent));
         }
       }
     }
@@ -182,8 +189,17 @@ public enum SubStatus {
     SST_FILE_DELTA_DAG_WALK,
     SST_FILE_DELTA_FULL_DIFF,
     OBJECT_ID_MAP_GEN_OBS,
+    @Deprecated
     OBJECT_ID_MAP_GEN_FSO,
-    DIFF_REPORT_GEN;
+    DIFF_REPORT_GEN,
+    PATH_RESOLUTION_FSO,
+    OBJECT_ID_MAP_GEN_FSO_FILE,
+    OBJECT_ID_MAP_GEN_FSO_DIR;
+
+    public boolean hasProgress() {
+      return this == OBJECT_ID_MAP_GEN_OBS || this == OBJECT_ID_MAP_GEN_FSO || 
this == OBJECT_ID_MAP_GEN_FSO_FILE
+          || this == OBJECT_ID_MAP_GEN_FSO_DIR;
+    }
 
     public static SubStatus 
fromProtoBuf(OzoneManagerProtocolProtos.SnapshotDiffResponse.SubStatus 
subStatusProto) {
       return SubStatus.valueOf(subStatusProto.name());
diff --git 
a/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/snapshot/TestSnapshotDiffResponse.java
 
b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/snapshot/TestSnapshotDiffResponse.java
index 74b7d73d3e2..14125911c15 100644
--- 
a/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/snapshot/TestSnapshotDiffResponse.java
+++ 
b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/snapshot/TestSnapshotDiffResponse.java
@@ -17,12 +17,16 @@
 
 package org.apache.hadoop.ozone.snapshot;
 
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 
 import java.util.Collections;
 import org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.JobStatus;
 import org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.SubStatus;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
+import org.junit.jupiter.params.provider.EnumSource.Mode;
 
 class TestSnapshotDiffResponse {
 
@@ -57,20 +61,48 @@ void testReportOnlyFailedMessageIncludesReason() {
     assertTrue(message.contains("resubmit the job without using the 
--get-report option"));
   }
 
-  @Test
-  void testReportOnlyInProgressIncludesSubStatusAndProgress() {
+  @ParameterizedTest
+  @EnumSource(value = SubStatus.class, names = {"OBJECT_ID_MAP_GEN_OBS", 
"OBJECT_ID_MAP_GEN_FSO",
+      "OBJECT_ID_MAP_GEN_FSO_FILE", "OBJECT_ID_MAP_GEN_FSO_DIR"})
+  void testInProgressWithMapGenSubStatusIncludesProgress(SubStatus subStatus) {
     SnapshotDiffResponse response = new SnapshotDiffResponse(createReport(),
         JobStatus.IN_PROGRESS, 1000L, true);
-    response.setSubStatus(SubStatus.OBJECT_ID_MAP_GEN_OBS);
+    response.setSubStatus(subStatus);
     response.setProgressPercent(55.5);
 
     String message = response.toString();
     assertTrue(message.contains("IN_PROGRESS"));
-    assertTrue(message.contains("OBJECT_ID_MAP_GEN_OBS"));
+    assertTrue(message.contains(subStatus.name()));
     assertTrue(message.contains("Keys Processed Estimated Percentage"));
     assertTrue(message.contains("55.5"));
   }
 
+  @ParameterizedTest
+  @EnumSource(value = SubStatus.class,
+      names = {"OBJECT_ID_MAP_GEN_OBS", "OBJECT_ID_MAP_GEN_FSO", 
"OBJECT_ID_MAP_GEN_FSO_FILE",
+          "OBJECT_ID_MAP_GEN_FSO_DIR"},
+      mode = Mode.EXCLUDE)
+  void 
testInProgressWithNonMapGenSubStatusRendersSubStatusButNotProgress(SubStatus 
subStatus) {
+    SnapshotDiffResponse response = new SnapshotDiffResponse(createReport(), 
JobStatus.IN_PROGRESS, 1000L, true);
+    response.setSubStatus(subStatus);
+    response.setProgressPercent(55.5);
+
+    String message = response.toString();
+    assertTrue(message.contains("IN_PROGRESS"));
+    assertTrue(message.contains(subStatus.name()));
+    assertFalse(message.contains("Keys Processed Estimated Percentage"));
+    assertFalse(message.contains("55.5"));
+  }
+
+  @Test
+  void testInProgressWithNullSubStatusOmitsSubStatusAndProgressLines() {
+    SnapshotDiffResponse response = new SnapshotDiffResponse(createReport(), 
JobStatus.IN_PROGRESS, 1000L, true);
+    String message = response.toString();
+    assertTrue(message.contains("IN_PROGRESS"));
+    assertFalse(message.contains("SubStatus"));
+    assertFalse(message.contains("Keys Processed Estimated Percentage"));
+  }
+
   private SnapshotDiffReportOzone createReport() {
     return new SnapshotDiffReportOzone("snapshotRoot", "vol", "bucket", 
"fromSnap",
         "toSnap", Collections.emptyList(), null);
diff --git 
a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto 
b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
index bccc6d11f11..2a27af60754 100644
--- a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
+++ b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
@@ -2254,8 +2254,11 @@ message SnapshotDiffResponse {
       SST_FILE_DELTA_DAG_WALK = 1;
       SST_FILE_DELTA_FULL_DIFF = 2;
       OBJECT_ID_MAP_GEN_OBS = 3;
-      OBJECT_ID_MAP_GEN_FSO = 4;
+      OBJECT_ID_MAP_GEN_FSO = 4 [deprecated = true];
       DIFF_REPORT_GEN = 5;
+      PATH_RESOLUTION_FSO = 6;
+      OBJECT_ID_MAP_GEN_FSO_FILE = 7;
+      OBJECT_ID_MAP_GEN_FSO_DIR = 8;
   }
 
   optional SnapshotDiffReportProto snapshotDiffReport = 1;
@@ -2263,6 +2266,7 @@ message SnapshotDiffResponse {
   optional int64 waitTimeInMs = 3;
   optional string reason = 4;
   optional SubStatus subStatus = 5;
+  optional double progressPercent = 6;
 }
 
 message SubmitSnapshotDiffResponse {
diff --git a/hadoop-ozone/interface-client/src/main/resources/proto.lock 
b/hadoop-ozone/interface-client/src/main/resources/proto.lock
index 5f6b5806361..88107a50557 100644
--- a/hadoop-ozone/interface-client/src/main/resources/proto.lock
+++ b/hadoop-ozone/interface-client/src/main/resources/proto.lock
@@ -1342,6 +1342,18 @@
               {
                 "name": "DIFF_REPORT_GEN",
                 "integer": 5
+              },
+              {
+                "name": "PATH_RESOLUTION_FSO",
+                "integer": 6
+              },
+              {
+                "name": "OBJECT_ID_MAP_GEN_FSO_FILE",
+                "integer": 7
+              },
+              {
+                "name": "OBJECT_ID_MAP_GEN_FSO_DIR",
+                "integer": 8
               }
             ]
           },
@@ -8103,6 +8115,12 @@
                 "name": "subStatus",
                 "type": "SubStatus",
                 "optional": true
+              },
+              {
+                "id": 6,
+                "name": "progressPercent",
+                "type": "double",
+                "optional": true
               }
             ]
           },
diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/snapshot/SnapshotDiffManager.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/snapshot/SnapshotDiffManager.java
index 9909f5792d7..d120973eaef 100644
--- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/snapshot/SnapshotDiffManager.java
+++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/snapshot/SnapshotDiffManager.java
@@ -41,6 +41,7 @@
 import static 
org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SNAPSHOT_FORCE_FULL_DIFF_DEFAULT;
 import static org.apache.hadoop.ozone.om.OmSnapshotManager.DELIMITER;
 import static org.apache.hadoop.ozone.om.codec.OMDBDefinition.DIRECTORY_TABLE;
+import static org.apache.hadoop.ozone.om.codec.OMDBDefinition.FILE_TABLE;
 import static 
org.apache.hadoop.ozone.om.snapshot.SnapshotUtils.checkSnapshotActive;
 import static 
org.apache.hadoop.ozone.om.snapshot.SnapshotUtils.dropColumnFamilyHandle;
 import static 
org.apache.hadoop.ozone.om.snapshot.SnapshotUtils.getSnapshotInfo;
@@ -60,8 +61,10 @@
 import static 
org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.JobStatus.QUEUED;
 import static 
org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.JobStatus.REJECTED;
 import static 
org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.SubStatus.DIFF_REPORT_GEN;
-import static 
org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.SubStatus.OBJECT_ID_MAP_GEN_FSO;
+import static 
org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.SubStatus.OBJECT_ID_MAP_GEN_FSO_DIR;
+import static 
org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.SubStatus.OBJECT_ID_MAP_GEN_FSO_FILE;
 import static 
org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.SubStatus.OBJECT_ID_MAP_GEN_OBS;
+import static 
org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.SubStatus.PATH_RESOLUTION_FSO;
 
 import com.google.common.annotations.VisibleForTesting;
 import com.google.common.collect.ImmutableMap;
@@ -143,6 +146,7 @@
 import org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.SubStatus;
 import org.apache.hadoop.ozone.snapshot.SubmitSnapshotDiffResponse;
 import org.apache.hadoop.ozone.util.ClosableIterator;
+import org.apache.hadoop.util.Time;
 import org.apache.logging.log4j.util.Strings;
 import org.apache.ozone.rocksdb.util.SstFileInfo;
 import org.apache.ozone.rocksdiff.RocksDBCheckpointDiffer;
@@ -1114,27 +1118,27 @@ void generateSnapshotDiffReport(final String jobKey,
       // repetition while constantly checking if the job is cancelled.
       Callable<Void>[] methodCalls = new Callable[]{
           () -> {
-            recordActivity(jobKey, OBJECT_ID_MAP_GEN_OBS);
             getDeltaFilesAndDiffKeysToObjectIdToKeyMap(fsKeyTable, tsKeyTable,
                 fsInfo, tsInfo, performNonNativeDiff, tablePrefixes,
                 objectIdToKeyNameMapForFromSnapshot,
                 objectIdToKeyNameMapForToSnapshot, objectIdToIsDirMap,
-                oldParentIds, newParentIds, deltaFileComputer, jobKey);
+                oldParentIds, newParentIds, deltaFileComputer, jobKey, jobId);
             return null;
           },
           () -> {
             if (bucketLayout.isFileSystemOptimized()) {
-              recordActivity(jobKey, OBJECT_ID_MAP_GEN_FSO);
               getDeltaFilesAndDiffKeysToObjectIdToKeyMap(fsDirTable, 
tsDirTable,
                   fsInfo, tsInfo, performNonNativeDiff, tablePrefixes,
                   objectIdToKeyNameMapForFromSnapshot,
                   objectIdToKeyNameMapForToSnapshot, objectIdToIsDirMap,
-                  oldParentIds, newParentIds, deltaFileComputer, jobKey);
+                  oldParentIds, newParentIds, deltaFileComputer, jobKey, 
jobId);
             }
             return null;
           },
           () -> {
             if (bucketLayout.isFileSystemOptimized()) {
+              recordActivity(jobKey, PATH_RESOLUTION_FSO);
+              long pathResolutionStart = Time.monotonicNow();
               long bucketId = toSnapshot.getMetadataManager()
                   .getBucketId(volumeName, bucketName);
               String tablePrefix = 
tablePrefixes.getTablePrefix(fromSnapshot.getMetadataManager()
@@ -1147,11 +1151,19 @@ void generateSnapshotDiffReport(final String jobKey,
                   tablePrefix, bucketId,
                   toSnapshot.getMetadataManager().getDirectoryTable())
                   .getAbsolutePathForObjectIDs(newParentIds, true));
+              if (LOG.isDebugEnabled()) {
+                LOG.debug("Completed FSO path resolution for snapshot diff, 
resolved {} out of {} parent IDs, " +
+                        "elapsed: {}ms, jobId: {}",
+                    oldParentIdPathMap.get().size() + 
newParentIdPathMap.get().size(),
+                    oldParentIds.get().size() + newParentIds.get().size(),
+                    Time.monotonicNow() - pathResolutionStart, jobId);
+              }
             }
             return null;
           },
           () -> {
             recordActivity(jobKey, DIFF_REPORT_GEN);
+            long reportGenStart = Time.monotonicNow();
             Pair<Long, String> reportEntries = generateDiffReport(jobId,
                 fsKeyTable,
                 tsKeyTable,
@@ -1168,6 +1180,10 @@ void generateSnapshotDiffReport(final String jobKey,
             if (reportEntries.getKey() >= 0 &&
                 areDiffJobAndSnapshotsActive(volumeName, bucketName,
                     fromSnapshotName, toSnapshotName)) {
+              if (LOG.isDebugEnabled()) {
+                LOG.debug("Generated snapshot diff report, entry count: {}, 
elapsed: {}ms, jobId: {}",
+                    reportEntries.getKey(), Time.monotonicNow() - 
reportGenStart, jobId);
+              }
               updateJobStatusToDone(jobKey, reportEntries.getKey(), 
reportEntries.getValue());
             }
             return null;
@@ -1223,17 +1239,24 @@ private void getDeltaFilesAndDiffKeysToObjectIdToKeyMap(
       final PersistentMap<byte[], byte[]> newObjIdToKeyMap,
       final PersistentMap<byte[], Boolean> objectIdToIsDirMap,
       final Optional<Set<Long>> oldParentIds, final Optional<Set<Long>> 
newParentIds,
-      final DeltaFileComputer deltaFileComputer, final String jobKey) throws 
IOException, RocksDBException {
+      final DeltaFileComputer deltaFileComputer, final String jobKey,
+      final String jobId) throws IOException, RocksDBException {
 
+    long deltaFilesStart = Time.monotonicNow();
     Set<String> tablesToLookUp = Collections.singleton(fsTable.getName());
     Collection<Pair<Path, SstFileInfo>> deltaFiles = 
deltaFileComputer.getDeltaFiles(fsInfo, tsInfo,
         tablesToLookUp);
     if (LOG.isDebugEnabled()) {
-      LOG.debug("Computed Delta SST File Set, Total count = {} ", 
deltaFiles.size());
+      LOG.debug("Computed Delta SST File Set for table '{}', file count: {}, 
elapsed: {}ms, jobId: {}",
+          fsTable.getName(), deltaFiles.size(), Time.monotonicNow() - 
deltaFilesStart, jobId);
     }
+    recordActivity(jobKey,
+        fsTable.getName().equals(FILE_TABLE) ? OBJECT_ID_MAP_GEN_FSO_FILE
+            : fsTable.getName().equals(DIRECTORY_TABLE) ? 
OBJECT_ID_MAP_GEN_FSO_DIR
+            : OBJECT_ID_MAP_GEN_OBS);
     addToObjectIdMap(fsTable, tsTable, 
deltaFiles.stream().map(Pair::getLeft).collect(Collectors.toList()),
         !skipNativeDiff, oldObjIdToKeyMap, newObjIdToKeyMap, 
objectIdToIsDirMap, oldParentIds,
-        newParentIds, tablePrefixes, jobKey);
+        newParentIds, tablePrefixes, jobKey, jobId);
   }
 
   @VisibleForTesting
@@ -1246,10 +1269,19 @@ void addToObjectIdMap(Table<String, ? extends 
WithParentObjectId> fsTable,
       PersistentMap<byte[], Boolean> objectIdToIsDirMap,
       Optional<Set<Long>> oldParentIds,
       Optional<Set<Long>> newParentIds,
-      TablePrefixInfo tablePrefixes, String jobKey) throws IOException, 
RocksDBException {
+      TablePrefixInfo tablePrefixes, String jobKey,
+      String jobId) throws IOException, RocksDBException {
     if (deltaFiles.isEmpty()) {
+      updateProgress(jobKey, 1.0);
+      if (LOG.isDebugEnabled()) {
+        LOG.debug("Skipped object ID map generation for table '{}' because 
there are no delta files, jobId: {}",
+            fsTable.getName(), jobId);
+      }
       return;
     }
+    updateProgress(jobKey, 0.0);
+    long objectIdMapStart = Time.monotonicNow();
+    AtomicLong keysProcessed = new AtomicLong(0);
     String tablePrefix = tablePrefixes.getTablePrefix(fsTable.getName());
     boolean isDirectoryTable = fsTable.getName().equals(DIRECTORY_TABLE);
     SstFileSetReader sstFileReader = new SstFileSetReader(deltaFiles);
@@ -1258,8 +1290,7 @@ void addToObjectIdMap(Table<String, ? extends 
WithParentObjectId> fsTable,
     String sstFileReaderLowerBound = tablePrefix;
     String sstFileReaderUpperBound = null;
     double stepIncreasePct = 0.1;
-    double[] checkpoint = new double[1];
-    checkpoint[0] = stepIncreasePct;
+    double checkpoint = stepIncreasePct;
     if (Strings.isNotEmpty(tablePrefix)) {
       sstFileReaderUpperBound = getLexicographicallyHigherString(tablePrefix);
     }
@@ -1268,15 +1299,14 @@ void addToObjectIdMap(Table<String, ? extends 
WithParentObjectId> fsTable,
         : sstFileReader.getKeyStream(sstFileReaderLowerBound, 
sstFileReaderUpperBound);
          TableMergeIterator<String, WithParentObjectId> tableMergeIterator = 
new TableMergeIterator<>(keysToCheck,
              tablePrefix, (Table<String, WithParentObjectId>) fsTable, 
(Table<String, WithParentObjectId>) tsTable)) {
-      AtomicLong keysProcessed = new AtomicLong(0);
       while (tableMergeIterator.hasNext()) {
         Table.KeyValue<String, List<WithParentObjectId>> kvs = 
tableMergeIterator.next();
         String key = kvs.getKey();
         if (totalEstimatedKeysToProcess > 0) {
           double progressPct = (double) keysProcessed.get() / 
totalEstimatedKeysToProcess;
-          if (progressPct >= checkpoint[0]) {
+          if (progressPct >= checkpoint) {
             updateProgress(jobKey, progressPct);
-            checkpoint[0] += stepIncreasePct;
+            checkpoint += stepIncreasePct;
           }
         }
 
@@ -1314,6 +1344,11 @@ void addToObjectIdMap(Table<String, ? extends 
WithParentObjectId> fsTable,
         }
       }
     }
+    updateProgress(jobKey, 1.0);
+    if (LOG.isDebugEnabled()) {
+      LOG.debug("Generated object ID map for table '{}', keys scanned: {}, 
elapsed: {}ms, jobId: {}",
+          fsTable.getName(), keysProcessed.get(), Time.monotonicNow() - 
objectIdMapStart, jobId);
+    }
   }
 
   private void validateEstimatedKeyChangesAreInLimits(
@@ -1683,7 +1718,8 @@ synchronized void updateProgress(String jobKey,
     snapshotDiffJob.setKeysProcessedPct(pct * 100);
     snapDiffJobTable.put(jobKey, snapshotDiffJob);
     if (LOG.isDebugEnabled()) {
-      LOG.debug("Completed processing {}% of keys for snapshot diff job {}", 
pct, jobKey);
+      LOG.debug("Completed processing {}% of keys for snapshot diff job {}",
+          snapshotDiffJob.getKeysProcessedPct(), jobKey);
     }
   }
 
diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/snapshot/diff/delta/CompositeDeltaDiffComputer.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/snapshot/diff/delta/CompositeDeltaDiffComputer.java
index 4ef17d84114..6c08a88db32 100644
--- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/snapshot/diff/delta/CompositeDeltaDiffComputer.java
+++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/snapshot/diff/delta/CompositeDeltaDiffComputer.java
@@ -86,6 +86,9 @@ Optional<Map<Path, Pair<Path, SstFileInfo>>> 
computeDeltaFiles(SnapshotInfo from
         updateActivity(SnapshotDiffResponse.SubStatus.SST_FILE_DELTA_DAG_WALK);
         deltaFiles = differComputer.computeDeltaFiles(fromSnapshotInfo, 
toSnapshotInfo, tablesToLookup,
             tablePrefixInfo).orElse(null);
+        if (deltaFiles == null) {
+          LOG.warn("DAG diff returned no result for tables {}, falling back to 
full diff.", tablesToLookup);
+        }
       }
     } catch (Exception e) {
       LOG.warn("Falling back to full diff.", e);
diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OzoneManagerRequestHandler.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OzoneManagerRequestHandler.java
index 701ecb6b38e..b585fa234c8 100644
--- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OzoneManagerRequestHandler.java
+++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OzoneManagerRequestHandler.java
@@ -1482,6 +1482,12 @@ private SnapshotDiffResponse snapshotDiff(
       builder.setSnapshotDiffReport(
           response.getSnapshotDiffReport().toProtobuf());
     }
+    if (response.getSubStatus() != null) {
+      builder.setSubStatus(response.getSubStatus().toProtoBuf());
+      if (response.getSubStatus().hasProgress()) {
+        builder.setProgressPercent(response.getProgressPercent());
+      }
+    }
 
     return builder.build();
   }
diff --git 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/snapshot/TestSnapshotDiffManager.java
 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/snapshot/TestSnapshotDiffManager.java
index 69507a6d749..14fc35756e6 100644
--- 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/snapshot/TestSnapshotDiffManager.java
+++ 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/snapshot/TestSnapshotDiffManager.java
@@ -158,6 +158,7 @@
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.Arguments;
 import org.junit.jupiter.params.provider.CsvSource;
+import org.junit.jupiter.params.provider.EnumSource;
 import org.junit.jupiter.params.provider.MethodSource;
 import org.mockito.Mock;
 import org.mockito.MockedConstruction;
@@ -494,12 +495,16 @@ public void testObjectIdMapWithTombstoneEntries(boolean 
nativeLibraryLoaded,
       Set<Long> oldParentIds = Sets.newHashSet();
       Set<Long> newParentIds = Sets.newHashSet();
 
+      SnapshotDiffJob dummyJob = new 
SnapshotDiffJob(System.currentTimeMillis(),
+          "", IN_PROGRESS, VOLUME_NAME, BUCKET_NAME, "from", "to", false, 
false, 0, null, 0.0, "");
+      db.get().put(snapDiffJobTable, codecRegistry.asRawData(""), 
codecRegistry.asRawData(dummyJob));
+
       snapshotDiffManager.addToObjectIdMap(toSnapshotTable,
           fromSnapshotTable, Sets.newHashSet(Paths.get("dummy.sst")),
           nativeLibraryLoaded, oldObjectIdKeyMap, newObjectIdKeyMap,
           objectIdsToCheck, Optional.of(oldParentIds),
           Optional.of(newParentIds),
-          new TablePrefixInfo(ImmutableMap.of(DIRECTORY_TABLE, "0", KEY_TABLE, 
"0", FILE_TABLE, "0")), "");
+          new TablePrefixInfo(ImmutableMap.of(DIRECTORY_TABLE, "0", KEY_TABLE, 
"0", FILE_TABLE, "0")), "", "");
 
       try (ClosableIterator<Map.Entry<byte[], byte[]>> oldObjectIdIter =
                oldObjectIdKeyMap.iterator()) {
@@ -1780,18 +1785,38 @@ private static Map<Long, Long> 
buildDirATreeParentMap(long bucketObjectId) {
         .build();
   }
 
+  @ParameterizedTest
+  @EnumSource(value = SnapshotDiffResponse.SubStatus.class,
+      names = {"OBJECT_ID_MAP_GEN_OBS", "OBJECT_ID_MAP_GEN_FSO", 
"OBJECT_ID_MAP_GEN_FSO_DIR",
+          "OBJECT_ID_MAP_GEN_FSO_FILE"})
+  public void 
testGetSnapshotDiffReportReportOnlyInProgressIncludesProgressDetails(
+      SnapshotDiffResponse.SubStatus subStatus) throws IOException {
+    SnapDiffTestContext ctx = setupRandomSnapDiffTestContext();
+    SnapshotDiffJob existing = new SnapshotDiffJob(0L, 
UUID.randomUUID().toString(),
+        IN_PROGRESS, ctx.volumeName, ctx.bucketName, ctx.fromSnapshotName, 
ctx.toSnapshotName,
+        false, false, 0L, subStatus, 55.5, null);
+    snapshotDiffManager.getSnapDiffJobTable().put(ctx.diffJobKey, existing);
+
+    SnapshotDiffResponse response = snapshotDiffManager.getSnapshotDiffReport(
+        ctx.volumeName, ctx.bucketName, ctx.fromSnapshotName, 
ctx.toSnapshotName, "", 1000);
+    assertEquals(IN_PROGRESS, response.getJobStatus());
+    assertEquals(subStatus, response.getSubStatus());
+    assertThat(response.getProgressPercent()).isEqualTo(55.5);
+  }
+
   @Test
-  public void 
testGetSnapshotDiffReportReportOnlyInProgressIncludesProgressDetails()
+  public void testGetSnapshotDiffReportInProgressWithPathResolutionSubStatus()
       throws IOException {
     SnapDiffTestContext ctx = setupRandomSnapDiffTestContext();
     SnapshotDiffJob existing = new SnapshotDiffJob(0L, 
UUID.randomUUID().toString(),
         IN_PROGRESS, ctx.volumeName, ctx.bucketName, ctx.fromSnapshotName, 
ctx.toSnapshotName,
-        false, false, 0L, 
SnapshotDiffResponse.SubStatus.OBJECT_ID_MAP_GEN_OBS, 55.5, null);
+        false, false, 0L, SnapshotDiffResponse.SubStatus.PATH_RESOLUTION_FSO, 
0.0, null);
     snapshotDiffManager.getSnapDiffJobTable().put(ctx.diffJobKey, existing);
 
     SnapshotDiffResponse response = snapshotDiffManager.getSnapshotDiffReport(
         ctx.volumeName, ctx.bucketName, ctx.fromSnapshotName, 
ctx.toSnapshotName, "", 1000);
     assertEquals(IN_PROGRESS, response.getJobStatus());
+    assertEquals(SnapshotDiffResponse.SubStatus.PATH_RESOLUTION_FSO, 
response.getSubStatus());
   }
 
   private static OmDirectoryInfo newDeletedDir(long objectId, long parentId) {
diff --git 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/snapshot/TestSnapshotDiffManagerMXBean.java
 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/snapshot/TestSnapshotDiffManagerMXBean.java
index afd75551322..10ab715594d 100644
--- 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/snapshot/TestSnapshotDiffManagerMXBean.java
+++ 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/snapshot/TestSnapshotDiffManagerMXBean.java
@@ -139,7 +139,7 @@ public void testMXBeanRegistration() throws Exception {
         org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.JobStatus.QUEUED,
         "vol", "bucket", "snap1", "snap2",
         false, false, 0,
-        
org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.SubStatus.OBJECT_ID_MAP_GEN_FSO,
+        
org.apache.hadoop.ozone.snapshot.SnapshotDiffResponse.SubStatus.OBJECT_ID_MAP_GEN_FSO_DIR,
         0, null);
     snapDiffJobTable.put("job-1", job);
 
@@ -151,6 +151,6 @@ public void testMXBeanRegistration() throws Exception {
     assertEquals("job-1", jobs[0].get("jobId"));
     assertEquals("snap1", jobs[0].get("fromSnapshot"));
     assertEquals("snap2", jobs[0].get("toSnapshot"));
-    assertEquals("OBJECT_ID_MAP_GEN_FSO", jobs[0].get("subStatus"));
+    assertEquals("OBJECT_ID_MAP_GEN_FSO_DIR", jobs[0].get("subStatus"));
   }
 }
diff --git 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOzoneManagerRequestHandler.java
 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOzoneManagerRequestHandler.java
index 35ee959236d..40e9e18e9bc 100644
--- 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOzoneManagerRequestHandler.java
+++ 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOzoneManagerRequestHandler.java
@@ -481,4 +481,69 @@ public void 
testSnapshotDiffRoutingUsesCorrectServerMethodBasedOnOptionalFlags()
         reportOnlyResponse.getSnapshotDiffResponse().getJobStatus());
     Assertions.assertEquals(0L, 
reportOnlyResponse.getSnapshotDiffResponse().getWaitTimeInMs());
   }
+
+  @Test
+  public void testSnapshotDiffHandlerSerializesSubStatusAndProgressIntoProto() 
throws IOException {
+    OzoneManagerRequestHandler handler = getRequestHandler(10);
+    OzoneManager ozoneManager = handler.getOzoneManager();
+
+    OMLayoutVersionManager lvm = Mockito.mock(OMLayoutVersionManager.class);
+    Mockito.when(lvm.isAllowed(Mockito.anyString())).thenReturn(true);
+    Mockito.when(ozoneManager.getVersionManager()).thenReturn(lvm);
+
+    SnapshotDiffResponse diffResponse =
+        new SnapshotDiffResponse(null, 
SnapshotDiffResponse.JobStatus.IN_PROGRESS, 60000L);
+    
diffResponse.setSubStatus(SnapshotDiffResponse.SubStatus.OBJECT_ID_MAP_GEN_FSO_DIR);
+    diffResponse.setProgressPercent(50.0);
+    Mockito.when(ozoneManager.snapshotDiff(Mockito.anyString(), 
Mockito.anyString(),
+            Mockito.anyString(), Mockito.anyString(), Mockito.anyString(), 
Mockito.anyInt()))
+        .thenReturn(diffResponse);
+
+    OzoneManagerProtocolProtos.OMRequest request =
+        OzoneManagerProtocolProtos.OMRequest.newBuilder()
+            .setCmdType(OzoneManagerProtocolProtos.Type.SnapshotDiff)
+            .setClientId("client")
+            
.setSnapshotDiffRequest(OzoneManagerProtocolProtos.SnapshotDiffRequest.newBuilder()
+                .setVolumeName("vol").setBucketName("buck")
+                .setFromSnapshot("s1").setToSnapshot("s2")
+                .setToken("t").setPageSize(10))
+            .build();
+
+    OzoneManagerProtocolProtos.SnapshotDiffResponse proto =
+        handler.handleReadRequest(request).getSnapshotDiffResponse();
+    Assertions.assertTrue(proto.hasSubStatus());
+    Assertions.assertEquals(
+        
OzoneManagerProtocolProtos.SnapshotDiffResponse.SubStatus.OBJECT_ID_MAP_GEN_FSO_DIR,
+        proto.getSubStatus());
+    Assertions.assertEquals(50.0, proto.getProgressPercent(), 1e-9);
+  }
+
+  @Test
+  public void testSnapshotDiffHandlerOmitsSubStatusFromProtoWhenNull() throws 
IOException {
+    OzoneManagerRequestHandler handler = getRequestHandler(10);
+    OzoneManager ozoneManager = handler.getOzoneManager();
+
+    OMLayoutVersionManager lvm = Mockito.mock(OMLayoutVersionManager.class);
+    Mockito.when(lvm.isAllowed(Mockito.anyString())).thenReturn(true);
+    Mockito.when(ozoneManager.getVersionManager()).thenReturn(lvm);
+
+    Mockito.when(ozoneManager.snapshotDiff(Mockito.anyString(), 
Mockito.anyString(),
+            Mockito.anyString(), Mockito.anyString(), Mockito.anyString(), 
Mockito.anyInt()))
+        .thenReturn(new SnapshotDiffResponse(null, 
SnapshotDiffResponse.JobStatus.IN_PROGRESS, 60000L));
+
+    OzoneManagerProtocolProtos.OMRequest request =
+        OzoneManagerProtocolProtos.OMRequest.newBuilder()
+            .setCmdType(OzoneManagerProtocolProtos.Type.SnapshotDiff)
+            .setClientId("client")
+            
.setSnapshotDiffRequest(OzoneManagerProtocolProtos.SnapshotDiffRequest.newBuilder()
+                .setVolumeName("vol").setBucketName("buck")
+                .setFromSnapshot("s1").setToSnapshot("s2")
+                .setToken("t").setPageSize(10))
+            .build();
+
+    OzoneManagerProtocolProtos.SnapshotDiffResponse proto =
+        handler.handleReadRequest(request).getSnapshotDiffResponse();
+    Assertions.assertFalse(proto.hasSubStatus());
+    Assertions.assertFalse(proto.hasProgressPercent());
+  }
 }


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

Reply via email to