This is an automated email from the ASF dual-hosted git repository.
ivandika3 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/master by this push:
new a394802ec93 HDDS-16314. Filter out blocks still in use before
enqueuing overwritten key versions in S3 MPU complete (#11135)
a394802ec93 is described below
commit a394802ec930cc1c8faaea058599cdeac1f3cecf
Author: Peter Lee <[email protected]>
AuthorDate: Fri Sep 11 09:49:57 2026 +0800
HDDS-16314. Filter out blocks still in use before enqueuing overwritten key
versions in S3 MPU complete (#11135)
---
.../S3MultipartUploadCompleteRequest.java | 9 ++
.../TestS3MultipartUploadCompleteRequest.java | 149 ++++++++++++++++++++-
...estS3MultipartUploadCompleteRequestWithFSO.java | 48 +++++++
3 files changed, 204 insertions(+), 2 deletions(-)
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCompleteRequest.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCompleteRequest.java
index 1eccefccb9a..88e4f26b7c9 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCompleteRequest.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCompleteRequest.java
@@ -343,6 +343,15 @@ public OMClientResponse
validateAndUpdateCache(OzoneManager ozoneManager, Execut
if (keyToDelete != null && !omBucketInfo.getIsVersionEnabled()) {
RepeatedOmKeyInfo oldKeyVersionsToDelete = getOldVersionsToCleanUp(
keyToDelete, omBucketInfo.getObjectID(), trxnLogIndex);
+ // Remove any block from oldKeyVersionsToDelete that shares the same
+ // container ID and local ID with omKeyInfo blocks'.
+ // Otherwise, it causes data loss once those shared blocks are added
+ // to deletedTable and processed by KeyDeletingService for deletion.
+ // Unlike OMKeyCommitRequest, the returned filtered-block sizes are
+ // not applied to quota: part bytes were already counted at
+ // commit-part, and this transaction does not charge the new key,
+ // so there is no double-charge to correct.
+ filterOutBlocksStillInUse(omKeyInfo, oldKeyVersionsToDelete);
allKeyInfoToRemove.addAll(oldKeyVersionsToDelete.getOmKeyInfoList());
usedBytesDiff -= keyToDelete.getReplicatedSize();
} else {
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCompleteRequest.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCompleteRequest.java
index 27b7217b2dc..5cb8c8a31f0 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCompleteRequest.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCompleteRequest.java
@@ -26,22 +26,32 @@
import java.io.IOException;
import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
import java.util.HashMap;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
+import java.util.Set;
import java.util.UUID;
+import java.util.stream.Collectors;
+import org.apache.hadoop.hdds.client.ContainerBlockID;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.hdds.utils.db.BatchOperation;
import org.apache.hadoop.hdds.utils.db.Table;
import org.apache.hadoop.hdds.utils.db.cache.CacheKey;
import org.apache.hadoop.ozone.OzoneConsts;
import org.apache.hadoop.ozone.om.helpers.OmBucketInfo;
import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
+import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfo;
+import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfoGroup;
import org.apache.hadoop.ozone.om.helpers.OmMultipartKeyInfo;
import org.apache.hadoop.ozone.om.helpers.RepeatedOmKeyInfo;
import org.apache.hadoop.ozone.om.request.OMRequestTestUtils;
import org.apache.hadoop.ozone.om.response.OMClientResponse;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
+import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.KeyLocation;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.Part;
import org.apache.hadoop.util.Time;
@@ -95,8 +105,9 @@ public void testValidateAndUpdateCacheSuccess() throws
Exception {
// Do it twice to test overwrite
uploadId = checkValidateAndUpdateCacheSuccess(volumeName, bucketName,
keyName, customMetadata, tags);
- // After overwrite, one entry must be in delete table
- checkDeleteTableCount(volumeName, bucketName, keyName, 1, uploadId);
+ // The overwritten key has no blocks, so after filtering out blocks still
+ // in use nothing is enqueued to the delete table.
+ checkDeleteTableCount(volumeName, bucketName, keyName, 0, uploadId);
}
public void checkDeleteTableCount(String volumeName,
@@ -207,6 +218,125 @@ private String checkValidateAndUpdateCacheSuccess(String
volumeName,
return multipartUploadID;
}
+ @Test
+ public void testOverwrittenKeySharedBlocksNotEnqueuedForDeletion()
+ throws Exception {
+ String volumeName = UUID.randomUUID().toString();
+ String bucketName = UUID.randomUUID().toString();
+ String keyName = getKeyName();
+
+ OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+ omMetadataManager, getBucketLayout());
+
+ OMRequest initiateMPURequest = doPreExecuteInitiateMPU(volumeName,
+ bucketName, keyName);
+
+ S3InitiateMultipartUploadRequest s3InitiateMultipartUploadRequest =
+ getS3InitiateMultipartUploadReq(initiateMPURequest);
+
+ OMClientResponse omClientResponse =
+ s3InitiateMultipartUploadRequest.validateAndUpdateCache(ozoneManager,
1L);
+
+ long clientID = Time.now();
+ String multipartUploadID = omClientResponse.getOMResponse()
+ .getInitiateMultiPartUploadResponse().getMultipartUploadID();
+
+ // The part being completed commits this block; the key table entry
+ // planted below references the same container ID and local ID.
+ KeyLocation sharedKeyLocation = createKeyLocation(1000L, 100L);
+ OmKeyLocationInfo sharedBlock =
+ OmKeyLocationInfo.getFromProtobuf(sharedKeyLocation);
+ OmKeyLocationInfo oldOnlyBlock =
+ OmKeyLocationInfo.getFromProtobuf(createKeyLocation(2000L, 200L));
+
+ OMRequest commitMultipartRequest = doPreExecuteCommitMPU(volumeName,
+ bucketName, keyName, clientID, multipartUploadID, 1,
+ Collections.singletonList(sharedKeyLocation));
+
+ S3MultipartUploadCommitPartRequest s3MultipartUploadCommitPartRequest =
+ getS3MultipartUploadCommitReq(commitMultipartRequest);
+
+ // Add the part key to the open key table with the block pre-allocated.
+ addKeyToTable(volumeName, bucketName, keyName, clientID,
+ Collections.singletonList(sharedBlock));
+
+ s3MultipartUploadCommitPartRequest.validateAndUpdateCache(ozoneManager,
2L);
+
+ // The key being overwritten holds the shared block plus one block that
+ // only the old key references.
+ addCommittedKeyToTable(volumeName, bucketName, keyName,
+ Arrays.asList(sharedBlock, oldOnlyBlock));
+
+ List<Part> partList = new ArrayList<>();
+
+ String eTag = s3MultipartUploadCommitPartRequest.getOmRequest()
+ .getCommitMultiPartUploadRequest()
+ .getKeyArgs()
+ .getMetadataList()
+ .stream()
+ .filter(keyValue -> keyValue.getKey().equals(OzoneConsts.ETAG))
+ .findFirst().get().getValue();
+ partList.add(Part.newBuilder().setETag(eTag).setPartName(eTag)
+ .setPartNumber(1).build());
+
+ OMRequest completeMultipartRequest = doPreExecuteCompleteMPU(volumeName,
+ bucketName, keyName, multipartUploadID, partList);
+
+ S3MultipartUploadCompleteRequest s3MultipartUploadCompleteRequest =
+ getS3MultipartUploadCompleteReq(completeMultipartRequest);
+
+ omClientResponse =
+ s3MultipartUploadCompleteRequest.validateAndUpdateCache(ozoneManager,
3L);
+
+ assertEquals(OzoneManagerProtocolProtos.Status.OK,
+ omClientResponse.getOMResponse().getStatus());
+
+ BatchOperation batchOperation
+ = omMetadataManager.getStore().initBatchOperation();
+ omClientResponse.checkAndUpdateDB(omMetadataManager, batchOperation);
+ omMetadataManager.getStore().commitBatchOperation(batchOperation);
+
+ // The completed key must reference the committed part's block.
+ OmKeyInfo completedKeyInfo = omMetadataManager
+ .getKeyTable(s3MultipartUploadCompleteRequest.getBucketLayout())
+ .get(getOzoneDBKey(volumeName, bucketName, keyName));
+ assertNotNull(completedKeyInfo);
+ assertEquals(
+ Collections.singletonList(
+ sharedBlock.getBlockID().getContainerBlockID()),
+ completedKeyInfo.getLatestVersionLocations().createLocationList()
+ .stream().map(loc -> loc.getBlockID().getContainerBlockID())
+ .collect(Collectors.toList()));
+
+ // The overwritten key's block that the completed key still references
+ // must not be enqueued for deletion, while its other block must be.
+ Set<ContainerBlockID> deletedBlockIds = new HashSet<>();
+ List<Table.KeyValue<String, RepeatedOmKeyInfo>> rangeKVs
+ = omMetadataManager.getDeletedTable().getRangeKVs(null, 100,
+ getMultipartKey(volumeName, bucketName, keyName, multipartUploadID));
+ for (Table.KeyValue<String, RepeatedOmKeyInfo> rangeKV : rangeKVs) {
+ for (OmKeyInfo deletedKeyInfo : rangeKV.getValue().getOmKeyInfoList()) {
+ for (OmKeyLocationInfoGroup group
+ : deletedKeyInfo.getKeyLocationVersions()) {
+ for (OmKeyLocationInfo loc : group.createLocationList()) {
+ deletedBlockIds.add(loc.getBlockID().getContainerBlockID());
+ }
+ }
+ }
+ }
+ assertThat(deletedBlockIds)
+ .contains(oldOnlyBlock.getBlockID().getContainerBlockID())
+ .doesNotContain(sharedBlock.getBlockID().getContainerBlockID());
+ }
+
+ private static KeyLocation createKeyLocation(long containerID, long localID)
{
+ return KeyLocation.newBuilder()
+ .setBlockID(HddsProtos.BlockID.newBuilder()
+ .setContainerBlockID(HddsProtos.ContainerBlockID.newBuilder()
+ .setContainerID(containerID).setLocalID(localID).build()))
+ .setOffset(0).setLength(100).setCreateVersion(0).build();
+ }
+
@Test
public void
testValidateAndUpdateCacheUsesSchemaVersionOneBeforeFinalization()
throws Exception {
@@ -423,6 +553,21 @@ protected void addKeyToTable(String volumeName, String
bucketName,
keyName, clientID, RatisReplicationConfig.getInstance(ONE),
omMetadataManager);
}
+ protected void addKeyToTable(String volumeName, String bucketName,
+ String keyName, long clientID, List<OmKeyLocationInfo> locationList)
+ throws Exception {
+ OMRequestTestUtils.addKeyToTable(true, true, volumeName, bucketName,
+ keyName, clientID, RatisReplicationConfig.getInstance(ONE), 0L,
+ omMetadataManager, locationList, 0L);
+ }
+
+ protected void addCommittedKeyToTable(String volumeName, String bucketName,
+ String keyName, List<OmKeyLocationInfo> locationList) throws Exception {
+ OMRequestTestUtils.addKeyToTable(false, false, volumeName, bucketName,
+ keyName, 0L, RatisReplicationConfig.getInstance(ONE), 100L,
+ omMetadataManager, locationList, 0L);
+ }
+
protected String getMultipartKey(String volumeName, String bucketName,
String keyName, String multipartUploadID) throws IOException {
return omMetadataManager.getMultipartKey(volumeName,
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCompleteRequestWithFSO.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCompleteRequestWithFSO.java
index 9f0f5cbcdfd..3771c2c61d7 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCompleteRequestWithFSO.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCompleteRequestWithFSO.java
@@ -22,11 +22,13 @@
import java.io.IOException;
import java.util.ArrayList;
+import java.util.List;
import java.util.UUID;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
import org.apache.hadoop.ozone.OzoneConsts;
import org.apache.hadoop.ozone.om.helpers.BucketLayout;
import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
+import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfo;
import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfoGroup;
import org.apache.hadoop.ozone.om.helpers.OzoneFSUtils;
import org.apache.hadoop.ozone.om.request.OMRequestTestUtils;
@@ -89,6 +91,52 @@ protected void addKeyToTable(String volumeName, String
bucketName,
omMetadataManager);
}
+ @Override
+ protected void addKeyToTable(String volumeName, String bucketName,
+ String keyName, long clientID, List<OmKeyLocationInfo> locationList)
+ throws Exception {
+ long parentID = getParentID(volumeName, bucketName, keyName);
+ long txnId = 2;
+ long objectId = parentID + 1;
+
+ OmKeyInfo omKeyInfoFSO =
+ OMRequestTestUtils.createOmKeyInfo(volumeName, bucketName, keyName,
+ RatisReplicationConfig.getInstance(ONE), new
OmKeyLocationInfoGroup(0L, new ArrayList<>(), true))
+ .setObjectID(objectId)
+ .setParentObjectID(parentID)
+ .setUpdateID(txnId)
+ .build();
+ omKeyInfoFSO.appendNewBlocks(locationList, false);
+
+ // add key to openFileTable
+ String fileName = OzoneFSUtils.getFileName(keyName);
+ omKeyInfoFSO.setKeyName(fileName);
+ OMRequestTestUtils.addFileToKeyTable(true, false,
+ fileName, omKeyInfoFSO, clientID, omKeyInfoFSO.getObjectID(),
+ omMetadataManager);
+ }
+
+ @Override
+ protected void addCommittedKeyToTable(String volumeName, String bucketName,
+ String keyName, List<OmKeyLocationInfo> locationList) throws Exception {
+ long parentID = getParentID(volumeName, bucketName, keyName);
+
+ OmKeyInfo omKeyInfoFSO =
+ OMRequestTestUtils.createOmKeyInfo(volumeName, bucketName, keyName,
+ RatisReplicationConfig.getInstance(ONE), new
OmKeyLocationInfoGroup(0L, new ArrayList<>(), false))
+ .setObjectID(parentID + 2)
+ .setParentObjectID(parentID)
+ .setUpdateID(0L)
+ .build();
+ omKeyInfoFSO.appendNewBlocks(locationList, false);
+
+ // add key to fileTable
+ String fileName = OzoneFSUtils.getFileName(keyName);
+ omKeyInfoFSO.setKeyName(fileName);
+ OMRequestTestUtils.addFileToKeyTable(false, false,
+ fileName, omKeyInfoFSO, 0L, 100L, omMetadataManager);
+ }
+
private long getParentID(String volumeName, String bucketName,
String keyName) throws IOException {
final long volumeId = omMetadataManager.getVolumeId(volumeName);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]