smengcl commented on code in PR #11009:
URL: https://github.com/apache/ozone/pull/11009#discussion_r4211952760
##########
hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/DirectoryDeletingService.java:
##########
@@ -670,44 +685,107 @@ void processDeletedDirsForStore(SnapshotInfo
currentSnapshotInfo, KeyManager key
UUID expectedPreviousSnapshotId = currentSnapshotInfo == null ?
snapshotChainManager.getLatestGlobalSnapshotId() :
SnapshotUtils.getPreviousSnapshotId(currentSnapshotInfo,
snapshotChainManager);
- Map<UUID, Pair<Long, Long>> exclusiveSizeMap = Maps.newConcurrentMap();
+ processedAllDeletedDirs = runDeletionWorkers(currentSnapshotInfo,
keyManager, dirSupplier, currentSnapshot,
+ expectedPreviousSnapshotId, exclusiveSizeMap, rnCnt, remainNum);
+ }
+
+ // If AOS or all directories have been processed for snapshot, update
snapshot size delta and deep clean flag
+ // if it is a snapshot. All snapshot DB iterators and handles have been
closed before this synchronous request.
+ if (processedAllDeletedDirs) {
+ List<OzoneManagerProtocolProtos.SetSnapshotPropertyRequest>
setSnapshotPropertyRequests = new ArrayList<>();
+
+ for (Map.Entry<UUID, Pair<Long, Long>> entry :
exclusiveSizeMap.entrySet()) {
+ UUID snapshotID = entry.getKey();
+ long exclusiveSize = entry.getValue().getLeft();
+ long exclusiveReplicatedSize = entry.getValue().getRight();
+
setSnapshotPropertyRequests.add(getSetSnapshotRequestUpdatingExclusiveSize(
+ exclusiveSize, exclusiveReplicatedSize, snapshotID));
+ }
+
+ // Updating directory deep clean flag of snapshot.
+ if (currentSnapshotInfo != null) {
+
setSnapshotPropertyRequests.add(OzoneManagerProtocolProtos.SetSnapshotPropertyRequest.newBuilder()
+ .setSnapshotKey(snapshotTableKey)
+ .setDeepCleanedDeletedDir(true)
+ .build());
+ }
+ submitSetSnapshotRequests(setSnapshotPropertyRequests);
+ }
+ }
- CompletableFuture<Boolean> processedAllDeletedDirs =
CompletableFuture.completedFuture(true);
- final int parallelThreads = numberOfParallelThreadsPerStore.get();
- for (int i = 0; i < parallelThreads; i++) {
+ /**
+ * Waits for snapshot workers to close DB handles before opening the
submission gate.
+ * The workers retain their GC locks through submission; the coordinator
closes its own iterator and handle.
+ * AOS has no task-owned snapshot handle, so its workers can close their
handles and submit independently.
+ */
+ @SuppressWarnings("checkstyle:ParameterNumber")
+ private boolean runDeletionWorkers(SnapshotInfo currentSnapshotInfo,
KeyManager keyManager,
+ DeletedDirSupplier dirSupplier,
UncheckedAutoCloseableSupplier<OmSnapshot> currentSnapshot,
+ UUID expectedPreviousSnapshotId, Map<UUID, Pair<Long, Long>>
exclusiveSizeMap, long rnCnt, int remainNum)
+ throws ExecutionException, InterruptedException {
+ int parallelThreads = numberOfParallelThreadsPerStore.get();
+ CountDownLatch snapshotDbHandlesClosed = currentSnapshotInfo == null ?
null : new CountDownLatch(parallelThreads);
+ CountDownLatch submitRequests = new CountDownLatch(1);
+ CompletableFuture<Boolean> workersCompleted =
CompletableFuture.completedFuture(true);
+ for (int i = 0; i < parallelThreads; i++) {
+ AtomicBoolean workerReady = new AtomicBoolean();
+ try {
CompletableFuture<Boolean> future = CompletableFuture.supplyAsync(()
-> {
try {
return processDeletedDirectories(currentSnapshotInfo,
keyManager, dirSupplier,
- expectedPreviousSnapshotId, exclusiveSizeMap, rnCnt,
remainNum);
+ expectedPreviousSnapshotId, exclusiveSizeMap, rnCnt,
remainNum, () -> {
+ if (snapshotDbHandlesClosed != null) {
+ signalWorkerReady(workerReady, snapshotDbHandlesClosed);
+ if (awaitUninterruptibly(submitRequests)) {
+ Thread.currentThread().interrupt();
+ }
+ }
+ });
} catch (Throwable e) {
return false;
+ } finally {
+ if (snapshotDbHandlesClosed != null) {
+ signalWorkerReady(workerReady, snapshotDbHandlesClosed);
+ }
}
}, isThreadPoolActive(deletionThreadPool) ? deletionThreadPool :
ForkJoinPool.commonPool());
Review Comment:
Thanks. Agreed. Removed the fallback so a stopped DDS executor reaches the
existing rejection handling. The regression now covers partial rejection and an
already-stopped executor with three configured workers. Both cases pass with
common-pool parallelism set to two.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]