Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.SortedMap;
import java.util.UUID;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
Expand Down Expand Up @@ -66,9 +67,13 @@
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.OmMultipartKeyInfo;
import org.apache.hadoop.ozone.om.helpers.OmMultipartPartInfo;
import org.apache.hadoop.ozone.om.helpers.QuotaUtil;
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.om.request.util.OMMultipartUploadUtils;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
import org.apache.hadoop.util.Time;
import org.apache.ratis.protocol.ClientId;
Expand All @@ -86,9 +91,9 @@ public class QuotaRepairTask {
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.
* snapshot DB deleted keys/dirs, multipart upload parts.
*/
private static final int QUOTA_REPAIR_SCAN_TASKS = 6;
private static final int QUOTA_REPAIR_SCAN_TASKS = 7;
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);
Expand Down Expand Up @@ -325,6 +330,7 @@ private void repairCount(
Map<String, CountPair> directoryCountMap = new ConcurrentHashMap<>();
Map<String, CountPair> snapshotDeletedKeyMap = new ConcurrentHashMap<>();
Map<String, CountPair> snapshotDeletedDirMap = new ConcurrentHashMap<>();
Map<String, CountPair> mpuCountMap = new ConcurrentHashMap<>();
try {
nameBucketInfoMap.keySet().stream().forEach(e -> keyCountMap.put(e,
new CountPair()));
Expand All @@ -334,6 +340,7 @@ private void repairCount(
new CountPair()));
nameBucketInfoMap.keySet().forEach(k -> snapshotDeletedKeyMap.put(k, new CountPair()));
idBucketInfoMap.keySet().forEach(k -> snapshotDeletedDirMap.put(k, new CountPair()));
nameBucketInfoMap.keySet().forEach(k -> mpuCountMap.put(k, new CountPair()));

List<Future<?>> tasks = new ArrayList<>();
tasks.add(executor.submit(() -> recalculateUsages(
Expand Down Expand Up @@ -362,6 +369,7 @@ private void repairCount(
throw new UncheckedIOException(ex);
}
}));
tasks.add(executor.submit(() -> recalculateMultipartUsages(metadataManager, mpuCountMap)));

for (Future<?> f : tasks) {
f.get();
Expand All @@ -378,6 +386,7 @@ private void repairCount(
updateCountToBucketInfo(nameBucketInfoMap, keyCountMap);
updateCountToBucketInfo(idBucketInfoMap, fileCountMap);
updateCountToBucketInfo(idBucketInfoMap, directoryCountMap);
updateCountToBucketInfo(nameBucketInfoMap, mpuCountMap);
mergeSnapshotDeletedTableCounts(nameBucketInfoMap, snapshotDeletedKeyMap);
mergeDeletedDirSnapshotNamespace(idBucketInfoMap, snapshotDeletedDirMap);
LOG.info("Completed quota repair counting for all keys, files and directories");
Expand Down Expand Up @@ -549,6 +558,43 @@ private void recalculateDeletedDirNamespace(
}
}

private void recalculateMultipartUsages(
OMMetadataManager metadataManager, Map<String, CountPair> mpuCountMap) throws UncheckedIOException {
LOG.info("Starting recalculate multipart upload usages");

int count = 0;
long startTime = Time.monotonicNow();
try (Table.KeyValueIterator<String, OmMultipartKeyInfo> keyIter
= metadataManager.getMultipartInfoTable().iterator()) {
while (keyIter.hasNext()) {
Table.KeyValue<String, OmMultipartKeyInfo> kv = keyIter.next();
count++;
CountPair usage = mpuCountMap.get(getVolumeBucketPrefix(kv.getKey()));
if (usage == null) {
continue;
}
OmMultipartKeyInfo multipartKeyInfo = kv.getValue();
long replicatedSize = 0;
if (multipartKeyInfo.getSchemaVersion() == OmMultipartKeyInfo.LEGACY_SCHEMA_VERSION) {
for (OzoneManagerProtocolProtos.PartKeyInfo partKeyInfo : multipartKeyInfo.getPartKeyInfoMap()) {
replicatedSize += QuotaUtil.getReplicatedSize(
partKeyInfo.getPartKeyInfo().getDataSize(), multipartKeyInfo.getReplicationConfig());
}
} else {
SortedMap<Integer, OmMultipartPartInfo> parts =
OMMultipartUploadUtils.scanParts(metadataManager, multipartKeyInfo.getUploadID());
replicatedSize = OMMultipartUploadUtils.getReplicatedSize(
parts, multipartKeyInfo.getReplicationConfig());
}
usage.incrSpace(replicatedSize);
}
LOG.info("Recalculate multipart upload usages completed, count {} time {}ms",
count, (Time.monotonicNow() - startTime));
} catch (IOException ex) {
throw new UncheckedIOException(ex);
}
}

private static synchronized void mergeSnapshotDeletedTableCounts(
Map<String, OmBucketInfo> nameBucketInfoMap,
Map<String, CountPair> counts) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,25 +29,36 @@
import static org.mockito.Mockito.when;

import java.io.IOException;
import java.util.UUID;
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.protocol.proto.HddsProtos;
import org.apache.hadoop.hdds.utils.db.BatchOperation;
import org.apache.hadoop.hdds.utils.db.cache.CacheKey;
import org.apache.hadoop.hdds.utils.db.cache.CacheValue;
import org.apache.hadoop.ozone.OzoneConsts;
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.OmMultipartKeyInfo;
import org.apache.hadoop.ozone.om.helpers.OmMultipartPartInfo;
import org.apache.hadoop.ozone.om.helpers.OmMultipartPartKey;
import org.apache.hadoop.ozone.om.helpers.OmMultipartUpload;
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.OMKeyRequestTests;
import org.apache.hadoop.ozone.om.request.s3.multipart.S3MultipartUploadAbortRequest;
import org.apache.hadoop.ozone.om.request.s3.multipart.S3MultipartUploadAbortRequestWithFSO;
import org.apache.hadoop.ozone.om.request.volume.OMQuotaRepairRequest;
import org.apache.hadoop.ozone.om.response.OMClientResponse;
import org.apache.hadoop.ozone.om.response.volume.OMQuotaRepairResponse;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.KeyInfo;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.PartKeyInfo;
import org.apache.hadoop.util.Time;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
Expand Down Expand Up @@ -336,4 +347,154 @@ private void zeroOutBucketUsedBytes(String volumeName, String bucketName,
CacheValue.get(trxnLogIndex, bucketInfo));
omMetadataManager.getBucketTable().put(dbKey, bucketInfo);
}

@Test
public void testQuotaRepairCountsMpuParts() throws Exception {
AtomicReference<OzoneManagerProtocolProtos.OMRequest> request = mockQuotaRepairRequest();
OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
omMetadataManager, BucketLayout.OBJECT_STORE);

String legacyKey = "legacyMpuKey";
String legacyUploadId = UUID.randomUUID().toString();
OmMultipartKeyInfo legacyInfo = OMRequestTestUtils.createOmMultipartKeyInfo(
legacyUploadId, Time.now(), HddsProtos.ReplicationType.RATIS,
HddsProtos.ReplicationFactor.ONE, 1001L);
legacyInfo.addPartKeyInfo(createPart(legacyKey, legacyUploadId, 1, 100L));
legacyInfo.addPartKeyInfo(createPart(legacyKey, legacyUploadId, 1, 300L));
legacyInfo.addPartKeyInfo(createPart(legacyKey, legacyUploadId, 2, 200L));
addMultipartInfo(legacyKey, legacyInfo, 1L);

String splitKey = "splitMpuKey";
String splitUploadId = UUID.randomUUID().toString();
OmMultipartKeyInfo splitInfo = OMRequestTestUtils.createOmMultipartKeyInfo(
splitUploadId, Time.now(), HddsProtos.ReplicationType.RATIS,
HddsProtos.ReplicationFactor.ONE, 1002L).toBuilder()
.setSchemaVersion(OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION).build();
addMultipartInfo(splitKey, splitInfo, 2L);
omMetadataManager.getMultipartPartsTable().put(OmMultipartPartKey.of(splitUploadId, 1),
createSplitPart(bucketName, splitKey, splitUploadId, 1, 111L));
omMetadataManager.getMultipartPartsTable().put(OmMultipartPartKey.of(splitUploadId, 1),
createSplitPart(bucketName, splitKey, splitUploadId, 1, 400L));
omMetadataManager.getMultipartPartsTable().put(OmMultipartPartKey.of(splitUploadId, 2),
createSplitPart(bucketName, splitKey, splitUploadId, 2, 500L));

String bucketKey = corruptBucketUsage(bucketName, 12345L, 99L, 3L);
applyQuotaRepair(request, 4L);

OmBucketInfo repaired = omMetadataManager.getBucketTable().get(bucketKey);
assertEquals(1400L, repaired.getUsedBytes());
assertEquals(0L, repaired.getUsedNamespace());

abortMpu(bucketName, legacyKey, legacyUploadId, BucketLayout.OBJECT_STORE, 1003L);
abortMpu(bucketName, splitKey, splitUploadId, BucketLayout.OBJECT_STORE, 1004L);
assertEquals(0L, omMetadataManager.getBucketTable().get(bucketKey).getUsedBytes());
}

