From b7af09c26402a6fb454f971c749a0423533eecb4 Mon Sep 17 00:00:00 2001 From: Siyao Meng <50227127+smengcl@users.noreply.github.com> Date: Wed, 7 Oct 2026 04:13:44 -0700 Subject: [PATCH] HDDS-16755. Uncommitted block cleanup can delete a block that a snapshot of an hsync'ed key still references A block committed by an hsync and left out of a later hsync or of the final commit is released as an uncommitted block when the key is closed. Its pseudo key in the deleted table had object ID 0, so ReclaimableKeyFilter never matched it with the key in the previous snapshot and KeyDeletingService deleted the block, even when a snapshot taken after the hsync still listed it. The pseudo key now keeps the object ID of the key. Together with the hsync metadata it already carries, this retains it while the previous snapshot holds an hsync'ed version of the key. The deleted table row name is still derived from the transaction index, so rows do not collide. Tested with TestKeyDeletingService (new testUncommittedBlockOfHsyncedKeyRetainedBySnapshot), TestOMKeyCommitRequest and TestOMKeyCommitRequestWithFSO. The new test fails without the fix, and also with the object ID kept but the hsync metadata removed from the pseudo key. Co-Authored-By: Claude Opus 5.5 --- .../ozone/om/request/key/OMKeyRequest.java | 11 +- .../request/key/TestOMKeyCommitRequest.java | 6 +- .../om/service/TestKeyDeletingService.java | 118 ++++++++++++++++-- 3 files changed, 118 insertions(+), 17 deletions(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java index d674976c1579..ffaf152e39da 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java @@ -21,7 +21,6 @@ import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.BlockTokenSecretProto.AccessModeProto.WRITE; import static org.apache.hadoop.ozone.OzoneAcl.AclScope.ACCESS; import static org.apache.hadoop.ozone.OzoneAcl.AclScope.DEFAULT; -import static org.apache.hadoop.ozone.OzoneConsts.OBJECT_ID_RECLAIM_BLOCKS; import static org.apache.hadoop.ozone.OzoneConsts.OZONE_URI_DELIMITER; import static org.apache.hadoop.ozone.om.exceptions.OMException.ResultCodes.BUCKET_NOT_FOUND; import static org.apache.hadoop.ozone.om.exceptions.OMException.ResultCodes.INVALID_KEY_NAME; @@ -1223,11 +1222,11 @@ protected OmKeyInfo wrapUncommittedBlocksAsPseudoKey( } LOG.debug("Detect allocated but uncommitted blocks {} in key {}.", uncommitted, omKeyInfo.getKeyName()); - OmKeyInfo pseudoKeyInfo = omKeyInfo.toBuilder() - .setObjectID(OBJECT_ID_RECLAIM_BLOCKS) - .build(); - // This is a special marker to indicate that SnapshotDeletingService - // can reclaim this key's blocks unconditionally. + // The pseudo key keeps the object ID and the hsync metadata of the key, so callers wrap before they remove the + // latter. A block committed by an earlier hsync may be referenced by a snapshot, and key deletion retains the + // pseudo key while the previous snapshot holds an hsync'ed version of the key with this object ID. That version + // may not list the block itself, as an hsync between two snapshots can drop a block that the older one references. + OmKeyInfo pseudoKeyInfo = omKeyInfo.toBuilder().build(); // TODO dataSize of pseudoKey is not real here List uncommittedGroups = new ArrayList<>(); // version not matters in the current logic of keyDeletingService, diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCommitRequest.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCommitRequest.java index b5152082c7d2..f14b0006aff2 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCommitRequest.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCommitRequest.java @@ -883,7 +883,9 @@ public void testValidateAndUpdateCacheOnOverwriteWithUncommittedBlocks() throws OMKeyCommitRequest omKeyCommitRequest = getOmKeyCommitRequest(modifiedOmRequest); - addKeyToOpenKeyTable(allocatedBlockList); + // A non-zero object ID, to tell whether the pseudo key of the uncommitted blocks keeps it + addKeyToOpenKeyTable(allocatedBlockList, OMRequestTestUtils.createOmKeyInfo(volumeName, bucketName, keyName, + replicationConfig, new OmKeyLocationInfoGroup(version, new ArrayList<>(), false)).setObjectID(100L)); OMClientResponse omClientResponse = omKeyCommitRequest.validateAndUpdateCache(ozoneManager, 102L); @@ -925,6 +927,8 @@ public void testValidateAndUpdateCacheOnOverwriteWithUncommittedBlocks() throws assertEquals(DEFAULT_COMMIT_BLOCK_SIZE, overwrittenKey.getLatestVersionLocations().createLocationList().size()); assertEquals(allocatedKeyLocationList.size() - committedKeyLocationList.size(), uncommittedPseudoKey.getLatestVersionLocations().createLocationList().size()); + // A snapshot's version of the key is matched by object ID + assertThat(uncommittedPseudoKey.getObjectID()).isNotZero().isEqualTo(omKeyInfo.getObjectID()); // flush response content to db BatchOperation batchOperation = omMetadataManager.getStore().initBatchOperation(); diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestKeyDeletingService.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestKeyDeletingService.java index 900d64315f6e..afdb3e8bd707 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestKeyDeletingService.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestKeyDeletingService.java @@ -20,6 +20,8 @@ import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_CONTAINER_REPORT_INTERVAL; import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE; import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_BLOCK_DELETING_SERVICE_INTERVAL; +import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_FS_HSYNC_ENABLED; +import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_HBASE_ENHANCEMENTS_ALLOWED; import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_MANAGER_STRIPED_LOCK_SIZE_PREFIX; import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_SNAPSHOT_DELETING_SERVICE_INTERVAL; import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_DIR_DELETING_SERVICE_INTERVAL; @@ -52,6 +54,7 @@ import java.io.IOException; import java.io.UncheckedIOException; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collection; import java.util.Collections; import java.util.HashMap; @@ -74,6 +77,7 @@ import java.util.stream.IntStream; import org.apache.commons.lang3.RandomStringUtils; import org.apache.hadoop.hdds.client.BlockID; +import org.apache.hadoop.hdds.client.ContainerBlockID; import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.client.StandaloneReplicationConfig; import org.apache.hadoop.hdds.conf.OzoneConfiguration; @@ -191,6 +195,8 @@ private void createConfig(File testDir, int delintervalMs) { conf.setTimeDuration(HDDS_CONTAINER_REPORT_INTERVAL, 200, TimeUnit.MILLISECONDS); conf.setBoolean(OZONE_SNAPSHOT_DEEP_CLEANING_ENABLED, true); + conf.setBoolean(OZONE_HBASE_ENHANCEMENTS_ALLOWED, true); + conf.setBoolean(OZONE_FS_HSYNC_ENABLED, true); conf.setQuietMode(false); } @@ -778,6 +784,67 @@ public void testKeyDeletingServiceWithDeepCleanedSnapshots() throws Exception { verify(omSnapshotManager, Mockito.never()).getActiveSnapshot(any(), any(), any()); } + /** + * A file is hsync'ed with blocks A and B and captured by a snapshot. The writer replaces B with block C and hsyncs, + * a second snapshot captures A and C, and the file is closed, which releases B as an uncommitted block. B is not + * reclaimed until the first snapshot is deleted, even while the previous snapshot is the one that does not list it. + */ + @Test + void testUncommittedBlockOfHsyncedKeyRetainedBySnapshot() throws Exception { + keyDeletingService.suspend(); + final String volumeName = getTestName(); + final String bucketName = uniqueObjectName("bucket"); + final String keyName = uniqueObjectName("key"); + createVolumeAndBucket(volumeName, bucketName, false); + OmKeyArgs keyArg = newKeyArgs(volumeName, bucketName, keyName); + OpenKeySession session = writeClient.openKey(keyArg); + OmKeyLocationInfo locationA = + session.getKeyInfo().getLatestVersionLocations().getBlocksLatestVersionOnly().get(0); + OmKeyLocationInfo locationB = writeClient.allocateBlock(keyArg, session.getId(), new ExcludeList()); + keyArg.setLocationInfoList(Arrays.asList(locationA, locationB)); + keyArg.setDataSize(locationA.getLength() + locationB.getLength()); + writeClient.hsyncKey(keyArg, session.getId()); + String snap1 = uniqueObjectName("snap"); + writeClient.createSnapshot(volumeName, bucketName, snap1); + OmKeyLocationInfo locationC = writeClient.allocateBlock(keyArg, session.getId(), new ExcludeList()); + keyArg.setLocationInfoList(Arrays.asList(locationA, locationC)); + keyArg.setDataSize(locationA.getLength() + locationC.getLength()); + writeClient.hsyncKey(keyArg, session.getId()); + String snap2 = uniqueObjectName("snap"); + writeClient.createSnapshot(volumeName, bucketName, snap2); + writeClient.commitKey(keyArg, session.getId()); + om.awaitDoubleBufferFlush(); + ContainerBlockID blockA = locationA.getBlockID().getContainerBlockID(); + ContainerBlockID blockB = locationB.getBlockID().getContainerBlockID(); + ContainerBlockID blockC = locationC.getBlockID().getContainerBlockID(); + assertThat(Arrays.asList(blockA, blockB, blockC)).doesNotHaveDuplicates(); + + String dbKey = metadataManager.getOzoneKey(volumeName, bucketName, keyName); + assertThat(getSnapshotKeyBlocks(volumeName, bucketName, snap1, dbKey)).containsExactly(blockA, blockB); + assertThat(getSnapshotKeyBlocks(volumeName, bucketName, snap2, dbKey)).containsExactly(blockA, blockC); + assertThat(getBlocksPendingDeletion(dbKey)).containsExactly(blockB); + + // The previous snapshot does not list B, but it holds an hsync'ed version of the key, so an older one may. + long runCount = getRunCount(); + keyDeletingService.resume(); + GenericTestUtils.waitFor(() -> getRunCount() > runCount + 5, 100, 10000); + assertThat(getBlocksPendingDeletion(dbKey)).containsExactly(blockB); + + // The first snapshot references B. + Table snapshotInfoTable = metadataManager.getSnapshotInfoTable(); + long snapshotCount = metadataManager.countRowsInTable(snapshotInfoTable); + writeClient.deleteSnapshot(volumeName, bucketName, snap2); + assertTableRowCount(snapshotInfoTable, snapshotCount - 1, metadataManager); + long runCountWithoutSnap2 = getRunCount(); + GenericTestUtils.waitFor(() -> getRunCount() > runCountWithoutSnap2 + 5, 100, 10000); + assertThat(getBlocksPendingDeletion(dbKey)).containsExactly(blockB); + + writeClient.deleteSnapshot(volumeName, bucketName, snap1); + GenericTestUtils.waitFor(() -> getBlocksPendingDeletion(dbKey).isEmpty(), 1000, 120000); + assertThat(getBlockIds(metadataManager.getKeyTable(BucketLayout.DEFAULT).get(dbKey))) + .containsExactly(blockA, blockC); + } + @Test void testSnapshotExclusiveSize() throws Exception { Table snapshotInfoTable = @@ -1635,6 +1702,19 @@ private void renameKey(String volumeName, writeClient.renameKey(keyArg, toKeyName); } + private static OmKeyArgs newKeyArgs(String volumeName, String bucketName, String keyName) { + return new OmKeyArgs.Builder() + .setVolumeName(volumeName) + .setBucketName(bucketName) + .setKeyName(keyName) + .setAcls(Collections.emptyList()) + .setReplicationConfig(RatisReplicationConfig.getInstance(THREE)) + .setDataSize(1000L) + .setLocationInfoList(new ArrayList<>()) + .setOwnerName("user" + RandomStringUtils.secure().nextNumeric(5)) + .build(); + } + private OmKeyArgs createAndCommitKey(String volumeName, String bucketName, String keyName, int numBlocks) throws IOException { return createAndCommitKey(volumeName, bucketName, keyName, numBlocks, 0, this.writeClient); @@ -1649,16 +1729,7 @@ private OmKeyArgs createAndCommitKey(String volumeName, String bucketName, String keyName, int numBlocks, int numUncommitted, OzoneManagerProtocol customWriteClient) throws IOException { - OmKeyArgs keyArg = new OmKeyArgs.Builder() - .setVolumeName(volumeName) - .setBucketName(bucketName) - .setKeyName(keyName) - .setAcls(Collections.emptyList()) - .setReplicationConfig(RatisReplicationConfig.getInstance(THREE)) - .setDataSize(1000L) - .setLocationInfoList(new ArrayList<>()) - .setOwnerName("user" + RandomStringUtils.secure().nextNumeric(5)) - .build(); + OmKeyArgs keyArg = newKeyArgs(volumeName, bucketName, keyName); // Open and Commit the Key in the Key Manager. OpenKeySession session = customWriteClient.openKey(keyArg); @@ -1694,6 +1765,33 @@ private OmKeyArgs createAndCommitKey(String volumeName, return keyArg; } + private static List getBlockIds(OmKeyInfo keyInfo) { + return keyInfo.getLatestVersionLocations().createLocationList().stream() + .map(location -> location.getBlockID().getContainerBlockID()).collect(Collectors.toList()); + } + + private List getSnapshotKeyBlocks(String volumeName, String bucketName, String snapshotName, + String dbKey) throws IOException { + try (UncheckedAutoCloseableSupplier snapshot = + om.getOmSnapshotManager().getSnapshot(volumeName, bucketName, snapshotName)) { + return getBlockIds(snapshot.get().getMetadataManager().getKeyTable(BucketLayout.DEFAULT).get(dbKey)); + } + } + + /** + * Returns the blocks in the deleted table entries of the given key. + */ + private List getBlocksPendingDeletion(String dbKey) { + try { + return metadataManager.getDeletedTable().getRangeKVs(null, 10, dbKey).stream() + .flatMap(kv -> kv.getValue().getOmKeyInfoList().stream()) + .flatMap(keyInfo -> getBlockIds(keyInfo).stream()) + .collect(Collectors.toList()); + } catch (IOException e) { + throw new UncheckedIOException(e); + } + } + private long getDeletedKeyCount() { final long count = keyDeletingService.getDeletedKeyCount().get(); LOG.debug("KeyDeletingService deleted keys: {}", count);