Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<OmKeyLocationInfoGroup> uncommittedGroups = new ArrayList<>();
// version not matters in the current logic of keyDeletingService,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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);
}

Expand Down Expand Up @@ -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<String, SnapshotInfo> 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<String, SnapshotInfo> snapshotInfoTable =
Expand Down Expand Up @@ -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);
Expand All @@ -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);
Expand Down Expand Up @@ -1694,6 +1765,33 @@ private OmKeyArgs createAndCommitKey(String volumeName,
return keyArg;
}

private static List<ContainerBlockID> getBlockIds(OmKeyInfo keyInfo) {
return keyInfo.getLatestVersionLocations().createLocationList().stream()
.map(location -> location.getBlockID().getContainerBlockID()).collect(Collectors.toList());
}

private List<ContainerBlockID> getSnapshotKeyBlocks(String volumeName, String bucketName, String snapshotName,
String dbKey) throws IOException {
try (UncheckedAutoCloseableSupplier<OmSnapshot> 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<ContainerBlockID> 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);
Expand Down