@Test
public void testQuotaRepairCountsMpuPartsFso() throws Exception {
AtomicReference<OzoneManagerProtocolProtos.OMRequest> request = mockQuotaRepairRequest();
String fsoBucketName = "fsomup" + bucketName;
OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, fsoBucketName,
omMetadataManager, BucketLayout.FILE_SYSTEM_OPTIMIZED);

String keyName = "fsoMpuFile";
String uploadId = UUID.randomUUID().toString();
OmMultipartKeyInfo mpuInfo = OMRequestTestUtils.createOmMultipartKeyInfo(
uploadId, Time.now(), HddsProtos.ReplicationType.RATIS,
HddsProtos.ReplicationFactor.THREE, 3001L).toBuilder()
.setSchemaVersion(OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION).build();
OmKeyInfo omKeyInfo = OMRequestTestUtils.createOmKeyInfo(volumeName, fsoBucketName,
keyName, RatisReplicationConfig.getInstance(THREE)).build();
OMRequestTestUtils.addMultipartInfoToTable(false, omKeyInfo, mpuInfo, 1L, omMetadataManager);
omMetadataManager.getMultipartPartsTable().put(OmMultipartPartKey.of(uploadId, 1),
createSplitPart(fsoBucketName, keyName, uploadId, 1, 100L));

