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

devmadhuu 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 69b1df11f3f HDDS-15335. Recon: parallelize NSSummaryTask sub-tasks and 
cache OmBucketInfo lookups (#10321)
69b1df11f3f is described below

commit 69b1df11f3f6f7057ddee875b8091021ddb150de
Author: Siyao Meng <[email protected]>
AuthorDate: Thu Jun 4 22:04:01 2026 -0700

    HDDS-15335. Recon: parallelize NSSummaryTask sub-tasks and cache 
OmBucketInfo lookups (#10321)
---
 .../hadoop/ozone/recon/tasks/NSSummaryTask.java    |  80 ++--
 .../recon/tasks/NSSummaryTaskDbEventHandler.java   |  42 ++
 .../ozone/recon/tasks/NSSummaryTaskWithLegacy.java |  20 +-
 .../ozone/recon/tasks/NSSummaryTaskWithOBS.java    |  31 +-
 .../ozone/recon/tasks/TestNSSummaryTask.java       | 442 ++++++++++++++++++++-
 5 files changed, 563 insertions(+), 52 deletions(-)

diff --git 
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTask.java
 
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTask.java
index 139190e4baa..53d07bca6b1 100644
--- 
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTask.java
+++ 
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTask.java
@@ -83,6 +83,16 @@ public class NSSummaryTask implements ReconOmTask {
   private final NSSummaryTaskWithLegacy nsSummaryTaskWithLegacy;
   private final NSSummaryTaskWithOBS nsSummaryTaskWithOBS;
 
+  // Shared executor for the three FSO/Legacy/OBS sub-tasks during process().
+  // The sub-tasks operate on disjoint slices of the event stream (filtered by
+  // table and bucket layout) and write to disjoint NSSummary entries, so they
+  // are safe to run in parallel.
+  private static final ExecutorService SUB_TASK_EXECUTOR =
+      Executors.newFixedThreadPool(3, new ThreadFactoryBuilder()
+          .setNameFormat("NSSummarySubTask-%d")
+          .setDaemon(true)
+          .build());
+
   /**
    * Rebuild state enum to track NSSummary tree rebuild status.
    */
@@ -172,37 +182,27 @@ public String getDescription() {
   @Override
   public TaskResult process(
       OMUpdateEventBatch events, Map<String, Integer> subTaskSeekPosMap) {
-    boolean anyFailure = false; // Track if any bucket fails
     Map<String, Integer> updatedSeekPositions = new HashMap<>();
 
-    // Process FSO bucket
-    Integer bucketSeek = subTaskSeekPosMap.getOrDefault(BucketType.FSO.name(), 
0);
-    Pair<Integer, Boolean> bucketResult = 
nsSummaryTaskWithFSO.processWithFSO(events, bucketSeek);
-    updatedSeekPositions.put(BucketType.FSO.name(), bucketResult.getLeft());
-    if (!bucketResult.getRight()) {
-      LOG.error("processWithFSO failed.");
-      anyFailure = true;
-    }
-
-    // Process Legacy bucket
-    bucketSeek = subTaskSeekPosMap.getOrDefault(BucketType.LEGACY.name(), 0);
-    bucketResult = nsSummaryTaskWithLegacy.processWithLegacy(events, 
bucketSeek);
-    updatedSeekPositions.put(BucketType.LEGACY.name(), bucketResult.getLeft());
-    if (!bucketResult.getRight()) {
-      LOG.error("processWithLegacy failed.");
-      anyFailure = true;
-    }
-
-    // Process OBS bucket
-    bucketSeek = subTaskSeekPosMap.getOrDefault(BucketType.OBS.name(), 0);
-    bucketResult = nsSummaryTaskWithOBS.processWithOBS(events, bucketSeek);
-    updatedSeekPositions.put(BucketType.OBS.name(), bucketResult.getLeft());
-    if (!bucketResult.getRight()) {
-      LOG.error("processWithOBS failed.");
-      anyFailure = true;
-    }
+    int fsoSeek = subTaskSeekPosMap.getOrDefault(BucketType.FSO.name(), 0);
+    int legacySeek = subTaskSeekPosMap.getOrDefault(BucketType.LEGACY.name(), 
0);
+    int obsSeek = subTaskSeekPosMap.getOrDefault(BucketType.OBS.name(), 0);
+
+    Future<Pair<Integer, Boolean>> fsoFuture = SUB_TASK_EXECUTOR.submit(
+        () -> nsSummaryTaskWithFSO.processWithFSO(events, fsoSeek));
+    Future<Pair<Integer, Boolean>> legacyFuture = SUB_TASK_EXECUTOR.submit(
+        () -> nsSummaryTaskWithLegacy.processWithLegacy(events, legacySeek));
+    Future<Pair<Integer, Boolean>> obsFuture = SUB_TASK_EXECUTOR.submit(
+        () -> nsSummaryTaskWithOBS.processWithOBS(events, obsSeek));
+
+    boolean anyFailure = false;
+    anyFailure |= !awaitSubTask("processWithFSO", BucketType.FSO,
+        fsoFuture, fsoSeek, updatedSeekPositions);
+    anyFailure |= !awaitSubTask("processWithLegacy", BucketType.LEGACY,
+        legacyFuture, legacySeek, updatedSeekPositions);
+    anyFailure |= !awaitSubTask("processWithOBS", BucketType.OBS,
+        obsFuture, obsSeek, updatedSeekPositions);
 
-    // Return task failure if any bucket failed, while keeping each bucket's 
latest seek position
     return new TaskResult.Builder()
         .setTaskName(getTaskName())
         .setSubTaskSeekPositions(updatedSeekPositions)
@@ -210,6 +210,30 @@ public TaskResult process(
         .build();
   }
 
+  private boolean awaitSubTask(String name, BucketType type,
+                               Future<Pair<Integer, Boolean>> future,
+                               int fallbackSeek,
+                               Map<String, Integer> updatedSeekPositions) {
+    try {
+      Pair<Integer, Boolean> result = future.get();
+      updatedSeekPositions.put(type.name(), result.getLeft());
+      if (!result.getRight()) {
+        LOG.error("{} failed.", name);
+        return false;
+      }
+      return true;
+    } catch (InterruptedException e) {
+      Thread.currentThread().interrupt();
+      LOG.error("{} interrupted.", name, e);
+      updatedSeekPositions.put(type.name(), fallbackSeek);
+      return false;
+    } catch (ExecutionException e) {
+      LOG.error("{} threw an exception.", name, e.getCause());
+      updatedSeekPositions.put(type.name(), fallbackSeek);
+      return false;
+    }
+  }
+
   @Override
   public TaskResult reprocess(OMMetadataManager omMetadataManager) {
     // Unified control for all NSS tree rebuild operations
diff --git 
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTaskDbEventHandler.java
 
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTaskDbEventHandler.java
index cd0d10c6f9e..d3ddf108f22 100644
--- 
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTaskDbEventHandler.java
+++ 
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTaskDbEventHandler.java
@@ -20,8 +20,10 @@
 import java.io.IOException;
 import java.util.Collection;
 import java.util.Collections;
+import java.util.HashMap;
 import java.util.Map;
 import org.apache.hadoop.hdds.utils.db.RDBBatchOperation;
+import org.apache.hadoop.ozone.om.helpers.OmBucketInfo;
 import org.apache.hadoop.ozone.om.helpers.OmDirectoryInfo;
 import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
 import org.apache.hadoop.ozone.recon.ReconUtils;
@@ -43,6 +45,18 @@ public class NSSummaryTaskDbEventHandler {
   private ReconNamespaceSummaryManager reconNamespaceSummaryManager;
   private ReconOMMetadataManager reconOMMetadataManager;
 
+  // Cache OmBucketInfo lookups across process() calls so the Legacy and OBS
+  // sub-tasks don't pay a RocksDB point read per event. A bucket's objectID 
and
+  // layout are stable while the bucket exists, but a bucket can be deleted and
+  // recreated under the same volume/bucket name with a new objectID (same DB
+  // key, different identity). A recreate is always preceded by a delete, so 
the
+  // sub-tasks call invalidateBucketCache() when they observe a bucketTable
+  // delete event; a recreated bucket is then re-read instead of served stale.
+  //
+  // Single-thread access only (each sub-task runs on its own thread and owns
+  // its own cache instance). HashMap is fine.
+  private final Map<String, OmBucketInfo> bucketInfoCache = new HashMap<>();
+
   public NSSummaryTaskDbEventHandler(ReconNamespaceSummaryManager
                                      reconNamespaceSummaryManager,
                                      ReconOMMetadataManager
@@ -51,6 +65,34 @@ public 
NSSummaryTaskDbEventHandler(ReconNamespaceSummaryManager
     this.reconOMMetadataManager = reconOMMetadataManager;
   }
 
+  /** Look up an {@link OmBucketInfo} via {@code getBucketTable().getSkipCache}
+   *  and cache the result. Bucket layout/object-id are stable while a bucket
+   *  exists, so a field-level cache avoids one RocksDB point read per event in
+   *  the per-event sub-task loops. Entries are dropped via
+   *  {@link #invalidateBucketCache(String)} when a bucketTable delete event is
+   *  seen, so a bucket deleted and recreated under the same name is not served
+   *  stale. */
+  protected OmBucketInfo lookupBucketCached(String bucketDBKey) throws 
IOException {
+    OmBucketInfo cached = bucketInfoCache.get(bucketDBKey);
+    if (cached != null) {
+      return cached;
+    }
+    OmBucketInfo info = 
reconOMMetadataManager.getBucketTable().getSkipCache(bucketDBKey);
+    if (info != null) {
+      bucketInfoCache.put(bucketDBKey, info);
+    }
+    return info;
+  }
+
+  /** Drop the cached {@link OmBucketInfo} for the given bucket DB key. Invoked
+   *  when a bucketTable delete event is observed so the next key event 
re-reads
+   *  the current bucket info. This matters when a bucket is deleted and
+   *  recreated under the same volume/bucket name, which assigns a new 
objectID;
+   *  the recreate always follows the delete, so invalidating on delete 
suffices. */
+  protected void invalidateBucketCache(String bucketDBKey) {
+    bucketInfoCache.remove(bucketDBKey);
+  }
+
   public ReconNamespaceSummaryManager getReconNamespaceSummaryManager() {
     return reconNamespaceSummaryManager;
   }
diff --git 
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTaskWithLegacy.java
 
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTaskWithLegacy.java
index 186a89e294a..cca44fd71fc 100644
--- 
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTaskWithLegacy.java
+++ 
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTaskWithLegacy.java
@@ -18,6 +18,7 @@
 package org.apache.hadoop.ozone.recon.tasks;
 
 import static org.apache.hadoop.ozone.OzoneConsts.OM_KEY_PREFIX;
+import static org.apache.hadoop.ozone.om.codec.OMDBDefinition.BUCKET_TABLE;
 import static org.apache.hadoop.ozone.om.codec.OMDBDefinition.KEY_TABLE;
 
 import java.io.IOException;
@@ -91,9 +92,20 @@ public Pair<Integer, Boolean> 
processWithLegacy(OMUpdateEventBatch events,
       OMDBUpdateEvent.OMDBUpdateAction action = omdbUpdateEvent.getAction();
       eventCounter++;
 
-      // we only process updates on OM's KeyTable
       String table = omdbUpdateEvent.getTable();
+      // A bucket can be deleted and recreated under the same name with a new
+      // objectID. A recreate is always preceded by a delete, so dropping the
+      // cached OmBucketInfo on the delete event is enough for a later key 
event
+      // to re-read the recreated bucket. Bucket property updates don't change
+      // objectID or layout, so they need not invalidate the cache.
+      if (table.equals(BUCKET_TABLE)) {
+        if (action == OMDBUpdateEvent.OMDBUpdateAction.DELETE) {
+          invalidateBucketCache(omdbUpdateEvent.getKey());
+        }
+        continue;
+      }
 
+      // we only process updates on OM's KeyTable
       if (!table.equals(KEY_TABLE)) {
         continue;
       }
@@ -363,8 +375,7 @@ private long setParentBucketId(OmKeyInfo keyInfo)
       throws IOException {
     String bucketKey = getReconOMMetadataManager()
         .getBucketKey(keyInfo.getVolumeName(), keyInfo.getBucketName());
-    OmBucketInfo parentBucketInfo =
-        getReconOMMetadataManager().getBucketTable().getSkipCache(bucketKey);
+    OmBucketInfo parentBucketInfo = lookupBucketCached(bucketKey);
 
     if (parentBucketInfo != null) {
       return parentBucketInfo.getObjectID();
@@ -388,8 +399,7 @@ private boolean isBucketLayoutValid(ReconOMMetadataManager 
metadataManager,
     String volumeName = keyInfo.getVolumeName();
     String bucketName = keyInfo.getBucketName();
     String bucketDBKey = metadataManager.getBucketKey(volumeName, bucketName);
-    OmBucketInfo omBucketInfo =
-        metadataManager.getBucketTable().getSkipCache(bucketDBKey);
+    OmBucketInfo omBucketInfo = lookupBucketCached(bucketDBKey);
 
     if (omBucketInfo.getBucketLayout() != LEGACY_BUCKET_LAYOUT) {
       LOG.debug(
diff --git 
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTaskWithOBS.java
 
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTaskWithOBS.java
index a7843961672..b30b837133d 100644
--- 
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTaskWithOBS.java
+++ 
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/NSSummaryTaskWithOBS.java
@@ -17,6 +17,7 @@
 
 package org.apache.hadoop.ozone.recon.tasks;
 
+import static org.apache.hadoop.ozone.om.codec.OMDBDefinition.BUCKET_TABLE;
 import static org.apache.hadoop.ozone.om.codec.OMDBDefinition.KEY_TABLE;
 
 import java.io.IOException;
@@ -201,10 +202,21 @@ public Pair<Integer, Boolean> 
processWithOBS(OMUpdateEventBatch events,
       OMDBUpdateEvent.OMDBUpdateAction action = omdbUpdateEvent.getAction();
       eventCounter++;
 
-      // We only process updates on OM's KeyTable
       String table = omdbUpdateEvent.getTable();
-      boolean updateOnKeyTable = table.equals(KEY_TABLE);
-      if (!updateOnKeyTable) {
+      // A bucket can be deleted and recreated under the same name with a new
+      // objectID. A recreate is always preceded by a delete, so dropping the
+      // cached OmBucketInfo on the delete event is enough for a later key 
event
+      // to re-read the recreated bucket. Bucket property updates don't change
+      // objectID or layout, so they need not invalidate the cache.
+      if (table.equals(BUCKET_TABLE)) {
+        if (action == OMDBUpdateEvent.OMDBUpdateAction.DELETE) {
+          invalidateBucketCache(omdbUpdateEvent.getKey());
+        }
+        continue;
+      }
+
+      // We only process updates on OM's KeyTable
+      if (!table.equals(KEY_TABLE)) {
         continue;
       }
 
@@ -234,15 +246,13 @@ public Pair<Integer, Boolean> 
processWithOBS(OMUpdateEventBatch events,
         String bucketName = updatedKeyInfo.getBucketName();
         String bucketDBKey =
             getReconOMMetadataManager().getBucketKey(volumeName, bucketName);
-        // Get bucket info from bucket table
-        OmBucketInfo omBucketInfo = 
getReconOMMetadataManager().getBucketTable()
-            .getSkipCache(bucketDBKey);
+        OmBucketInfo omBucketInfo = lookupBucketCached(bucketDBKey);
 
         if (omBucketInfo.getBucketLayout() != BUCKET_LAYOUT) {
           continue;
         }
 
-        long parentObjectID = getKeyParentID(updatedKeyInfo);
+        long parentObjectID = omBucketInfo.getObjectID();
 
         switch (action) {
         case PUT:
@@ -253,9 +263,10 @@ public Pair<Integer, Boolean> 
processWithOBS(OMUpdateEventBatch events,
           break;
         case UPDATE:
           if (oldKeyInfo != null) {
-            // delete first, then put
-            long oldKeyParentObjectID = getKeyParentID(oldKeyInfo);
-            handleDeleteKeyEvent(oldKeyInfo, nsSummaryMap, 
oldKeyParentObjectID);
+            // For OBS, parent is always the bucket, so same parentObjectID
+            // applies to old and new (a key cannot move between buckets via
+            // an UPDATE event — that would be a delete+put).
+            handleDeleteKeyEvent(oldKeyInfo, nsSummaryMap, parentObjectID);
           } else {
             LOG.warn("Update event does not have the old keyInfo for {}.",
                 updatedKey);
diff --git 
a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/tasks/TestNSSummaryTask.java
 
b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/tasks/TestNSSummaryTask.java
index 7de5f54b881..9d7d05405bb 100644
--- 
a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/tasks/TestNSSummaryTask.java
+++ 
b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/tasks/TestNSSummaryTask.java
@@ -21,16 +21,23 @@
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 
 import java.io.File;
 import java.io.IOException;
+import java.lang.reflect.Field;
 import java.util.Arrays;
 import java.util.Collections;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
 import java.util.Set;
+import org.apache.hadoop.ozone.om.helpers.OmBucketInfo;
 import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
 import org.apache.hadoop.ozone.recon.ReconConstants;
 import org.apache.hadoop.ozone.recon.api.types.NSSummary;
+import org.apache.hadoop.ozone.recon.recovery.ReconOMMetadataManager;
 import org.junit.jupiter.api.BeforeAll;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Nested;
@@ -39,10 +46,8 @@
 import org.junit.jupiter.api.io.TempDir;
 
 /**
- * Test for NSSummaryTask. Create one bucket of each layout
- * and test process and reprocess. Currently, there is no
- * support for OBS buckets. Check that the NSSummary
- * for the OBS bucket is null.
+ * Test for NSSummaryTask. Creates one bucket of each layout (FSO, Legacy, OBS)
+ * and exercises reprocess, parallel process(), and bucket-cache invalidation.
  */
 @TestInstance(TestInstance.Lifecycle.PER_CLASS)
 public class TestNSSummaryTask extends AbstractNSSummaryTaskTest {
@@ -66,6 +71,24 @@ void setUp(@TempDir File tmpDir) throws Exception {
     );
   }
 
+  @Test
+  public void testTaskInstancesReuseSharedProcessExecutor() throws Exception {
+    NSSummaryTask anotherTask = new NSSummaryTask(
+        getReconNamespaceSummaryManager(),
+        getReconOMMetadataManager(),
+        getOmConfiguration());
+
+    assertSame(getSubTaskExecutor(nSSummaryTask),
+        getSubTaskExecutor(anotherTask));
+  }
+
+  private Object getSubTaskExecutor(NSSummaryTask task) throws Exception {
+    Field executorField = NSSummaryTask.class.getDeclaredField(
+        "SUB_TASK_EXECUTOR");
+    executorField.setAccessible(true);
+    return executorField.get(task);
+  }
+
   /**
    * Nested class for testing NSSummaryTaskWithLegacy reprocess.
    */
@@ -136,11 +159,15 @@ public class TestProcess {
 
     private NSSummary nsSummaryForBucket1;
     private NSSummary nsSummaryForBucket2;
+    private NSSummary nsSummaryForBucket3;
+    private ReconOmTask.TaskResult processResult;
 
     @BeforeEach
     public void setUp() throws IOException {
       nSSummaryTask.reprocess(getReconOMMetadataManager());
-      nSSummaryTask.process(processEventBatch(), Collections.emptyMap());
+      // Exercise process() across all three bucket layouts in a single batch
+      // so the parallel sub-task dispatch is covered end-to-end.
+      processResult = nSSummaryTask.process(processEventBatch(), 
Collections.emptyMap());
 
       nsSummaryForBucket1 =
           getReconNamespaceSummaryManager().getNSSummary(BUCKET_ONE_OBJECT_ID);
@@ -148,12 +175,13 @@ public void setUp() throws IOException {
       nsSummaryForBucket2 =
           getReconNamespaceSummaryManager().getNSSummary(BUCKET_TWO_OBJECT_ID);
       assertNotNull(nsSummaryForBucket2);
-      NSSummary nsSummaryForBucket3 = 
getReconNamespaceSummaryManager().getNSSummary(BUCKET_THREE_OBJECT_ID);
+      nsSummaryForBucket3 =
+          
getReconNamespaceSummaryManager().getNSSummary(BUCKET_THREE_OBJECT_ID);
       assertNotNull(nsSummaryForBucket3);
     }
 
     private OMUpdateEventBatch processEventBatch() throws IOException {
-      // put file5 under bucket 2
+      // PUT file5 under bucket 2 (Legacy)
       String omPutKey =
           OM_KEY_PREFIX + VOL +
               OM_KEY_PREFIX + BUCKET_TWO +
@@ -169,7 +197,7 @@ private OMUpdateEventBatch processEventBatch() throws 
IOException {
                                       
.setAction(OMDBUpdateEvent.OMDBUpdateAction.PUT)
                                       .build();
 
-      // delete file 1 under bucket 1
+      // DELETE file1 under bucket 1 (FSO)
       String omDeleteKey = BUCKET_ONE_OBJECT_ID + OM_KEY_PREFIX + FILE_ONE;
       OmKeyInfo omDeleteInfo = buildOmKeyInfo(
           VOL, BUCKET_ONE, KEY_ONE, FILE_ONE,
@@ -183,7 +211,24 @@ private OMUpdateEventBatch processEventBatch() throws 
IOException {
                                       
.setAction(OMDBUpdateEvent.OMDBUpdateAction.DELETE)
                                       .build();
 
-      return new OMUpdateEventBatch(Arrays.asList(keyEvent1, keyEvent2), 0L);
+      // PUT file4 under bucket 3 (OBS) — exercises the OBS sub-task path so a
+      // regression in the OBS branch (e.g. missed events) is caught.
+      String omObsPutKey =
+          OM_KEY_PREFIX + VOL +
+              OM_KEY_PREFIX + BUCKET_THREE +
+              OM_KEY_PREFIX + KEY_FOUR;
+      OmKeyInfo omObsPutKeyInfo = buildOmKeyInfo(VOL, BUCKET_THREE, KEY_FOUR,
+          KEY_FOUR, KEY_FOUR_OBJECT_ID, BUCKET_THREE_OBJECT_ID, KEY_FOUR_SIZE);
+      OMDBUpdateEvent keyEvent3 = new OMDBUpdateEvent.
+                                          OMUpdateEventBuilder<String, 
OmKeyInfo>()
+                                      .setKey(omObsPutKey)
+                                      .setValue(omObsPutKeyInfo)
+                                      
.setTable(getOmMetadataManager().getKeyTable(getOBSBucketLayout())
+                                                    .getName())
+                                      
.setAction(OMDBUpdateEvent.OMDBUpdateAction.PUT)
+                                      .build();
+
+      return new OMUpdateEventBatch(Arrays.asList(keyEvent1, keyEvent2, 
keyEvent3), 0L);
     }
 
     @Test
@@ -216,5 +261,384 @@ public void testProcessBucket() throws IOException {
         assertEquals(0, fileSizeDist[i]);
       }
     }
+
+    @Test
+    public void testProcessObsBucket() {
+      // bucket 3 (OBS) had file3 from reprocess; the batch added file4.
+      assertEquals(2, nsSummaryForBucket3.getNumOfFiles());
+      assertEquals(KEY_THREE_SIZE + KEY_FOUR_SIZE,
+          nsSummaryForBucket3.getSizeOfFiles());
+    }
+
+    @Test
+    public void testProcessTaskResult() {
+      // Sub-task seek positions must be reported for all three layouts so the
+      // dispatcher can resume each sub-task independently on retry.
+      assertNotNull(processResult);
+      assertTrue(processResult.isTaskSuccess());
+      
assertNotNull(processResult.getSubTaskSeekPositions().get(NSSummaryTask.BucketType.FSO.name()));
+      
assertNotNull(processResult.getSubTaskSeekPositions().get(NSSummaryTask.BucketType.LEGACY.name()));
+      
assertNotNull(processResult.getSubTaskSeekPositions().get(NSSummaryTask.BucketType.OBS.name()));
+    }
+  }
+
+  /**
+   * Exercises parallel FSO / Legacy / OBS sub-task dispatch beyond the happy
+   * path: interleaved multi-mutation batches, per-layout filtering in a shared
+   * event stream, and independent seek positions per sub-task on retry.
+   */
+  @Nested
+  public class TestProcessParallelMixedBatch {
+
+    @BeforeEach
+    public void reprocessBaseline() throws IOException {
+      nSSummaryTask.reprocess(getReconOMMetadataManager());
+    }
+
+    @Test
+    public void testInterleavedBatchAppliesAllLayoutMutations() throws 
IOException {
+      // Deliberately scramble layout order: OBS, Legacy, FSO, OBS, Legacy.
+      ReconOmTask.TaskResult result = nSSummaryTask.process(
+          new OMUpdateEventBatch(Arrays.asList(
+              obsDeleteEvent(KEY_THREE, KEY_THREE_OBJECT_ID, KEY_THREE_SIZE),
+              legacyPutEvent(KEY_FIVE, KEY_FIVE_OBJECT_ID, KEY_FIVE_SIZE),
+              fsoDeleteEvent(KEY_ONE, KEY_ONE_OBJECT_ID, KEY_ONE_SIZE),
+              obsPutEvent(KEY_FOUR, KEY_FOUR_OBJECT_ID, KEY_FOUR_SIZE),
+              legacyUpdateEvent(KEY_TWO, KEY_TWO_OBJECT_ID, KEY_TWO_SIZE,
+                  KEY_TWO_UPDATE_SIZE)),
+              0L),
+          Collections.emptyMap());
+
+      assertTrue(result.isTaskSuccess());
+
+      NSSummary fsoBucket =
+          getReconNamespaceSummaryManager().getNSSummary(BUCKET_ONE_OBJECT_ID);
+      assertNotNull(fsoBucket);
+      assertEquals(0, fsoBucket.getNumOfFiles());
+      assertEquals(0, fsoBucket.getSizeOfFiles());
+
+      NSSummary legacyBucket =
+          getReconNamespaceSummaryManager().getNSSummary(BUCKET_TWO_OBJECT_ID);
+      assertNotNull(legacyBucket);
+      assertEquals(2, legacyBucket.getNumOfFiles());
+      assertEquals(KEY_TWO_UPDATE_SIZE + KEY_FIVE_SIZE,
+          legacyBucket.getSizeOfFiles());
+
+      NSSummary obsBucket =
+          
getReconNamespaceSummaryManager().getNSSummary(BUCKET_THREE_OBJECT_ID);
+      assertNotNull(obsBucket);
+      assertEquals(1, obsBucket.getNumOfFiles());
+      assertEquals(KEY_FOUR_SIZE, obsBucket.getSizeOfFiles());
+    }
+
+    @Test
+    public void testLayoutFilteringInMixedBatch() throws IOException {
+      // OBS-only keyTable events in a batch that still runs all three 
sub-tasks
+      // in parallel. Legacy and FSO must leave their buckets at the reprocess
+      // baseline while OBS applies delete+put.
+      ReconOmTask.TaskResult result = nSSummaryTask.process(
+          new OMUpdateEventBatch(Arrays.asList(
+              obsDeleteEvent(KEY_THREE, KEY_THREE_OBJECT_ID, KEY_THREE_SIZE),
+              obsPutEvent(KEY_FOUR, KEY_FOUR_OBJECT_ID, KEY_FOUR_SIZE)),
+              0L),
+          Collections.emptyMap());
+
+      assertTrue(result.isTaskSuccess());
+
+      assertEquals(1, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_ONE_OBJECT_ID).getNumOfFiles());
+      assertEquals(KEY_ONE_SIZE, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_ONE_OBJECT_ID).getSizeOfFiles());
+      assertEquals(1, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_TWO_OBJECT_ID).getNumOfFiles());
+      assertEquals(KEY_TWO_SIZE, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_TWO_OBJECT_ID).getSizeOfFiles());
+      assertEquals(1, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_THREE_OBJECT_ID).getNumOfFiles());
+      assertEquals(KEY_FOUR_SIZE, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_THREE_OBJECT_ID).getSizeOfFiles());
+    }
+
+    @Test
+    public void testIndependentSubTaskSeekPositions() throws IOException {
+      // Shared stream (5 events): OBS delete, Legacy put, FSO delete, OBS put,
+      // Legacy update. Legacy seek=3 skips the first three stream positions
+      // (including its own put at position 2) while FSO/OBS start at 0.
+      OMUpdateEventBatch batch = new OMUpdateEventBatch(Arrays.asList(
+          obsDeleteEvent(KEY_THREE, KEY_THREE_OBJECT_ID, KEY_THREE_SIZE),
+          legacyPutEvent(KEY_FIVE, KEY_FIVE_OBJECT_ID, KEY_FIVE_SIZE),
+          fsoDeleteEvent(KEY_ONE, KEY_ONE_OBJECT_ID, KEY_ONE_SIZE),
+          obsPutEvent(KEY_FOUR, KEY_FOUR_OBJECT_ID, KEY_FOUR_SIZE),
+          legacyUpdateEvent(KEY_TWO, KEY_TWO_OBJECT_ID, KEY_TWO_SIZE,
+              KEY_TWO_UPDATE_SIZE)),
+          0L);
+
+      Map<String, Integer> legacySeekOnly = new HashMap<>();
+      legacySeekOnly.put(NSSummaryTask.BucketType.LEGACY.name(), 3);
+
+      ReconOmTask.TaskResult result =
+          nSSummaryTask.process(batch, legacySeekOnly);
+
+      assertTrue(result.isTaskSuccess());
+      // Each sub-task reports its resume position independently. Legacy was
+      // told to skip the first three stream indices; FSO/OBS consumed from 0.
+      assertNotNull(result.getSubTaskSeekPositions());
+      assertEquals(0, result.getSubTaskSeekPositions()
+          .get(NSSummaryTask.BucketType.FSO.name()).intValue());
+      assertEquals(3, result.getSubTaskSeekPositions()
+          .get(NSSummaryTask.BucketType.LEGACY.name()).intValue());
+      assertEquals(0, result.getSubTaskSeekPositions()
+          .get(NSSummaryTask.BucketType.OBS.name()).intValue());
+
+      // FSO/OBS processed from the start; Legacy only picked up the update.
+      assertEquals(0, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_ONE_OBJECT_ID).getNumOfFiles());
+      assertEquals(1, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_TWO_OBJECT_ID).getNumOfFiles());
+      assertEquals(KEY_TWO_UPDATE_SIZE, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_TWO_OBJECT_ID).getSizeOfFiles());
+      assertEquals(1, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_THREE_OBJECT_ID).getNumOfFiles());
+      assertEquals(KEY_FOUR_SIZE, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_THREE_OBJECT_ID).getSizeOfFiles());
+    }
+
+    @Test
+    public void testLegacySeekZeroAppliesAllLegacyEvents() throws IOException {
+      // Contrast with testIndependentSubTaskSeekPositions: without a Legacy
+      // seek offset, the interleaved batch applies both Legacy mutations.
+      OMUpdateEventBatch batch = new OMUpdateEventBatch(Arrays.asList(
+          obsDeleteEvent(KEY_THREE, KEY_THREE_OBJECT_ID, KEY_THREE_SIZE),
+          legacyPutEvent(KEY_FIVE, KEY_FIVE_OBJECT_ID, KEY_FIVE_SIZE),
+          fsoDeleteEvent(KEY_ONE, KEY_ONE_OBJECT_ID, KEY_ONE_SIZE),
+          obsPutEvent(KEY_FOUR, KEY_FOUR_OBJECT_ID, KEY_FOUR_SIZE),
+          legacyUpdateEvent(KEY_TWO, KEY_TWO_OBJECT_ID, KEY_TWO_SIZE,
+              KEY_TWO_UPDATE_SIZE)),
+          0L);
+
+      assertTrue(nSSummaryTask.process(batch, 
Collections.emptyMap()).isTaskSuccess());
+      assertEquals(2, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_TWO_OBJECT_ID).getNumOfFiles());
+      assertEquals(KEY_TWO_UPDATE_SIZE + KEY_FIVE_SIZE,
+          getReconNamespaceSummaryManager()
+              .getNSSummary(BUCKET_TWO_OBJECT_ID).getSizeOfFiles());
+    }
+
+    @Test
+    public void testSequentialMixedBatches() throws IOException {
+      ReconOmTask.TaskResult batch1 = nSSummaryTask.process(
+          new OMUpdateEventBatch(Collections.singletonList(
+              fsoDeleteEvent(KEY_ONE, KEY_ONE_OBJECT_ID, KEY_ONE_SIZE)), 0L),
+          Collections.emptyMap());
+      assertTrue(batch1.isTaskSuccess());
+
+      ReconOmTask.TaskResult batch2 = nSSummaryTask.process(
+          new OMUpdateEventBatch(Arrays.asList(
+              legacyPutEvent(KEY_FIVE, KEY_FIVE_OBJECT_ID, KEY_FIVE_SIZE),
+              obsPutEvent(KEY_FOUR, KEY_FOUR_OBJECT_ID, KEY_FOUR_SIZE)),
+              0L),
+          Collections.emptyMap());
+      assertTrue(batch2.isTaskSuccess());
+
+      assertEquals(0, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_ONE_OBJECT_ID).getNumOfFiles());
+      assertEquals(2, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_TWO_OBJECT_ID).getNumOfFiles());
+      assertEquals(2, getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_THREE_OBJECT_ID).getNumOfFiles());
+      assertEquals(KEY_THREE_SIZE + KEY_FOUR_SIZE, 
getReconNamespaceSummaryManager()
+          .getNSSummary(BUCKET_THREE_OBJECT_ID).getSizeOfFiles());
+    }
+
+    private OMDBUpdateEvent fsoDeleteEvent(String fileName, long objectId,
+                                           long dataSize) {
+      String key = BUCKET_ONE_OBJECT_ID + OM_KEY_PREFIX + fileName;
+      OmKeyInfo keyInfo = buildOmKeyInfo(
+          VOL, BUCKET_ONE, fileName, fileName, objectId, BUCKET_ONE_OBJECT_ID,
+          dataSize);
+      return new OMDBUpdateEvent.OMUpdateEventBuilder<String, OmKeyInfo>()
+          .setKey(key)
+          .setValue(keyInfo)
+          
.setTable(getOmMetadataManager().getKeyTable(getFSOBucketLayout()).getName())
+          .setAction(OMDBUpdateEvent.OMDBUpdateAction.DELETE)
+          .build();
+    }
+
+    private OMDBUpdateEvent legacyPutEvent(String keyName, long objectId,
+                                           long dataSize) {
+      String key = OM_KEY_PREFIX + VOL + OM_KEY_PREFIX + BUCKET_TWO
+          + OM_KEY_PREFIX + keyName;
+      OmKeyInfo keyInfo = buildOmKeyInfo(
+          VOL, BUCKET_TWO, keyName, keyName, objectId, BUCKET_TWO_OBJECT_ID,
+          dataSize);
+      return new OMDBUpdateEvent.OMUpdateEventBuilder<String, OmKeyInfo>()
+          .setKey(key)
+          .setValue(keyInfo)
+          
.setTable(getOmMetadataManager().getKeyTable(getLegacyBucketLayout()).getName())
+          .setAction(OMDBUpdateEvent.OMDBUpdateAction.PUT)
+          .build();
+    }
+
+    private OMDBUpdateEvent legacyUpdateEvent(String keyName, long objectId,
+                                              long oldSize, long newSize) {
+      String key = OM_KEY_PREFIX + VOL + OM_KEY_PREFIX + BUCKET_TWO
+          + OM_KEY_PREFIX + keyName;
+      OmKeyInfo oldInfo = buildOmKeyInfo(
+          VOL, BUCKET_TWO, keyName, keyName, objectId, BUCKET_TWO_OBJECT_ID,
+          oldSize);
+      OmKeyInfo newInfo = buildOmKeyInfo(
+          VOL, BUCKET_TWO, keyName, keyName, objectId, BUCKET_TWO_OBJECT_ID,
+          newSize);
+      return new OMDBUpdateEvent.OMUpdateEventBuilder<String, OmKeyInfo>()
+          .setKey(key)
+          .setValue(newInfo)
+          .setOldValue(oldInfo)
+          
.setTable(getOmMetadataManager().getKeyTable(getLegacyBucketLayout()).getName())
+          .setAction(OMDBUpdateEvent.OMDBUpdateAction.UPDATE)
+          .build();
+    }
+
+    private OMDBUpdateEvent obsPutEvent(String keyName, long objectId,
+                                        long dataSize) {
+      String key = OM_KEY_PREFIX + VOL + OM_KEY_PREFIX + BUCKET_THREE
+          + OM_KEY_PREFIX + keyName;
+      OmKeyInfo keyInfo = buildOmKeyInfo(
+          VOL, BUCKET_THREE, keyName, keyName, objectId, 
BUCKET_THREE_OBJECT_ID,
+          dataSize);
+      return new OMDBUpdateEvent.OMUpdateEventBuilder<String, OmKeyInfo>()
+          .setKey(key)
+          .setValue(keyInfo)
+          
.setTable(getOmMetadataManager().getKeyTable(getOBSBucketLayout()).getName())
+          .setAction(OMDBUpdateEvent.OMDBUpdateAction.PUT)
+          .build();
+    }
+
+    private OMDBUpdateEvent obsDeleteEvent(String keyName, long objectId,
+                                           long dataSize) {
+      String key = OM_KEY_PREFIX + VOL + OM_KEY_PREFIX + BUCKET_THREE
+          + OM_KEY_PREFIX + keyName;
+      OmKeyInfo keyInfo = buildOmKeyInfo(
+          VOL, BUCKET_THREE, keyName, keyName, objectId, 
BUCKET_THREE_OBJECT_ID,
+          dataSize);
+      return new OMDBUpdateEvent.OMUpdateEventBuilder<String, OmKeyInfo>()
+          .setKey(key)
+          .setValue(keyInfo)
+          
.setTable(getOmMetadataManager().getKeyTable(getOBSBucketLayout()).getName())
+          .setAction(OMDBUpdateEvent.OMDBUpdateAction.DELETE)
+          .build();
+    }
+  }
+
+  /**
+   * Regression test for the field-level OmBucketInfo cache in
+   * {@link NSSummaryTaskDbEventHandler}. A bucket deleted and recreated under
+   * the same volume/bucket name gets a new objectID but reuses the same DB 
key.
+   * The cached entry must be invalidated when the bucketTable event is 
observed,
+   * otherwise key events for the recreated bucket would be attributed to the
+   * stale (old) bucket objectID.
+   */
+  @Nested
+  public class TestProcessBucketRecreate {
+
+    private static final String RECREATE_BUCKET = "bucketrecreate";
+    private static final long OLD_BUCKET_OBJECT_ID = 100L;
+    private static final long NEW_BUCKET_OBJECT_ID = 200L;
+    private static final long FIRST_KEY_OBJECT_ID = 101L;
+    private static final long SECOND_KEY_OBJECT_ID = 201L;
+    private static final long RECREATE_KEY_SIZE = 1024L;
+
+    @Test
+    public void testBucketCacheInvalidatedOnRecreate() throws IOException {
+      ReconOMMetadataManager reconMetadataMgr = getReconOMMetadataManager();
+      String bucketKey = reconMetadataMgr.getBucketKey(VOL, RECREATE_BUCKET);
+
+      // The bucket initially exists as an OBS bucket with the old objectID.
+      reconMetadataMgr.getBucketTable()
+          .put(bucketKey, buildObsBucketInfo(OLD_BUCKET_OBJECT_ID));
+
+      // A fresh task starts with an empty bucket cache and exercises the full
+      // parallel sub-task dispatch.
+      NSSummaryTask task = new NSSummaryTask(getReconNamespaceSummaryManager(),
+          reconMetadataMgr, getOmConfiguration());
+
+      // Batch 1: a key lands in the original bucket. This warms the OBS
+      // sub-task's bucket cache with the old objectID.
+      OMDBUpdateEvent firstKeyPut =
+          obsKeyPutEvent(bucketKey, KEY_ONE, FIRST_KEY_OBJECT_ID, 
OLD_BUCKET_OBJECT_ID);
+      task.process(new OMUpdateEventBatch(
+          Collections.singletonList(firstKeyPut), 0L), Collections.emptyMap());
+
+      NSSummary oldBucketSummary =
+          getReconNamespaceSummaryManager().getNSSummary(OLD_BUCKET_OBJECT_ID);
+      assertNotNull(oldBucketSummary);
+      assertEquals(1, oldBucketSummary.getNumOfFiles());
+      assertEquals(RECREATE_KEY_SIZE, oldBucketSummary.getSizeOfFiles());
+
+      // The bucket is deleted and recreated under the same name with a new
+      // objectID (same DB key, different identity). The manual bucketTable put
+      // below simulates OM DB state after the recreate: synthetic events in 
this
+      // harness do not apply bucketTable writes to RocksDB (unlike production,
+      // where Recon commits the OM batch before tasks run).
+      reconMetadataMgr.getBucketTable()
+          .put(bucketKey, buildObsBucketInfo(NEW_BUCKET_OBJECT_ID));
+
+      // Batch 2 mirrors the real recreate sequence: the old bucket is deleted,
+      // a new bucket is created under the same name, then a key lands in it.
+      // The bucketTable delete event must invalidate the stale cache entry so
+      // the key event re-reads the recreated bucket's objectID.
+      OMDBUpdateEvent bucketDelete = new OMDBUpdateEvent
+          .OMUpdateEventBuilder<String, OmBucketInfo>()
+          .setKey(bucketKey)
+          .setValue(buildObsBucketInfo(OLD_BUCKET_OBJECT_ID))
+          .setTable(reconMetadataMgr.getBucketTable().getName())
+          .setAction(OMDBUpdateEvent.OMDBUpdateAction.DELETE)
+          .build();
+      OMDBUpdateEvent bucketPut = new OMDBUpdateEvent
+          .OMUpdateEventBuilder<String, OmBucketInfo>()
+          .setKey(bucketKey)
+          .setValue(buildObsBucketInfo(NEW_BUCKET_OBJECT_ID))
+          .setTable(reconMetadataMgr.getBucketTable().getName())
+          .setAction(OMDBUpdateEvent.OMDBUpdateAction.PUT)
+          .build();
+      OMDBUpdateEvent secondKeyPut =
+          obsKeyPutEvent(bucketKey, KEY_TWO, SECOND_KEY_OBJECT_ID, 
NEW_BUCKET_OBJECT_ID);
+      task.process(new OMUpdateEventBatch(
+          Arrays.asList(bucketDelete, bucketPut, secondKeyPut), 0L), 
Collections.emptyMap());
+
+      // The new key must be attributed to the recreated bucket's objectID...
+      NSSummary newBucketSummary =
+          getReconNamespaceSummaryManager().getNSSummary(NEW_BUCKET_OBJECT_ID);
+      assertNotNull(newBucketSummary);
+      assertEquals(1, newBucketSummary.getNumOfFiles());
+      assertEquals(RECREATE_KEY_SIZE, newBucketSummary.getSizeOfFiles());
+
+      // ...and the stale (old) bucket must not have absorbed it.
+      NSSummary staleBucketSummary =
+          getReconNamespaceSummaryManager().getNSSummary(OLD_BUCKET_OBJECT_ID);
+      assertEquals(1, staleBucketSummary.getNumOfFiles());
+      assertEquals(RECREATE_KEY_SIZE, staleBucketSummary.getSizeOfFiles());
+    }
+
+    private OmBucketInfo buildObsBucketInfo(long objectId) {
+      return OmBucketInfo.newBuilder()
+          .setVolumeName(VOL)
+          .setBucketName(RECREATE_BUCKET)
+          .setObjectID(objectId)
+          .setBucketLayout(getOBSBucketLayout())
+          .build();
+    }
+
+    private OMDBUpdateEvent obsKeyPutEvent(String bucketKey, String keyName,
+                                           long keyObjectId, long 
parentObjectId) {
+      String omKey = bucketKey + OM_KEY_PREFIX + keyName;
+      OmKeyInfo keyInfo = buildOmKeyInfo(VOL, RECREATE_BUCKET, keyName, 
keyName,
+          keyObjectId, parentObjectId, RECREATE_KEY_SIZE);
+      return new OMDBUpdateEvent.OMUpdateEventBuilder<String, OmKeyInfo>()
+          .setKey(omKey)
+          .setValue(keyInfo)
+          
.setTable(getReconOMMetadataManager().getKeyTable(getOBSBucketLayout()).getName())
+          .setAction(OMDBUpdateEvent.OMDBUpdateAction.PUT)
+          .build();
+    }
   }
 }


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


Reply via email to