rich7420 commented on code in PR #11009:
URL: https://github.com/apache/ozone/pull/11009#discussion_r4129051178
##########
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:
This fallback can deadlock the snapshot coordinator during shutdown. With
three workers and common-pool parallelism two, two workers wait at the
submission gate while the third remains queued. I reproduced this with a
stopped DDS executor and a real snapshot. Please submit directly to the DDS
executor and use the rejection handling below; the same test passes with that
change.
##########
hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestDirectoryDeletingService.java:
##########
@@ -383,4 +413,254 @@ void testPurgeDirectoriesGroupedByBucketPerTransaction()
throws Exception {
bucketIdA, bucketIdB);
}
+ /**
+ * Regression test for HDDS-16164. Snapshot workers must wait for the
task-owned handle to close before submitting.
+ * AOS workers own their handles and can finish scanning while a sibling
waits for double-buffer capacity.
+ */
+ @ParameterizedTest(name = "activeStore={0}")
+ @ValueSource(booleans = {false, true})
+ @DisplayName("DirectoryDeletingService must not deadlock at the unflushed
transaction limit")
+ void testNoUnflushedTransactionDeadlock(boolean activeStore) throws
Exception {
+ int threadCount = 3;
+ OzoneConfiguration conf = createConfAndInitValues(threadCount);
+ conf.setInt(OMConfigKeys.OZONE_OM_UNFLUSHED_TRANSACTION_MAX_COUNT, 1);
+ conf.setInt(OZONE_PATH_DELETING_LIMIT_PER_TASK, 1);
+ conf.setInt(OZONE_MANAGER_STRIPED_LOCK_SIZE_PREFIX + "snapshot_db_lock",
1);
+ conf.setTimeDuration(OZONE_BLOCK_DELETING_SERVICE_INTERVAL, 1,
TimeUnit.HOURS);
+ conf.setTimeDuration(OZONE_SNAPSHOT_DELETING_SERVICE_INTERVAL, 1,
TimeUnit.HOURS);
+ conf.setTimeDuration(OZONE_DIR_DELETING_SERVICE_INTERVAL, 1,
TimeUnit.HOURS);
+
+ OmTestManagers managers = new OmTestManagers(conf);
+ om = managers.getOzoneManager();
+ KeyManager keyManager = managers.getKeyManager();
+ OMMetadataManager metadataManager = managers.getMetadataManager();
+ OzoneManagerProtocol writeClient = managers.getWriteClient();
+ keyManager.getDirDeletingService().shutdown();
+
+ OMRequestTestUtils.addVolumeToOM(metadataManager,
+
OmVolumeArgs.newBuilder().setOwnerName("o").setAdminName("a").setVolume(volumeName).build());
+ OMRequestTestUtils.addBucketToOM(metadataManager,
+
OmBucketInfo.newBuilder().setVolumeName(volumeName).setBucketName(bucketName)
+
.setObjectID(1).setBucketLayout(BucketLayout.FILE_SYSTEM_OPTIMIZED).build());
+
+ String purgeSnapshotName = "purge";
+ writeClient.createSnapshot(volumeName, bucketName, purgeSnapshotName);
+ String previousSnapshotName = "previous";
+ writeClient.createSnapshot(volumeName, bucketName, previousSnapshotName);
+ String deepCleanSnapshotName = "deep-clean";
+ if (activeStore) {
+ writeClient.createSnapshot(volumeName, bucketName,
deepCleanSnapshotName);
+ }
+ om.awaitDoubleBufferFlush();
+
+ addDeletedDirectories(metadataManager, threadCount);
+ if (!activeStore) {
+ writeClient.createSnapshot(volumeName, bucketName,
deepCleanSnapshotName);
+ }
+ om.awaitDoubleBufferFlush();
+
+ String deepCleanKey = SnapshotInfo.getTableKey(volumeName, bucketName,
deepCleanSnapshotName);
+ SnapshotInfo deepCleanInfo =
metadataManager.getSnapshotInfoTable().get(deepCleanKey);
+ assertNotNull(deepCleanInfo);
+ awaitSnapshotFlushed(metadataManager, deepCleanInfo);
+
+ writeClient.deleteSnapshot(volumeName, bucketName, purgeSnapshotName);
+ om.awaitDoubleBufferFlush();
+ String purgeKey = SnapshotInfo.getTableKey(volumeName, bucketName,
purgeSnapshotName);
+ assertNotNull(metadataManager.getSnapshotInfoTable().get(purgeKey));
+
+ CountDownLatch ddsAtSubmit = new CountDownLatch(1);
+ CountDownLatch ddsMayProceed = new CountDownLatch(1);
+ DirectoryDeletingService service = Mockito.spy(new
DirectoryDeletingService(1, TimeUnit.HOURS, 1,
+ om, conf, threadCount, true) {
+ @Override
+ protected OzoneManagerProtocolProtos.OMResponse submitRequest(
+ OzoneManagerProtocolProtos.OMRequest omRequest) throws
ServiceException {
+ assertThat(omRequest.getCmdType()).isIn(
+ OzoneManagerProtocolProtos.Type.PurgeDirectories,
+ OzoneManagerProtocolProtos.Type.SetSnapshotProperty);
+ ddsAtSubmit.countDown();
+ try {
+ ddsMayProceed.await();
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new ServiceException(e);
+ }
+ return super.submitRequest(omRequest);
+ }
+ });
+
+ CountDownLatch workersPrepared = new CountDownLatch(threadCount);
+ CountDownLatch finishScanning = new CountDownLatch(1);
+ AtomicBoolean pauseScan = new AtomicBoolean();
+ Mockito.doAnswer(invocation -> {
+ workersPrepared.countDown();
+ if (pauseScan.compareAndSet(false, true)) {
+ // Keep this worker's previous-snapshot handle open while the other
workers finish scanning.
+ assertThat(metadataManager.getLock().getReadHoldCount(SNAPSHOT_DB_LOCK,
+ deepCleanInfo.getSnapshotId().toString())).isPositive();
+ assertTrue(finishScanning.await(30, TimeUnit.SECONDS));
+ }
+ return invocation.callRealMethod();
+ }).when(service).optimizeDirDeletesAndSubmitRequest(Mockito.anyLong(),
Mockito.anyLong(), Mockito.anyLong(),
+ Mockito.anyList(), Mockito.anyList(), Mockito.any(),
Mockito.anyLong(), Mockito.any(), Mockito.any(),
+ Mockito.any(), Mockito.anyMap(), Mockito.any(), Mockito.anyLong(),
Mockito.any(), Mockito.any());
+
+ OzoneManagerDoubleBuffer doubleBuffer =
+
om.getOmRatisServer().getOmStateMachine().getOzoneManagerDoubleBuffer();
+ Semaphore unflushedTransactions =
+ (Semaphore) HddsWhiteboxTestUtils.getInternalState(doubleBuffer,
"unFlushedTransactions");
+ ThreadPoolExecutor deletionThreadPool = service.getDeletionThreadPool();
+ long completedTaskCountBefore = deletionThreadPool.getCompletedTaskCount();
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+ try {
+ Future<?> ddsFuture = service.getExecutorService().submit(
+ () -> service.new DirDeletingTask(activeStore ? null :
deepCleanInfo.getSnapshotId()).call());
+ assertTrue(workersPrepared.await(30, TimeUnit.SECONDS), "Expected all
configured workers to scan in parallel");
+ if (!activeStore) {
+ assertThat(ddsAtSubmit.await(200, TimeUnit.MILLISECONDS))
+ .as("snapshot workers must wait for the task-owned snapshot handle
to close").isFalse();
+ finishScanning.countDown();
+ }
+ assertTrue(ddsAtSubmit.await(30, TimeUnit.SECONDS), "DirDeletingTask
never reached submitRequest");
+
+ CountDownLatch nextStoreStarted = new CountDownLatch(1);
+ Future<?> nextStore =
service.getExecutorService().submit(nextStoreStarted::countDown);
+ assertThat(nextStoreStarted.await(200, TimeUnit.MILLISECONDS))
+ .as("the next store must not open DB handles while workers are
submitting requests").isFalse();
+
+ AtomicBoolean purgeDone = new AtomicBoolean();
+ OzoneManagerProtocolProtos.OMRequest purgeRequest =
OzoneManagerProtocolProtos.OMRequest.newBuilder()
+ .setCmdType(OzoneManagerProtocolProtos.Type.SnapshotPurge)
+
.setSnapshotPurgeRequest(OzoneManagerProtocolProtos.SnapshotPurgeRequest.newBuilder()
+ .addSnapshotDBKeys(purgeKey))
+ .setClientId(ClientId.randomId().toString())
+ .build();
+ Future<?> purgeFuture = executor.submit(() -> {
+ OzoneManagerRatisUtils.submitRequest(om, purgeRequest,
ClientId.randomId(), 1L);
+ purgeDone.set(true);
+ return null;
+ });
+ GenericTestUtils.waitFor(
+ () -> unflushedTransactions.availablePermits() == 0 ||
purgeDone.get(), 50, 30000);
+
+ ddsMayProceed.countDown();
+ if (activeStore) {
+ GenericTestUtils.waitFor(unflushedTransactions::hasQueuedThreads, 50,
10000);
+ }
+ // An AOS worker can finish scanning and close its handle even while a
sibling waits for flush capacity.
+ finishScanning.countDown();
+ try {
+ ddsFuture.get(60, TimeUnit.SECONDS);
+ } catch (TimeoutException e) {
+ fail("DirectoryDeletingService retained a snapshot DB read lock across
its synchronous OM request");
+ }
+ purgeFuture.get(60, TimeUnit.SECONDS);
+ nextStore.get(10, TimeUnit.SECONDS);
+ om.awaitDoubleBufferFlush();
+ GenericTestUtils.waitFor(() ->
deletionThreadPool.getCompletedTaskCount() ==
+ completedTaskCountBefore + threadCount, 100, 5000);
+ assertThat(deletionThreadPool.getCompletedTaskCount() -
completedTaskCountBefore)
+ .as("store processing should use all configured deletion
workers").isEqualTo(threadCount);
+ assertThat(service.getDeletedDirsCount()).isEqualTo(threadCount);
+ } finally {
+ finishScanning.countDown();
+ ddsMayProceed.countDown();
+ executor.shutdownNow();
+ service.shutdown();
+ }
+ }
+
+ @Test
+ @DisplayName("DirectoryDeletingService releases snapshot workers when
submission is rejected")
+ void testRejectedSnapshotWorkerSubmission() throws Exception {
+ OzoneConfiguration conf = createConfAndInitValues(2);
+ OmTestManagers managers = new OmTestManagers(conf);
+ om = managers.getOzoneManager();
+ KeyManager keyManager = managers.getKeyManager();
+ OMMetadataManager metadataManager = managers.getMetadataManager();
+ OzoneManagerProtocol writeClient = managers.getWriteClient();
+ keyManager.getDirDeletingService().shutdown();
+
+ OMRequestTestUtils.addVolumeToOM(metadataManager,
+
OmVolumeArgs.newBuilder().setOwnerName("o").setAdminName("a").setVolume(volumeName).build());
+ OMRequestTestUtils.addBucketToOM(metadataManager,
+
OmBucketInfo.newBuilder().setVolumeName(volumeName).setBucketName(bucketName)
+
.setObjectID(1).setBucketLayout(BucketLayout.FILE_SYSTEM_OPTIMIZED).build());
+ String snapshotName = "snapshot";
+ writeClient.createSnapshot(volumeName, bucketName, snapshotName);
+ om.awaitDoubleBufferFlush();
+
+ String snapshotKey = SnapshotInfo.getTableKey(volumeName, bucketName,
snapshotName);
+ SnapshotInfo snapshotInfo =
metadataManager.getSnapshotInfoTable().get(snapshotKey);
+ assertNotNull(snapshotInfo);
+ awaitSnapshotFlushed(metadataManager, snapshotInfo);
+
+ DirectoryDeletingService service = new DirectoryDeletingService(1,
TimeUnit.HOURS, 1,
+ om, conf, 2, true);
+ RejectAfterFirstExecutor rejectionExecutor = new
RejectAfterFirstExecutor();
+ HddsWhiteboxTestUtils.setInternalState(service, "deletionThreadPool",
rejectionExecutor);
+ ExecutorService taskExecutor = Executors.newSingleThreadExecutor();
+ try {
+ Future<?> processFuture = taskExecutor.submit(() -> {
+ try (UncheckedAutoCloseableSupplier<OmSnapshot> snapshot =
om.getOmSnapshotManager().getActiveSnapshot(
+ volumeName, bucketName, snapshotName)) {
+ service.new
DirDeletingTask(snapshotInfo.getSnapshotId()).processDeletedDirsForStore(snapshotInfo,
+ snapshot.get().getKeyManager(), snapshot, 1, 1);
+ }
+ return null;
+ });
+ processFuture.get(10, TimeUnit.SECONDS);
+ assertThat(rejectionExecutor.getCompletedTaskCount()).isEqualTo(1);
Review Comment:
Please wait for the completed-task count with `GenericTestUtils.waitFor`, as
the other tests here do. The future can complete before the executor increments
the counter. Holding the worker in `afterExecute` reproduces `expected: 1, but
was: 0`.
--
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]