This is an automated email from the ASF dual-hosted git repository.
jojochuang 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 d14585c8ce7 HDDS-15204. Quota repair includes snapshot pending-delete
usage. (#10217)
d14585c8ce7 is described below
commit d14585c8ce73544416448f27fc756e3cfc565e69
Author: Wei-Chiu Chuang <[email protected]>
AuthorDate: Mon Jun 22 12:11:55 2026 -0700
HDDS-15204. Quota repair includes snapshot pending-delete usage. (#10217)
Generated-by: Cursor <[email protected]>
Co-authored-by: sadanand48 <[email protected]>
---
.../hadoop/ozone/om/helpers/OmBucketInfo.java | 4 +-
.../src/main/proto/OmClientProtocol.proto | 2 +
.../om/request/volume/OMQuotaRepairRequest.java | 6 +
.../hadoop/ozone/om/service/QuotaRepairTask.java | 267 ++++++++++++++++++++-
.../ozone/om/service/TestQuotaRepairTask.java | 142 ++++++++++-
5 files changed, 408 insertions(+), 13 deletions(-)
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmBucketInfo.java
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmBucketInfo.java
index 0d4eabf0bc2..8da2be2755a 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmBucketInfo.java
+++
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmBucketInfo.java
@@ -267,7 +267,7 @@ public void decrUsedBytes(long bytes, boolean
increasePendingDeleteBytes) {
}
}
- private void incrSnapshotUsedBytes(long bytes) {
+ public void incrSnapshotUsedBytes(long bytes) {
this.snapshotUsedBytes += bytes;
}
@@ -282,7 +282,7 @@ public void decrUsedNamespace(long namespaceToUse, boolean
increasePendingDelete
}
}
- private void incrSnapshotUsedNamespace(long namespaceToUse) {
+ public void incrSnapshotUsedNamespace(long namespaceToUse) {
this.snapshotUsedNamespace += namespaceToUse;
}
diff --git
a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
index d71ff27f302..f50d812e51f 100644
--- a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
+++ b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
@@ -2381,6 +2381,8 @@ message BucketQuotaCount {
required int64 diffUsedBytes = 3;
required int64 diffUsedNamespace = 4;
required bool supportOldQuota = 5 [default=false];
+ optional int64 diffSnapshotUsedBytes = 6;
+ optional int64 diffSnapshotUsedNamespace = 7;
}
message QuotaRepairResponse {
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/volume/OMQuotaRepairRequest.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/volume/OMQuotaRepairRequest.java
index 08b38cb2174..c58b700529f 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/volume/OMQuotaRepairRequest.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/volume/OMQuotaRepairRequest.java
@@ -124,6 +124,12 @@ private void updateBucketInfo(
}
bucketInfo.incrUsedBytes(bucketCountInfo.getDiffUsedBytes());
bucketInfo.incrUsedNamespace(bucketCountInfo.getDiffUsedNamespace());
+ if (bucketCountInfo.hasDiffSnapshotUsedBytes()) {
+
bucketInfo.incrSnapshotUsedBytes(bucketCountInfo.getDiffSnapshotUsedBytes());
+ }
+ if (bucketCountInfo.hasDiffSnapshotUsedNamespace()) {
+
bucketInfo.incrSnapshotUsedNamespace(bucketCountInfo.getDiffSnapshotUsedNamespace());
+ }
if (bucketCountInfo.getSupportOldQuota()) {
OmBucketInfo.Builder builder = bucketInfo.toBuilder();
if (bucketInfo.getQuotaInBytes() == OLD_QUOTA_DEFAULT) {
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/QuotaRepairTask.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/QuotaRepairTask.java
index cba6ad4aec2..2805f84c837 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/QuotaRepairTask.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/QuotaRepairTask.java
@@ -21,6 +21,7 @@
import static org.apache.hadoop.hdds.utils.db.IteratorType.KEY_ONLY;
import static org.apache.hadoop.ozone.OzoneConsts.OLD_QUOTA_DEFAULT;
import static org.apache.hadoop.ozone.OzoneConsts.OM_KEY_PREFIX;
+import static
org.apache.hadoop.ozone.om.helpers.SnapshotInfo.SnapshotStatus.SNAPSHOT_ACTIVE;
import com.google.common.util.concurrent.UncheckedExecutionException;
import com.google.protobuf.ServiceException;
@@ -31,8 +32,12 @@
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
+import java.util.HashSet;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
+import java.util.Set;
+import java.util.UUID;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.CompletableFuture;
@@ -52,15 +57,22 @@
import org.apache.hadoop.hdds.utils.db.TableIterator;
import org.apache.hadoop.ozone.om.OMMetadataManager;
import org.apache.hadoop.ozone.om.OmMetadataManagerImpl;
+import org.apache.hadoop.ozone.om.OmSnapshot;
+import org.apache.hadoop.ozone.om.OmSnapshotManager;
import org.apache.hadoop.ozone.om.OzoneManager;
+import org.apache.hadoop.ozone.om.SnapshotChainInfo;
+import org.apache.hadoop.ozone.om.SnapshotChainManager;
import org.apache.hadoop.ozone.om.exceptions.OMException;
import org.apache.hadoop.ozone.om.helpers.BucketLayout;
import org.apache.hadoop.ozone.om.helpers.OmBucketInfo;
import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
+import org.apache.hadoop.ozone.om.helpers.RepeatedOmKeyInfo;
+import org.apache.hadoop.ozone.om.helpers.SnapshotInfo;
import org.apache.hadoop.ozone.om.ratis.utils.OzoneManagerRatisUtils;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
import org.apache.hadoop.util.Time;
import org.apache.ratis.protocol.ClientId;
+import org.apache.ratis.util.function.UncheckedAutoCloseableSupplier;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -72,6 +84,11 @@ public class QuotaRepairTask {
QuotaRepairTask.class);
private static final int BATCH_SIZE = 5000;
private static final int TASK_THREAD_CNT = 3;
+ /**
+ * Parallel full-table scans: OBS keys, FSO files, dirs, active deleted
keys/dirs,
+ * snapshot DB deleted keys/dirs.
+ */
+ private static final int QUOTA_REPAIR_SCAN_TASKS = 6;
private static final AtomicBoolean IN_PROGRESS = new AtomicBoolean(false);
private static final RepairStatus REPAIR_STATUS = new RepairStatus();
private static final AtomicLong RUN_CNT = new AtomicLong(0);
@@ -102,8 +119,8 @@ public static String getStatus() {
private boolean repairTask(List<String> buckets) {
LOG.info("Starting quota repair task {}", REPAIR_STATUS);
- // thread pool with 3 Table type * (1 task each + 3 thread for each task)
- executor = Executors.newFixedThreadPool(3 * (1 + TASK_THREAD_CNT));
+ // thread pool: scan task types * (1 coordinator + worker threads per task)
+ executor = Executors.newFixedThreadPool(QUOTA_REPAIR_SCAN_TASKS * (1 +
TASK_THREAD_CNT));
try (OMMetadataManager activeMetaManager =
createActiveDBCheckpoint(om.getMetadataManager(),
om.getConfiguration())) {
OzoneManagerProtocolProtos.QuotaRepairRequest.Builder builder
@@ -111,8 +128,6 @@ private boolean repairTask(List<String> buckets) {
// repair active db
repairActiveDb(activeMetaManager, builder, buckets);
- // TODO: repair snapshots for quota
-
// submit request to update
ClientId clientId = ClientId.randomId();
OzoneManagerProtocolProtos.OMRequest omRequest =
OzoneManagerProtocolProtos.OMRequest.newBuilder()
@@ -175,6 +190,10 @@ private void repairActiveDb(
bucketCountBuilder.setDiffUsedBytes(updatedBuckedInfo.getUsedBytes() -
oriBucketInfo.getUsedBytes());
bucketCountBuilder.setDiffUsedNamespace(
updatedBuckedInfo.getUsedNamespace() -
oriBucketInfo.getUsedNamespace());
+ bucketCountBuilder.setDiffSnapshotUsedBytes(
+ updatedBuckedInfo.getSnapshotUsedBytes() -
oriBucketInfo.getSnapshotUsedBytes());
+ bucketCountBuilder.setDiffSnapshotUsedNamespace(
+ updatedBuckedInfo.getSnapshotUsedNamespace() -
oriBucketInfo.getSnapshotUsedNamespace());
bucketCountBuilder.setSupportOldQuota(oldQuota);
builder.addBucketCount(bucketCountBuilder.build());
}
@@ -253,17 +272,28 @@ private static void populateBucket(
oriBucketInfoMap.put(bucketNameKey, bucketInfo.copyObject());
bucketInfo.decrUsedBytes(bucketInfo.getUsedBytes(), false);
bucketInfo.decrUsedNamespace(bucketInfo.getUsedNamespace(), false);
+ resetSnapshotBucketQuota(bucketInfo);
nameBucketInfoMap.put(bucketNameKey, bucketInfo);
idBucketInfoMap.put(buildIdPath(metadataManager.getVolumeId(bucketInfo.getVolumeName()),
bucketInfo.getObjectID()), bucketInfo);
}
private boolean isChange(OmBucketInfo lBucketInfo, OmBucketInfo rBucketInfo)
{
- if (lBucketInfo.getUsedNamespace() != rBucketInfo.getUsedNamespace()
- || lBucketInfo.getUsedBytes() != rBucketInfo.getUsedBytes()) {
- return true;
+ return lBucketInfo.getUsedNamespace() != rBucketInfo.getUsedNamespace()
+ || lBucketInfo.getUsedBytes() != rBucketInfo.getUsedBytes()
+ || lBucketInfo.getSnapshotUsedBytes() !=
rBucketInfo.getSnapshotUsedBytes()
+ || lBucketInfo.getSnapshotUsedNamespace() !=
rBucketInfo.getSnapshotUsedNamespace();
+ }
+
+ private static void resetSnapshotBucketQuota(OmBucketInfo bucketInfo) {
+ long snapBytes = bucketInfo.getSnapshotUsedBytes();
+ if (snapBytes != 0) {
+ bucketInfo.purgeSnapshotUsedBytes(snapBytes);
+ }
+ long snapNs = bucketInfo.getSnapshotUsedNamespace();
+ if (snapNs != 0) {
+ bucketInfo.purgeSnapshotUsedNamespace(snapNs);
}
- return false;
}
private static String buildNamePath(String volumeName, String bucketName) {
@@ -293,6 +323,8 @@ private void repairCount(
Map<String, CountPair> keyCountMap = new ConcurrentHashMap<>();
Map<String, CountPair> fileCountMap = new ConcurrentHashMap<>();
Map<String, CountPair> directoryCountMap = new ConcurrentHashMap<>();
+ Map<String, CountPair> snapshotDeletedKeyMap = new ConcurrentHashMap<>();
+ Map<String, CountPair> snapshotDeletedDirMap = new ConcurrentHashMap<>();
try {
nameBucketInfoMap.keySet().stream().forEach(e -> keyCountMap.put(e,
new CountPair()));
@@ -300,7 +332,9 @@ private void repairCount(
new CountPair()));
idBucketInfoMap.keySet().stream().forEach(e -> directoryCountMap.put(e,
new CountPair()));
-
+ nameBucketInfoMap.keySet().forEach(k -> snapshotDeletedKeyMap.put(k, new
CountPair()));
+ idBucketInfoMap.keySet().forEach(k -> snapshotDeletedDirMap.put(k, new
CountPair()));
+
List<Future<?>> tasks = new ArrayList<>();
tasks.add(executor.submit(() -> recalculateUsages(
metadataManager.getKeyTable(BucketLayout.OBJECT_STORE),
@@ -311,6 +345,24 @@ private void repairCount(
tasks.add(executor.submit(() -> recalculateUsages(
metadataManager.getDirectoryTable(),
directoryCountMap, "Directory usages", false)));
+
+ Map<Long, OmBucketInfo> bucketById =
buildBucketByObjectId(idBucketInfoMap);
+
+ tasks.add(executor.submit(() -> recalculateDeletedKeyUsages(
+ metadataManager.getDeletedTable(), bucketById, snapshotDeletedKeyMap,
+ "active DB checkpoint")));
+ tasks.add(executor.submit(() -> recalculateDeletedDirNamespace(
+ metadataManager.getDeletedDirTable(), snapshotDeletedDirMap,
+ "active DB checkpoint")));
+ tasks.add(executor.submit(() -> {
+ try {
+ recalculateSnapshotDbPendingDeleteQuota(nameBucketInfoMap,
bucketById, metadataManager,
+ snapshotDeletedKeyMap, snapshotDeletedDirMap);
+ } catch (IOException ex) {
+ throw new UncheckedIOException(ex);
+ }
+ }));
+
for (Future<?> f : tasks) {
f.get();
}
@@ -326,9 +378,200 @@ private void repairCount(
updateCountToBucketInfo(nameBucketInfoMap, keyCountMap);
updateCountToBucketInfo(idBucketInfoMap, fileCountMap);
updateCountToBucketInfo(idBucketInfoMap, directoryCountMap);
+ mergeSnapshotDeletedTableCounts(nameBucketInfoMap, snapshotDeletedKeyMap);
+ mergeDeletedDirSnapshotNamespace(idBucketInfoMap, snapshotDeletedDirMap);
LOG.info("Completed quota repair counting for all keys, files and
directories");
}
+ private static Map<Long, OmBucketInfo> buildBucketByObjectId(
+ Map<String, OmBucketInfo> idBucketInfoMap) {
+ Map<Long, OmBucketInfo> bucketById = new HashMap<>();
+ for (OmBucketInfo bucketInfo : idBucketInfoMap.values()) {
+ bucketById.putIfAbsent(bucketInfo.getObjectID(), bucketInfo);
+ }
+ return bucketById;
+ }
+
+ /**
+ * Recompute pending-delete quota from each ACTIVE snapshot DB on repaired
buckets' path chains.
+ * Runs as an executor task in parallel with active-table scans.
+ */
+ private void recalculateSnapshotDbPendingDeleteQuota(
+ Map<String, OmBucketInfo> nameBucketInfoMap,
+ Map<Long, OmBucketInfo> bucketById,
+ OMMetadataManager activeMetaManager,
+ Map<String, CountPair> snapshotDeletedKeyMap,
+ Map<String, CountPair> snapshotDeletedDirMap) throws IOException {
+ LOG.info("Starting recalculate snapshot pending-delete from snapshot DBs");
+ OMMetadataManager liveMetaManager = om.getMetadataManager();
+ SnapshotChainManager chain =
+ ((OmMetadataManagerImpl) liveMetaManager).getSnapshotChainManager();
+ OmSnapshotManager snapshotManager = om.getOmSnapshotManager();
+ Set<UUID> scannedSnapshotIds = new HashSet<>();
+
+ for (OmBucketInfo bucket : nameBucketInfoMap.values()) {
+ String snapshotPath = buildSnapshotPath(bucket.getVolumeName(),
bucket.getBucketName());
+ LinkedHashMap<UUID, SnapshotChainInfo> pathChain;
+ try {
+ pathChain = chain.getSnapshotChainPath(snapshotPath);
+ } catch (IOException ex) {
+ throw new IOException("Failed to read snapshot chain for path " +
snapshotPath, ex);
+ }
+ if (pathChain == null || pathChain.isEmpty()) {
+ continue;
+ }
+ for (UUID snapshotId : pathChain.keySet()) {
+ if (!scannedSnapshotIds.add(snapshotId)) {
+ continue;
+ }
+ SnapshotInfo snapshotInfo = loadActiveSnapshot(activeMetaManager,
chain, snapshotId);
+ if (snapshotInfo == null) {
+ continue;
+ }
+ if (!snapshotInfo.getVolumeName().equals(bucket.getVolumeName())
+ || !snapshotInfo.getBucketName().equals(bucket.getBucketName())) {
+ continue;
+ }
+ if (!OmSnapshotManager.isSnapshotFlushedToDB(liveMetaManager,
snapshotInfo)) {
+ LOG.warn("Skipping snapshot {} for quota repair: create txn not
flushed to active DB",
+ snapshotInfo.getTableKey());
+ continue;
+ }
+ String sourceLabel = "snapshot DB " + snapshotInfo.getTableKey();
+ try (UncheckedAutoCloseableSupplier<OmSnapshot> snapshotRef =
+ snapshotManager.getSnapshot(snapshotId)) {
+ scanDeletedTables(snapshotRef.get().getMetadataManager(), bucketById,
+ snapshotDeletedKeyMap, snapshotDeletedDirMap, sourceLabel);
+ }
+ }
+ }
+ LOG.info("Recalculate snapshot pending-delete from snapshot DBs completed,
snapshots scanned: {}",
+ scannedSnapshotIds.size());
+ }
+
+ private static String buildSnapshotPath(String volumeName, String
bucketName) {
+ return volumeName + OM_KEY_PREFIX + bucketName;
+ }
+
+ private static SnapshotInfo loadActiveSnapshot(
+ OMMetadataManager metadataManager,
+ SnapshotChainManager chain,
+ UUID snapshotId) throws IOException {
+ String tableKey = chain.getTableKey(snapshotId);
+ if (tableKey == null) {
+ LOG.warn("Snapshot id {} is not present in snapshot chain table-key
map", snapshotId);
+ return null;
+ }
+ SnapshotInfo snapshotInfo =
metadataManager.getSnapshotInfoTable().get(tableKey);
+ if (snapshotInfo == null) {
+ LOG.warn("Snapshot {} not found in snapshotInfoTable during quota
repair", tableKey);
+ return null;
+ }
+ if (snapshotInfo.getSnapshotStatus() != SNAPSHOT_ACTIVE) {
+ return null;
+ }
+ return snapshotInfo;
+ }
+
+ private void scanDeletedTables(
+ OMMetadataManager metadataManager,
+ Map<Long, OmBucketInfo> bucketById,
+ Map<String, CountPair> snapshotDeletedKeyMap,
+ Map<String, CountPair> snapshotDeletedDirMap,
+ String sourceLabel) throws UncheckedIOException {
+ recalculateDeletedKeyUsages(metadataManager.getDeletedTable(), bucketById,
+ snapshotDeletedKeyMap, sourceLabel);
+ recalculateDeletedDirNamespace(metadataManager.getDeletedDirTable(),
+ snapshotDeletedDirMap, sourceLabel);
+ }
+
+ private void recalculateDeletedKeyUsages(
+ Table<String, RepeatedOmKeyInfo> deletedTable,
+ Map<Long, OmBucketInfo> bucketById,
+ Map<String, CountPair> snapshotCountByBucketNameKey,
+ String sourceLabel)
+ throws UncheckedIOException {
+ LOG.info("Starting recalculate snapshot usages from deletedTable ({})",
sourceLabel);
+
+ int count = 0;
+ long startTime = Time.monotonicNow();
+ try (Table.KeyValueIterator<String, RepeatedOmKeyInfo> keyIter
+ = deletedTable.iterator()) {
+ while (keyIter.hasNext()) {
+ Table.KeyValue<String, RepeatedOmKeyInfo> kv = keyIter.next();
+ count++;
+ RepeatedOmKeyInfo val = kv.getValue();
+ OmBucketInfo bucket = bucketById.get(val.getBucketId());
+ if (bucket == null) {
+ continue;
+ }
+ String nameKey = buildNamePath(bucket.getVolumeName(),
bucket.getBucketName());
+ CountPair usage = snapshotCountByBucketNameKey.get(nameKey);
+ if (usage == null) {
+ continue;
+ }
+ usage.incrSpace(val.getTotalSize().getRight());
+ usage.incrNamespace(val.getOmKeyInfoList().size());
+ }
+ LOG.info("Recalculate snapshot usages from deletedTable ({}) completed,
count {} time {}ms",
+ sourceLabel, count, (Time.monotonicNow() - startTime));
+ } catch (IOException ex) {
+ throw new UncheckedIOException(ex);
+ }
+ }
+
+ private void recalculateDeletedDirNamespace(
+ Table<String, OmKeyInfo> deletedDirTable,
+ Map<String, CountPair> snapshotDirNsByIdPrefix,
+ String sourceLabel)
+ throws UncheckedIOException {
+ LOG.info("Starting recalculate snapshot namespace from
deletedDirectoryTable ({})",
+ sourceLabel);
+
+ int count = 0;
+ long startTime = Time.monotonicNow();
+ try (Table.KeyValueIterator<String, OmKeyInfo> keyIter
+ = deletedDirTable.iterator()) {
+ while (keyIter.hasNext()) {
+ Table.KeyValue<String, OmKeyInfo> kv = keyIter.next();
+ count++;
+ String prefix = getVolumeBucketPrefix(kv.getKey());
+ CountPair usage = snapshotDirNsByIdPrefix.get(prefix);
+ if (usage != null) {
+ usage.incrNamespace(1L);
+ }
+ }
+ LOG.info(
+ "Recalculate snapshot namespace from deletedDirectoryTable ({})
completed, count {} time {}ms",
+ sourceLabel, count, (Time.monotonicNow() - startTime));
+ } catch (IOException ex) {
+ throw new UncheckedIOException(ex);
+ }
+ }
+
+ private static synchronized void mergeSnapshotDeletedTableCounts(
+ Map<String, OmBucketInfo> nameBucketInfoMap,
+ Map<String, CountPair> counts) {
+ for (Map.Entry<String, CountPair> entry : counts.entrySet()) {
+ OmBucketInfo bucket = nameBucketInfoMap.get(entry.getKey());
+ if (bucket != null) {
+ bucket.incrSnapshotUsedBytes(entry.getValue().getSpace());
+ bucket.incrSnapshotUsedNamespace(entry.getValue().getNamespace());
+ }
+ }
+ }
+
+ private static synchronized void mergeDeletedDirSnapshotNamespace(
+ Map<String, OmBucketInfo> idBucketInfoMap,
+ Map<String, CountPair> counts) {
+ for (Map.Entry<String, CountPair> entry : counts.entrySet()) {
+ OmBucketInfo bucket = idBucketInfoMap.get(entry.getKey());
+ if (bucket != null) {
+ bucket.incrSnapshotUsedNamespace(entry.getValue().getNamespace());
+ }
+ }
+ }
+
private <VALUE> void recalculateUsages(
Table<String, VALUE> table, Map<String, CountPair> prefixUsageMap,
String strType, boolean haveValue) throws UncheckedIOException,
@@ -500,6 +743,12 @@ public void
updateStatus(OzoneManagerProtocolProtos.QuotaRepairRequest.Builder b
ConcurrentHashMap<String, Long> diffCountMap = new
ConcurrentHashMap<>();
diffCountMap.put("DiffUsedBytes", quotaCount.getDiffUsedBytes());
diffCountMap.put("DiffUsedNamespace",
quotaCount.getDiffUsedNamespace());
+ if (quotaCount.hasDiffSnapshotUsedBytes()) {
+ diffCountMap.put("DiffSnapshotUsedBytes",
quotaCount.getDiffSnapshotUsedBytes());
+ }
+ if (quotaCount.hasDiffSnapshotUsedNamespace()) {
+ diffCountMap.put("DiffSnapshotUsedNamespace",
quotaCount.getDiffSnapshotUsedNamespace());
+ }
bucketCountDiffMap.put(bucketKey, diffCountMap);
}
}
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestQuotaRepairTask.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestQuotaRepairTask.java
index de950e8a5a1..cf8ddfbba56 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestQuotaRepairTask.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestQuotaRepairTask.java
@@ -20,6 +20,7 @@
import static
org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.ONE;
import static
org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyLong;
@@ -29,6 +30,7 @@
import java.io.IOException;
import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
import org.apache.hadoop.hdds.utils.db.BatchOperation;
@@ -38,6 +40,7 @@
import org.apache.hadoop.ozone.om.helpers.OmBucketInfo;
import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
import org.apache.hadoop.ozone.om.helpers.OmVolumeArgs;
+import org.apache.hadoop.ozone.om.helpers.RepeatedOmKeyInfo;
import org.apache.hadoop.ozone.om.ratis.OzoneManagerRatisServer;
import org.apache.hadoop.ozone.om.request.OMRequestTestUtils;
import org.apache.hadoop.ozone.om.request.key.TestOMKeyRequest;
@@ -47,12 +50,21 @@
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
import org.apache.hadoop.util.Time;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
/**
* Test class for quota repair.
*/
+@Timeout(120)
public class TestQuotaRepairTask extends TestOMKeyRequest {
+ /** Seconds; must match {@link Timeout} on this class. */
+ private static final int REPAIR_TEST_TIMEOUT_SECONDS = 120;
+
+ private static Boolean awaitRepair(CompletableFuture<Boolean> repair) throws
Exception {
+ return repair.get(REPAIR_TEST_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ }
+
@Test
public void testQuotaRepair() throws Exception {
OzoneManagerProtocolProtos.OMResponse respMock =
mock(OzoneManagerProtocolProtos.OMResponse.class);
@@ -110,7 +122,7 @@ public void testQuotaRepair() throws Exception {
QuotaRepairTask quotaRepairTask = new QuotaRepairTask(ozoneManager);
CompletableFuture<Boolean> repair = quotaRepairTask.repair();
- Boolean repairStatus = repair.get();
+ Boolean repairStatus = awaitRepair(repair);
assertTrue(repairStatus);
OMQuotaRepairRequest omQuotaRepairRequest = new
OMQuotaRepairRequest(ref.get());
@@ -170,7 +182,7 @@ public void testQuotaRepairForOldVersionVolumeBucket()
throws Exception {
QuotaRepairTask quotaRepairTask = new QuotaRepairTask(ozoneManager);
CompletableFuture<Boolean> repair = quotaRepairTask.repair();
- Boolean repairStatus = repair.get();
+ Boolean repairStatus = awaitRepair(repair);
assertTrue(repairStatus);
OMQuotaRepairRequest omQuotaRepairRequest = new
OMQuotaRepairRequest(ref.get());
@@ -187,6 +199,132 @@ public void testQuotaRepairForOldVersionVolumeBucket()
throws Exception {
assertEquals(-1, volArgsVerify.getQuotaInNamespace());
}
+ @Test
+ public void testQuotaRepairDeletedTableSnapshotQuota() throws Exception {
+ OzoneManagerProtocolProtos.OMResponse respMock =
mock(OzoneManagerProtocolProtos.OMResponse.class);
+ when(respMock.getSuccess()).thenReturn(true);
+ OzoneManagerRatisServer ratisServerMock =
mock(OzoneManagerRatisServer.class);
+ AtomicReference<OzoneManagerProtocolProtos.OMRequest> ref = new
AtomicReference<>();
+ doAnswer(invocation -> {
+ ref.set(invocation.getArgument(0,
OzoneManagerProtocolProtos.OMRequest.class));
+ return respMock;
+ }).when(ratisServerMock).submitRequest(any(), any(), anyLong());
+ when(ozoneManager.getOmRatisServer()).thenReturn(ratisServerMock);
+
+ OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+ omMetadataManager, BucketLayout.OBJECT_STORE);
+
+ String keyName = "/user/snapKey";
+ OMRequestTestUtils.addKeyToTableAndCache(volumeName, bucketName,
+ keyName, -1, RatisReplicationConfig.getInstance(THREE), 1L,
omMetadataManager);
+
+ String ozoneKey = omMetadataManager.getOzoneKey(volumeName, bucketName,
keyName);
+ OmBucketInfo bucketInfo = omMetadataManager.getBucketTable().get(
+ omMetadataManager.getBucketKey(volumeName, bucketName));
+ long bucketObjId = bucketInfo.getObjectID();
+
+ OMRequestTestUtils.deleteKey(ozoneKey, bucketObjId, omMetadataManager, 2L);
+
+ RepeatedOmKeyInfo deletedEntry =
omMetadataManager.getDeletedTable().get(ozoneKey);
+ long expectedSnapNs = deletedEntry.getOmKeyInfoList().size();
+
+ bucketInfo = omMetadataManager.getBucketTable().get(
+ omMetadataManager.getBucketKey(volumeName, bucketName));
+ OmBucketInfo corruptedSnapshot = bucketInfo.toBuilder()
+ .setSnapshotUsedBytes(7L)
+ .setSnapshotUsedNamespace(99L)
+ .build();
+ String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName);
+ omMetadataManager.getBucketTable().put(bucketKey, corruptedSnapshot);
+ omMetadataManager.getBucketTable().addCacheEntry(
+ new CacheKey<>(bucketKey), CacheValue.get(3L, corruptedSnapshot));
+
+ QuotaRepairTask quotaRepairTask = new QuotaRepairTask(ozoneManager);
+ CompletableFuture<Boolean> repair = quotaRepairTask.repair();
+ assertTrue(awaitRepair(repair));
+
+ OMQuotaRepairRequest omQuotaRepairRequest = new
OMQuotaRepairRequest(ref.get());
+ OMClientResponse omClientResponse =
omQuotaRepairRequest.validateAndUpdateCache(ozoneManager, 1);
+ BatchOperation batchOperation =
omMetadataManager.getStore().initBatchOperation();
+ ((OMQuotaRepairResponse) omClientResponse).addToDBBatch(omMetadataManager,
batchOperation);
+ omMetadataManager.getStore().commitBatchOperation(batchOperation);
+
+ OmBucketInfo repaired = omMetadataManager.getBucketTable().get(bucketKey);
+ assertEquals(0, repaired.getUsedBytes());
+ assertEquals(0, repaired.getUsedNamespace());
+ assertEquals(expectedSnapNs, repaired.getSnapshotUsedNamespace());
+ assertTrue(repaired.getSnapshotUsedBytes() > 0,
+ "Snapshot pending-delete bytes must be recomputed from deletedTable");
+ }
+
+ @Test
+ public void testQuotaRepairSnapshotDbDeletedTableQuota() throws Exception {
+ OzoneManagerProtocolProtos.OMResponse respMock =
mock(OzoneManagerProtocolProtos.OMResponse.class);
+ when(respMock.getSuccess()).thenReturn(true);
+ OzoneManagerRatisServer ratisServerMock =
mock(OzoneManagerRatisServer.class);
+ AtomicReference<OzoneManagerProtocolProtos.OMRequest> ref = new
AtomicReference<>();
+ doAnswer(invocation -> {
+ ref.set(invocation.getArgument(0,
OzoneManagerProtocolProtos.OMRequest.class));
+ return respMock;
+ }).when(ratisServerMock).submitRequest(any(), any(), anyLong());
+ when(ozoneManager.getOmRatisServer()).thenReturn(ratisServerMock);
+
+ OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+ omMetadataManager, BucketLayout.OBJECT_STORE);
+
+ String keyName = "/user/snapKey";
+ OMRequestTestUtils.addKeyToTableAndCache(volumeName, bucketName,
+ keyName, -1, RatisReplicationConfig.getInstance(THREE), 1L,
omMetadataManager);
+
+ String ozoneKey = omMetadataManager.getOzoneKey(volumeName, bucketName,
keyName);
+ OmKeyInfo omKeyInfo =
omMetadataManager.getKeyTable(BucketLayout.OBJECT_STORE).get(ozoneKey);
+ long keyBytes = omKeyInfo.getReplicatedSize();
+
+ OmBucketInfo bucketInfo = omMetadataManager.getBucketTable().get(
+ omMetadataManager.getBucketKey(volumeName, bucketName));
+ OMRequestTestUtils.deleteKey(ozoneKey, bucketInfo.getObjectID(),
omMetadataManager, 2L);
+
+ String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName);
+ OmBucketInfo afterDelete = bucketInfo.toBuilder()
+ .setUsedBytes(0)
+ .setUsedNamespace(0)
+ .setSnapshotUsedBytes(keyBytes)
+ .setSnapshotUsedNamespace(1)
+ .build();
+ omMetadataManager.getBucketTable().put(bucketKey, afterDelete);
+
+ when(ozoneManager.getDefaultReplicationConfig())
+ .thenReturn(RatisReplicationConfig.getInstance(THREE));
+ createSnapshot("snap1");
+
+ assertNull(omMetadataManager.getDeletedTable().get(ozoneKey),
+ "Deleted key should move out of active deletedTable after snapshot");
+ assertEquals(0,
omMetadataManager.countRowsInTable(omMetadataManager.getDeletedTable()));
+
+ OmBucketInfo corrupted = afterDelete.toBuilder()
+ .setSnapshotUsedBytes(7L)
+ .build();
+ omMetadataManager.getBucketTable().put(bucketKey, corrupted);
+ omMetadataManager.getBucketTable().addCacheEntry(
+ new CacheKey<>(bucketKey), CacheValue.get(3L, corrupted));
+
+ QuotaRepairTask quotaRepairTask = new QuotaRepairTask(ozoneManager);
+ CompletableFuture<Boolean> repair = quotaRepairTask.repair();
+ assertTrue(awaitRepair(repair));
+
+ OMQuotaRepairRequest omQuotaRepairRequest = new
OMQuotaRepairRequest(ref.get());
+ OMClientResponse omClientResponse =
omQuotaRepairRequest.validateAndUpdateCache(ozoneManager, 1);
+ BatchOperation batchOperation =
omMetadataManager.getStore().initBatchOperation();
+ ((OMQuotaRepairResponse) omClientResponse).addToDBBatch(omMetadataManager,
batchOperation);
+ omMetadataManager.getStore().commitBatchOperation(batchOperation);
+
+ OmBucketInfo repaired = omMetadataManager.getBucketTable().get(bucketKey);
+ assertEquals(0, repaired.getUsedBytes());
+ assertEquals(0, repaired.getUsedNamespace());
+ assertEquals(keyBytes, repaired.getSnapshotUsedBytes());
+ assertEquals(1, repaired.getSnapshotUsedNamespace());
+ }
+
private void zeroOutBucketUsedBytes(String volumeName, String bucketName,
long trxnLogIndex)
throws IOException {
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]