String bucketKey = corruptBucketUsage(fsoBucketName, 777L, 88L, 2L);
applyQuotaRepair(request, 3L);

OmBucketInfo repaired = omMetadataManager.getBucketTable().get(bucketKey);
assertEquals(300L, repaired.getUsedBytes());
assertEquals(0L, repaired.getUsedNamespace());
abortMpu(fsoBucketName, keyName, uploadId, BucketLayout.FILE_SYSTEM_OPTIMIZED, 3002L);
assertEquals(0L, omMetadataManager.getBucketTable().get(bucketKey).getUsedBytes());
}

private AtomicReference<OzoneManagerProtocolProtos.OMRequest> mockQuotaRepairRequest() throws Exception {
OzoneManagerProtocolProtos.OMResponse response = mock(OzoneManagerProtocolProtos.OMResponse.class);
when(response.getSuccess()).thenReturn(true);
OzoneManagerRatisServer ratisServer = mock(OzoneManagerRatisServer.class);
AtomicReference<OzoneManagerProtocolProtos.OMRequest> request = new AtomicReference<>();
doAnswer(invocation -> {
request.set(invocation.getArgument(0, OzoneManagerProtocolProtos.OMRequest.class));
return response;
}).when(ratisServer).submitRequest(any(), any(), anyLong());
when(ozoneManager.getOmRatisServer()).thenReturn(ratisServer);
return request;
}

