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 @@ -618,13 +618,12 @@ public void testAOSKeyDeletingWithSnapshotCreateParallelExecution()
assertTableRowCount(deletedDirTable, initialDeletedCount + 1);
assertTableRowCount(renameTable, initialRenameCount + 1);
Mockito.doAnswer(i -> {
List<OzoneManagerProtocolProtos.PurgePathRequest> purgePathRequestList = i.getArgument(4);
List<OzoneManagerProtocolProtos.PurgePathRequest> purgePathRequestList = i.getArgument(1);
for (OzoneManagerProtocolProtos.PurgePathRequest purgeRequest : purgePathRequestList) {
Assertions.assertNotEquals(deletePathKey, purgeRequest.getDeletedDir());
}
return null;
}).when(service).optimizeDirDeletesAndSubmitRequest(anyLong(), anyLong(),
anyLong(), anyList(), anyList(), eq(null), anyLong(), any(),
}).when(service).optimizeDirDeletesAndSubmitRequest(anyList(), anyList(), eq(null), anyLong(), any(),
any(ReclaimableDirFilter.class), any(ReclaimableKeyFilter.class), anyMap(), any(),
anyLong(), any(AtomicInteger.class));

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -322,7 +322,6 @@ ThreadPoolExecutor getDeletionThreadPool() {

@SuppressWarnings("checkstyle:ParameterNumber")
void optimizeDirDeletesAndSubmitRequest(
long dirNum, long subDirNum, long subFileNum,
List<Pair<String, OmKeyInfo>> allSubDirList,
List<PurgePathRequest> purgePathRequestList,
String snapTableKey, long startTime,
Expand All @@ -334,7 +333,7 @@ void optimizeDirDeletesAndSubmitRequest(

// Optimization to handle delete sub-dir and keys to remove quickly
// This case will be useful to handle when depth of directory is high
int subdirDelNum = 0;
int initialRequestCount = purgePathRequestList.size();
int subDirRecursiveCnt = 0;
while (subDirRecursiveCnt < allSubDirList.size() && remainNum.get() > 0) {
try {
Expand All @@ -347,25 +346,33 @@ void optimizeDirDeletesAndSubmitRequest(
if (!request.isPresent()) {
continue;
}
PurgePathRequest requestVal = request.get();
purgePathRequestList.add(requestVal);
// Count up the purgeDeletedDir, subDirs and subFiles
if (requestVal.hasDeletedDir() && !StringUtils.isBlank(requestVal.getDeletedDir())) {
subdirDelNum++;
}
subDirNum += requestVal.getMarkDeletedSubDirsCount();
subFileNum += requestVal.getDeletedSubFilesCount();
purgePathRequestList.add(request.get());
} catch (IOException e) {
LOG.error("Error while running delete directories and files " +
"background task. Will retry at next run for subset.", e);
break;
}
}
if (!purgePathRequestList.isEmpty()) {
submitPurgePathsWithBatching(purgePathRequestList, snapTableKey, expectedPreviousSnapshotId, bucketNameInfoMap);
List<PurgePathRequest> submitted = submitPurgePathsWithBatching(purgePathRequestList, snapTableKey,
expectedPreviousSnapshotId, bucketNameInfoMap);
long dirNum = 0;
long subdirDelNum = 0;
long subDirNum = 0;
long subFileNum = 0;
for (int i = 0; i < submitted.size(); i++) {
PurgePathRequest request = submitted.get(i);
if (request.hasDeletedDir() && !StringUtils.isBlank(request.getDeletedDir())) {
if (i < initialRequestCount) {
dirNum++;
} else {
subdirDelNum++;
}
}
subDirNum += request.getMarkDeletedSubDirsCount();
subFileNum += request.getDeletedSubFilesCount();
}

if (dirNum != 0 || subDirNum != 0 || subFileNum != 0) {
if (dirNum != 0 || subdirDelNum != 0 || subDirNum != 0 || subFileNum != 0) {
long subdirMoved = subDirNum - subdirDelNum;
deletedDirsCount.addAndGet(dirNum + subdirDelNum);
movedDirsCount.addAndGet(subdirMoved);
Expand Down Expand Up @@ -533,10 +540,10 @@ private OzoneManagerProtocolProtos.PurgePathRequest wrapPurgeRequest(
return purgePathsRequest.build();
}

private List<OzoneManagerProtocolProtos.OMResponse> submitPurgePathsWithBatching(List<PurgePathRequest> requests,
private List<PurgePathRequest> submitPurgePathsWithBatching(List<PurgePathRequest> requests,
String snapTableKey, UUID expectedPreviousSnapshotId, Map<VolumeBucketId, BucketNameInfo> bucketNameInfoMap) {

List<OzoneManagerProtocolProtos.OMResponse> responses = new ArrayList<>();
List<PurgePathRequest> submitted = new ArrayList<>();
List<PurgePathRequest> purgePathRequestBatch = new ArrayList<>();
long batchBytes = 0;

Expand All @@ -547,10 +554,10 @@ private List<OzoneManagerProtocolProtos.OMResponse> submitPurgePathsWithBatching
if (batchBytes + reqSize > ratisByteLimit && !purgePathRequestBatch.isEmpty()) {
OzoneManagerProtocolProtos.OMResponse resp =
submitPurgeRequest(snapTableKey, expectedPreviousSnapshotId, bucketNameInfoMap, purgePathRequestBatch);
if (!resp.getSuccess()) {
return Collections.emptyList();
if (resp == null || !resp.getSuccess()) {
return submitted;
}
responses.add(resp);
submitted.addAll(purgePathRequestBatch);
purgePathRequestBatch.clear();
batchBytes = 0;
}
Expand All @@ -564,13 +571,13 @@ private List<OzoneManagerProtocolProtos.OMResponse> submitPurgePathsWithBatching
if (!purgePathRequestBatch.isEmpty()) {
OzoneManagerProtocolProtos.OMResponse resp =
submitPurgeRequest(snapTableKey, expectedPreviousSnapshotId, bucketNameInfoMap, purgePathRequestBatch);
if (!resp.getSuccess()) {
return Collections.emptyList();
if (resp == null || !resp.getSuccess()) {
return submitted;
}
responses.add(resp);
submitted.addAll(purgePathRequestBatch);
}

return responses;
return submitted;
}

@VisibleForTesting
Expand Down Expand Up @@ -724,9 +731,6 @@ private boolean processDeletedDirectories(SnapshotInfo currentSnapshotInfo, KeyM
ReclaimableKeyFilter reclaimableFileFilter = new ReclaimableKeyFilter(getOzoneManager(),
omSnapshotManager, snapshotChainManager, currentSnapshotInfo, keyManager, lock)) {
long startTime = Time.monotonicNow();
long dirNum = 0L;
long subDirNum = 0L;
long subFileNum = 0L;
List<PurgePathRequest> purgePathRequestList = new ArrayList<>();
Map<VolumeBucketId, BucketNameInfo> bucketNameInfos = new HashMap<>();
AtomicInteger remainNum = new AtomicInteger(remaining);
Expand Down Expand Up @@ -755,18 +759,10 @@ private boolean processDeletedDirectories(SnapshotInfo currentSnapshotInfo, KeyM
if (!request.isPresent()) {
continue;
}
PurgePathRequest purgePathRequest = request.get();
purgePathRequestList.add(purgePathRequest);
// Count up the purgeDeletedDir, subDirs and subFiles
if (purgePathRequest.hasDeletedDir() && !StringUtils.isBlank(purgePathRequest.getDeletedDir())) {
dirNum++;
}
subDirNum += purgePathRequest.getMarkDeletedSubDirsCount();
subFileNum += purgePathRequest.getDeletedSubFilesCount();
purgePathRequestList.add(request.get());
}

optimizeDirDeletesAndSubmitRequest(dirNum, subDirNum,
subFileNum, allSubDirList, purgePathRequestList, snapshotTableKey,
optimizeDirDeletesAndSubmitRequest(allSubDirList, purgePathRequestList, snapshotTableKey,
startTime, getOzoneManager().getKeyManager(),
reclaimableDirFilter, reclaimableFileFilter, bucketNameInfos, expectedPreviousSnapshotId,
runCount, remainNum);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertTrue;

import com.google.protobuf.ServiceException;
import java.io.File;
import java.io.IOException;
import java.nio.file.Files;
Expand All @@ -40,14 +41,19 @@
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import org.apache.commons.lang3.StringUtils;
import org.apache.commons.lang3.tuple.Pair;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
import org.apache.hadoop.hdds.client.StandaloneReplicationConfig;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.conf.StorageUnit;
import org.apache.hadoop.hdds.server.ServerUtils;
import org.apache.hadoop.hdds.utils.db.DBConfigFromFile;
import org.apache.hadoop.ozone.ClientVersion;
import org.apache.hadoop.ozone.om.DeletingServiceMetrics;
import org.apache.hadoop.ozone.om.KeyManager;
import org.apache.hadoop.ozone.om.OMConfigKeys;
import org.apache.hadoop.ozone.om.OMMetadataManager;
import org.apache.hadoop.ozone.om.OmTestManagers;
import org.apache.hadoop.ozone.om.OzoneManager;
import org.apache.hadoop.ozone.om.helpers.BucketLayout;
Expand All @@ -58,6 +64,9 @@
import org.apache.hadoop.ozone.om.protocol.OzoneManagerProtocol;
import org.apache.hadoop.ozone.om.request.OMRequestTestUtils;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMResponse;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.PurgePathRequest;
import org.apache.hadoop.util.Time;
import org.apache.ozone.test.GenericTestUtils;
import org.apache.ratis.util.ExitUtils;
Expand All @@ -66,6 +75,9 @@
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.CsvSource;
import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.Mockito;

/**
Expand Down Expand Up @@ -259,6 +271,151 @@ void testUpdateAndRestart() throws Exception {
.isEqualTo(newInterval.toMillis());
}

@ParameterizedTest
@CsvSource({"1, true, false", "2, true, false", "3, true, false",
"1, false, false", "2, false, false", "3, false, false",
"1, true, true", "2, true, true", "3, true, true",
"1, false, true", "2, false, true", "3, false, true"})
void testPurgeDirectoriesSubmitFailure(int failedBatch, boolean throwException, boolean snapshot) throws Exception {
List<PurgePathRequest> purgeList = new ArrayList<>();
String directoryName = StringUtils.repeat("d", 1200);
OzoneManagerProtocolProtos.KeyInfo keyInfo =
OMRequestTestUtils.createOmKeyInfo("volume", "bucket", "key", RatisReplicationConfig.getInstance(ONE))
.build().getProtobuf(ClientVersion.CURRENT_VERSION);
for (int i = 0; i < 3; i++) {
purgeList.add(PurgePathRequest.newBuilder().setVolumeId(1).setBucketId(2)
.setDeletedDir(directoryName + i).addMarkDeletedSubDirs(keyInfo)
.addDeletedSubFiles(keyInfo).addDeletedSubFiles(keyInfo).build());
}

OzoneConfiguration conf = createConfAndInitValues(1);
// Each path fits in one batch, but two paths exceed the byte limit.
conf.setStorageSize(OMConfigKeys.OZONE_OM_RATIS_LOG_APPENDER_QUEUE_BYTE_LIMIT, 2304, StorageUnit.BYTES);
OmTestManagers managers = new OmTestManagers(conf);
try {
KeyManager keyManager = managers.getKeyManager();
DirectoryDeletingService service = keyManager.getDirDeletingService();
service.suspend();
DirectoryDeletingService subject = Mockito.spy(service);
DeletingServiceMetrics metrics = Mockito.spy(subject.getMetrics());
Mockito.doReturn(metrics).when(subject).getMetrics();
long initialDirs = metrics.getNumDirsSentForPurge();
long initialSubDirs = metrics.getNumSubDirsSentForPurge();
long initialSubFiles = metrics.getNumSubFilesSentForPurge();
String snapshotKey = snapshot ? "snapshot" : null;
subject.getTasks();

OMResponse success = OMResponse.newBuilder().setCmdType(OzoneManagerProtocolProtos.Type.PurgeDirectories)
.setStatus(OzoneManagerProtocolProtos.Status.OK).setSuccess(true).build();
List<OMRequest> submitted = new ArrayList<>();
Mockito.doAnswer(invocation -> {
submitted.add(invocation.getArgument(0));
if (submitted.size() == failedBatch) {
if (throwException) {
throw new ServiceException("Transient purge submit failure");
}
return success.toBuilder().setStatus(OzoneManagerProtocolProtos.Status.INTERNAL_ERROR)
.setSuccess(false).build();
}
return success;
}).when(subject).submitRequest(Mockito.any(OMRequest.class));

subject.optimizeDirDeletesAndSubmitRequest(Collections.emptyList(), purgeList, snapshotKey,
Time.monotonicNow(), keyManager, kv -> true, kv -> true, Collections.emptyMap(), null, 1,
new AtomicInteger(0));

assertThat(submitted).hasSize(failedBatch);
int successfulBatches = failedBatch - 1;
assertThat(subject.getDeletedDirsCount()).isEqualTo(successfulBatches);
assertThat(subject.getMovedDirsCount()).isEqualTo(successfulBatches);
assertThat(subject.getMovedFilesCount()).isEqualTo(2L * successfulBatches);
assertThat(metrics.getNumDirsSentForPurge() - initialDirs).isEqualTo(successfulBatches);
assertThat(metrics.getNumSubDirsSentForPurge() - initialSubDirs).isEqualTo(successfulBatches);
assertThat(metrics.getNumSubFilesSentForPurge() - initialSubFiles).isEqualTo(2L * successfulBatches);
subject.execTaskCompletion();
Mockito.verify(metrics).updateAosDdsLastRunMetrics(snapshot ? 0 : successfulBatches,
snapshot ? 0 : successfulBatches, snapshot ? 0 : 2L * successfulBatches);
Mockito.verify(metrics).updateSnapDdsLastRunMetrics(snapshot ? successfulBatches : 0,
snapshot ? successfulBatches : 0, snapshot ? 2L * successfulBatches : 0);

// Successfully purged paths are no longer pending on the next run.
List<PurgePathRequest> remaining = new ArrayList<>(purgeList.subList(successfulBatches, purgeList.size()));
Mockito.clearInvocations(metrics);
subject.getTasks();
subject.optimizeDirDeletesAndSubmitRequest(Collections.emptyList(), remaining, snapshotKey,
Time.monotonicNow(), keyManager, kv -> true, kv -> true, Collections.emptyMap(), null, 2,
new AtomicInteger(0));

assertThat(submitted).hasSize(failedBatch + remaining.size());
for (int i = 0; i < remaining.size(); i++) {
OMRequest request = submitted.get(failedBatch + i);
assertThat(request.getCmdType()).isEqualTo(OzoneManagerProtocolProtos.Type.PurgeDirectories);
assertThat(request.getPurgeDirectoriesRequest().getDeletedPathList()).containsExactly(remaining.get(i));
}
assertThat(subject.getDeletedDirsCount()).isEqualTo(3);
assertThat(subject.getMovedDirsCount()).isEqualTo(3);
assertThat(subject.getMovedFilesCount()).isEqualTo(6);
assertThat(metrics.getNumDirsSentForPurge() - initialDirs).isEqualTo(3);
assertThat(metrics.getNumSubDirsSentForPurge() - initialSubDirs).isEqualTo(3);
assertThat(metrics.getNumSubFilesSentForPurge() - initialSubFiles).isEqualTo(6);
subject.execTaskCompletion();
Mockito.verify(metrics).updateAosDdsLastRunMetrics(snapshot ? 0 : remaining.size(),
snapshot ? 0 : remaining.size(), snapshot ? 0 : 2L * remaining.size());
Mockito.verify(metrics).updateSnapDdsLastRunMetrics(snapshot ? remaining.size() : 0,
snapshot ? remaining.size() : 0, snapshot ? 2L * remaining.size() : 0);
} finally {
managers.stop();
}
}

@ParameterizedTest
@ValueSource(ints = {0, 1, 2})
void testPurgeDirectoriesRecursiveAccounting(int failedBatch) throws Exception {
OzoneConfiguration conf = createConfAndInitValues(1);
conf.setStorageSize(OMConfigKeys.OZONE_OM_RATIS_LOG_APPENDER_QUEUE_BYTE_LIMIT, 2304, StorageUnit.BYTES);
OmTestManagers managers = new OmTestManagers(conf);
try {
KeyManager keyManager = managers.getKeyManager();
DirectoryDeletingService service = keyManager.getDirDeletingService();
service.suspend();
DirectoryDeletingService subject = Mockito.spy(service);
OMMetadataManager metadataManager = managers.getMetadataManager();
OmKeyInfo child = OMRequestTestUtils.createOmKeyInfo("volume", "bucket", StringUtils.repeat("d", 1200),
RatisReplicationConfig.getInstance(ONE)).setObjectID(3).setParentObjectID(2).build();
String childDeleteKey = metadataManager.getOzoneDeletePathKey(child.getObjectID(),
metadataManager.getOzonePathKey(1, 2, child.getParentObjectID(), child.getFileName()));
List<Pair<String, OmKeyInfo>> subDirs = new ArrayList<>();
subDirs.add(Pair.of(childDeleteKey, child));
PurgePathRequest parent = PurgePathRequest.newBuilder().setVolumeId(1).setBucketId(2)
.setDeletedDir("parent").addMarkDeletedSubDirs(child.getProtobuf(ClientVersion.CURRENT_VERSION)).build();
List<PurgePathRequest> purgeList = new ArrayList<>();
purgeList.add(parent);
List<OMRequest> submitted = new ArrayList<>();
Mockito.doAnswer(invocation -> {
submitted.add(invocation.getArgument(0));
boolean success = submitted.size() != failedBatch;
return OMResponse.newBuilder().setCmdType(OzoneManagerProtocolProtos.Type.PurgeDirectories)
.setStatus(success ? OzoneManagerProtocolProtos.Status.OK
: OzoneManagerProtocolProtos.Status.INTERNAL_ERROR)
.setSuccess(success).build();
}).when(subject).submitRequest(Mockito.any(OMRequest.class));

subject.optimizeDirDeletesAndSubmitRequest(subDirs, purgeList, null, Time.monotonicNow(), keyManager,
kv -> true, kv -> true, Collections.emptyMap(), null, 1, new AtomicInteger(10));

assertThat(submitted).hasSize(failedBatch == 1 ? 1 : 2);
assertThat(subject.getDeletedDirsCount()).isEqualTo(failedBatch == 0 ? 2 : failedBatch - 1);
assertThat(subject.getMovedDirsCount()).isEqualTo(failedBatch == 2 ? 1 : 0);
assertThat(subject.getMovedFilesCount()).isZero();
if (failedBatch != 1) {
assertThat(submitted.get(1).getPurgeDirectoriesRequest().getDeletedPathList())
.singleElement().satisfies(request -> assertThat(request.getDeletedDir()).isEqualTo(childDeleteKey));
}
} finally {
managers.stop();
}
}

@Test
@DisplayName("DirectoryDeletingService batches PurgeDirectories by Ratis byte limit (via submitRequest spy)")
void testPurgeDirectoriesBatching() throws Exception {
Expand Down Expand Up @@ -310,12 +467,13 @@ void testPurgeDirectoriesBatching() throws Exception {
bucketNameInfoMap = new HashMap<>();
bucketNameInfoMap.put(vbId, bni);

dds.optimizeDirDeletesAndSubmitRequest(0L, 0L, 0L, new ArrayList<>(), purgeList, null, Time.monotonicNow(), km,
dds.optimizeDirDeletesAndSubmitRequest(new ArrayList<>(), purgeList, null, Time.monotonicNow(), km,
kv -> true, kv -> true, bucketNameInfoMap, null, 1L, new AtomicInteger(Integer.MAX_VALUE));

assertThat(captured.size())
.as("Expect batching to respect Ratis byte limit")
.isBetween(3, 5);
assertThat(dds.getDeletedDirsCount()).isEqualTo(purgeList.size());

for (OzoneManagerProtocolProtos.OMRequest omReq : captured) {
assertThat(omReq.getCmdType()).isEqualTo(OzoneManagerProtocolProtos.Type.PurgeDirectories);
Expand Down
Loading