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]