private void addMultipartInfo(String keyName, OmMultipartKeyInfo multipartInfo, long transactionIndex)
throws IOException {
OmKeyInfo omKeyInfo = OMRequestTestUtils.createOmKeyInfo(volumeName, bucketName,
keyName, RatisReplicationConfig.getInstance(ONE)).build();
OMRequestTestUtils.addMultipartInfoToTable(false, omKeyInfo, multipartInfo, transactionIndex,
omMetadataManager);
}

private String corruptBucketUsage(String bucket, long usedBytes, long usedNamespace, long transactionIndex)
throws IOException {
String bucketKey = omMetadataManager.getBucketKey(volumeName, bucket);
OmBucketInfo corrupted = omMetadataManager.getBucketTable().get(bucketKey).toBuilder()
.setUsedBytes(usedBytes).setUsedNamespace(usedNamespace).build();
omMetadataManager.getBucketTable().put(bucketKey, corrupted);
omMetadataManager.getBucketTable().addCacheEntry(
new CacheKey<>(bucketKey), CacheValue.get(transactionIndex, corrupted));
return bucketKey;
}

private void applyQuotaRepair(AtomicReference<OzoneManagerProtocolProtos.OMRequest> request,
long transactionIndex) throws Exception {
QuotaRepairTask quotaRepairTask = new QuotaRepairTask(ozoneManager);
assertTrue(awaitRepair(quotaRepairTask.repair()));
OMClientResponse response = new OMQuotaRepairRequest(request.get())
.validateAndUpdateCache(ozoneManager, transactionIndex);
BatchOperation batchOperation = omMetadataManager.getStore().initBatchOperation();
((OMQuotaRepairResponse) response).addToDBBatch(omMetadataManager, batchOperation);
omMetadataManager.getStore().commitBatchOperation(batchOperation);
}

private void abortMpu(String bucket, String keyName, String uploadId, BucketLayout layout,
long transactionIndex) throws Exception {
OzoneManagerProtocolProtos.OMRequest abortRequest = OMRequestTestUtils.createAbortMPURequest(
volumeName, bucket, keyName, uploadId);
OzoneManagerProtocolProtos.OMRequest preExecuted;
OMClientResponse response;
if (layout == BucketLayout.FILE_SYSTEM_OPTIMIZED) {
preExecuted = new S3MultipartUploadAbortRequestWithFSO(abortRequest, layout).preExecute(ozoneManager);
response = new S3MultipartUploadAbortRequestWithFSO(preExecuted, layout)
.validateAndUpdateCache(ozoneManager, transactionIndex);
} else {
preExecuted = new S3MultipartUploadAbortRequest(abortRequest, layout).preExecute(ozoneManager);
response = new S3MultipartUploadAbortRequest(preExecuted, layout)
.validateAndUpdateCache(ozoneManager, transactionIndex);
}
assertTrue(response.getOMResponse().getSuccess());
}

private PartKeyInfo createPart(String keyName, String uploadId, int partNumber, long dataSize) {
return PartKeyInfo.newBuilder().setPartNumber(partNumber)
.setPartName(OmMultipartUpload.getDbKey(volumeName, bucketName, keyName, uploadId))
.setPartKeyInfo(KeyInfo.newBuilder().setVolumeName(volumeName).setBucketName(bucketName)
.setKeyName(keyName).setDataSize(dataSize).setCreationTime(Time.now()).setModificationTime(Time.now())
.setType(HddsProtos.ReplicationType.RATIS).setFactor(ONE).build())
.build();
}

private OmMultipartPartInfo createSplitPart(String bucket, String keyName, String uploadId,
int partNumber, long dataSize) {
OmKeyInfo partKeyInfo = OMRequestTestUtils.createOmKeyInfo(volumeName, bucket, keyName,
RatisReplicationConfig.getInstance(ONE)).setDataSize(dataSize)
.addMetadata(OzoneConsts.ETAG, "etag-" + partNumber).build();
return OmMultipartPartInfo.from(OmMultipartUpload.getDbKey(volumeName, bucket, keyName, uploadId),
partNumber, partKeyInfo);
}
}