diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/service/TestDirectoryDeletingServiceWithFSO.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/service/TestDirectoryDeletingServiceWithFSO.java index 58d816fbabd..5bd8ab1bc85 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/service/TestDirectoryDeletingServiceWithFSO.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/service/TestDirectoryDeletingServiceWithFSO.java @@ -676,13 +676,12 @@ public void testAOSKeyDeletingWithSnapshotCreateParallelExecution() assertTableRowCount(deletedDirTable, initialDeletedCount + 1); assertTableRowCount(renameTable, initialRenameCount + 1); Mockito.doAnswer(i -> { - List purgePathRequestList = i.getArgument(4); + List 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)); diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/DirectoryDeletingService.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/DirectoryDeletingService.java index ef3bd91a3e2..4c19f0055c9 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/DirectoryDeletingService.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/DirectoryDeletingService.java @@ -33,12 +33,14 @@ import java.util.Collection; import java.util.Collections; import java.util.HashMap; +import java.util.HashSet; import java.util.Iterator; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Objects; import java.util.Optional; +import java.util.Set; import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; @@ -323,7 +325,6 @@ ThreadPoolExecutor getDeletionThreadPool() { @SuppressWarnings("checkstyle:ParameterNumber") void optimizeDirDeletesAndSubmitRequest( - long dirNum, long subDirNum, long subFileNum, List> allSubDirList, List purgePathRequestList, String snapTableKey, long startTime, @@ -335,7 +336,8 @@ 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; + // Bucket grouping can submit recursive requests before some of the initial requests. + Set initialRequests = new HashSet<>(purgePathRequestList); int subDirRecursiveCnt = 0; while (subDirRecursiveCnt < allSubDirList.size() && remainNum.get() > 0) { try { @@ -348,25 +350,32 @@ 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 submitted = submitPurgePathsWithBatching(purgePathRequestList, snapTableKey, + expectedPreviousSnapshotId, bucketNameInfoMap); + long dirNum = 0; + long subdirDelNum = 0; + long subDirNum = 0; + long subFileNum = 0; + for (PurgePathRequest request : submitted) { + if (request.hasDeletedDir() && !StringUtils.isBlank(request.getDeletedDir())) { + if (initialRequests.contains(request)) { + 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); @@ -534,7 +543,7 @@ private OzoneManagerProtocolProtos.PurgePathRequest wrapPurgeRequest( return purgePathsRequest.build(); } - private List submitPurgePathsWithBatching(List requests, + private List submitPurgePathsWithBatching(List requests, String snapTableKey, UUID expectedPreviousSnapshotId, Map bucketNameInfoMap) { // Group purge paths by their owning bucket so that every submitted purge transaction contains paths from a single @@ -546,7 +555,7 @@ private List submitPurgePathsWithBatching k -> new ArrayList<>()).add(req); } - List responses = new ArrayList<>(); + List submitted = new ArrayList<>(); for (List bucketRequests : requestsByBucket.values()) { List purgePathRequestBatch = new ArrayList<>(); long batchBytes = 0; @@ -559,9 +568,9 @@ private List submitPurgePathsWithBatching OzoneManagerProtocolProtos.OMResponse resp = submitPurgeRequest(snapTableKey, expectedPreviousSnapshotId, bucketNameInfoMap, purgePathRequestBatch); if (resp == null || !resp.getSuccess()) { - return Collections.emptyList(); + return submitted; } - responses.add(resp); + submitted.addAll(purgePathRequestBatch); purgePathRequestBatch.clear(); batchBytes = 0; } @@ -576,13 +585,13 @@ private List submitPurgePathsWithBatching OzoneManagerProtocolProtos.OMResponse resp = submitPurgeRequest(snapTableKey, expectedPreviousSnapshotId, bucketNameInfoMap, purgePathRequestBatch); if (resp == null || !resp.getSuccess()) { - return Collections.emptyList(); + return submitted; } - responses.add(resp); + submitted.addAll(purgePathRequestBatch); } } - return responses; + return submitted; } @VisibleForTesting @@ -736,9 +745,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 purgePathRequestList = new ArrayList<>(); Map bucketNameInfos = new HashMap<>(); AtomicInteger remainNum = new AtomicInteger(remaining); @@ -767,18 +773,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); diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestDirectoryDeletingService.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestDirectoryDeletingService.java index 4f3e2a2eeec..c931e1c6694 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestDirectoryDeletingService.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestDirectoryDeletingService.java @@ -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.Path; @@ -39,14 +40,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; @@ -57,6 +63,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; @@ -65,6 +74,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; /** @@ -258,6 +270,157 @@ 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 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 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 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, 3}) + 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> 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 purgeList = new ArrayList<>(); + purgeList.add(parent); + PurgePathRequest otherBucket = PurgePathRequest.newBuilder().setVolumeId(1).setBucketId(3) + .setDeletedDir("other-parent").build(); + purgeList.add(otherBucket); + List 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 == 0 ? 3 : failedBatch); + assertThat(subject.getDeletedDirsCount()).isEqualTo(failedBatch == 0 ? 3 : 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)); + } + if (failedBatch == 0 || failedBatch == 3) { + assertThat(submitted.get(2).getPurgeDirectoriesRequest().getDeletedPathList()).containsExactly(otherBucket); + } + } finally { + managers.stop(); + } + } + @Test @DisplayName("DirectoryDeletingService batches PurgeDirectories by Ratis byte limit (via submitRequest spy)") void testPurgeDirectoriesBatching() throws Exception { @@ -308,12 +471,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); @@ -367,7 +531,7 @@ void testPurgeDirectoriesGroupedByBucketPerTransaction() throws Exception { .setVolumeId(volumeId).setBucketId(bucketIdB).setDeletedDir("b-dir-" + i).build()); } - 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).as("Expect one transaction per bucket when the byte limit is not hit").hasSize(2);