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);