This is an automated email from the ASF dual-hosted git repository.
sarvekshayr 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 ca0715a1a0a HDDS-15949. Handle split MPU part counting for abort batch
sizing and container-key-mapping CLI (#10849)
ca0715a1a0a is described below
commit ca0715a1a0ad460e54967756f458b67222579720
Author: Abhishek Pal <[email protected]>
AuthorDate: Tue Jul 28 10:11:40 2026 +0530
HDDS-15949. Handle split MPU part counting for abort batch sizing and
container-key-mapping CLI (#10849)
---
.../ozone/debug/om/ContainerToKeyMapping.java | 63 +++++-
.../ozone/debug/om/TestContainerToKeyMapping.java | 59 ++++++
.../hadoop/ozone/om/OmMetadataManagerImpl.java | 9 +-
.../om/request/util/OMMultipartUploadUtils.java | 10 +
.../ozone/om/service/KeyLifecycleService.java | 20 +-
.../hadoop/ozone/om/TestOmMetadataManager.java | 225 +++++++++++++++++++++
.../ozone/om/service/TestKeyLifecycleService.java | 126 ++++++++++++
7 files changed, 501 insertions(+), 11 deletions(-)
diff --git
a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/om/ContainerToKeyMapping.java
b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/om/ContainerToKeyMapping.java
index 589bd98aace..14559dad3f1 100644
---
a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/om/ContainerToKeyMapping.java
+++
b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/om/ContainerToKeyMapping.java
@@ -52,7 +52,11 @@
import org.apache.hadoop.ozone.om.helpers.OmBucketInfo;
import org.apache.hadoop.ozone.om.helpers.OmDirectoryInfo;
import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
+import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfoGroup;
import org.apache.hadoop.ozone.om.helpers.OmMultipartKeyInfo;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartPartInfo;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartPartKey;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartUpload;
import org.apache.hadoop.ozone.om.helpers.OmVolumeArgs;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.PartKeyInfo;
import picocli.CommandLine;
@@ -100,6 +104,7 @@ public class ContainerToKeyMapping extends
AbstractSubcommand implements Callabl
private Table<String, OmKeyInfo> openFileTable;
private Table<String, OmKeyInfo> openKeyTable;
private Table<String, OmMultipartKeyInfo> multipartInfoTable;
+ private Table<OmMultipartPartKey, OmMultipartPartInfo> multipartPartsTable;
private DBStore dirTreeDbStore;
private Table<Long, String> dirTreeTable;
// Cache volume IDs to avoid repeated lookups
@@ -138,6 +143,7 @@ public Void call() throws Exception {
openFileTable = OMDBDefinition.OPEN_FILE_TABLE_DEF.getTable(omDbStore,
CacheType.NO_CACHE);
openKeyTable = OMDBDefinition.OPEN_KEY_TABLE_DEF.getTable(omDbStore,
CacheType.NO_CACHE);
multipartInfoTable =
OMDBDefinition.MULTIPART_INFO_TABLE_DEF.getTable(omDbStore, CacheType.NO_CACHE);
+ multipartPartsTable =
OMDBDefinition.MULTIPART_PARTS_TABLE_DEF.getTable(omDbStore,
CacheType.NO_CACHE);
retrieve(dbPath, writer, containerIDs);
} catch (Exception e) {
@@ -313,11 +319,18 @@ private void processMultipartUpload(Set<Long>
containerIds, Map<Long, List<Strin
String dbKey = entry.getKey();
OmMultipartKeyInfo mpuInfo = entry.getValue();
- // Collect all target containers that have parts of this MPU
+ // Collect all target containers that have parts of this MPU. Legacy
+ // (schemaVersion 0) MPUs embed their parts in the multipartInfoTable
+ // value, whereas split (schemaVersion 1) MPUs store each part as a
+ // separate row in the multipartPartsTable.
Set<Long> matchedContainers = new HashSet<>();
- for (PartKeyInfo partKeyInfo : mpuInfo.getPartKeyInfoMap()) {
- OmKeyInfo partKey =
OmKeyInfo.getFromProtobuf(partKeyInfo.getPartKeyInfo());
- matchedContainers.addAll(getKeyContainers(partKey, containerIds));
+ if (mpuInfo.getSchemaVersion() ==
OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION) {
+ matchedContainers.addAll(getSplitPartContainers(dbKey,
containerIds));
+ } else {
+ for (PartKeyInfo partKeyInfo : mpuInfo.getPartKeyInfoMap()) {
+ OmKeyInfo partKey =
OmKeyInfo.getFromProtobuf(partKeyInfo.getPartKeyInfo());
+ matchedContainers.addAll(getKeyContainers(partKey, containerIds));
+ }
}
if (!matchedContainers.isEmpty()) {
@@ -332,8 +345,48 @@ private void processMultipartUpload(Set<Long>
containerIds, Map<Long, List<Strin
}
private Set<Long> getKeyContainers(OmKeyInfo keyInfo, Set<Long>
targetContainerIds) {
+ return getContainers(keyInfo.getKeyLocationVersions(), targetContainerIds);
+ }
+
+ /**
+ * Scans the split multipartPartsTable for all parts belonging to the given
+ * multipart upload (its uploadId is the last path component of the
+ * multipartInfoTable db key) and returns the target containers referenced by
+ * those parts' block locations.
+ */
+ private Set<Long> getSplitPartContainers(String multipartInfoDbKey,
Set<Long> targetContainerIds) {
+ Set<Long> matchedContainers = new HashSet<>();
+ String uploadId;
+ try {
+ uploadId = OmMultipartUpload.from(multipartInfoDbKey).getUploadId();
+ } catch (IllegalArgumentException e) {
+ err().println("Invalid multipartInfoTable key " + multipartInfoDbKey +
", " + e);
+ return matchedContainers;
+ }
+ OmMultipartPartKey prefix = OmMultipartPartKey.prefix(uploadId);
+ try (TableIterator<OmMultipartPartKey, Table.KeyValue<OmMultipartPartKey,
OmMultipartPartInfo>>
+ partIterator = multipartPartsTable.iterator(prefix)) {
+ while (partIterator.hasNext()) {
+ Table.KeyValue<OmMultipartPartKey, OmMultipartPartInfo> partEntry =
partIterator.next();
+ OmMultipartPartKey partKey = partEntry.getKey();
+ // Prefix iteration can overshoot into the next upload's rows; stop
then.
+ if (!uploadId.equals(partKey.getUploadId())) {
+ break;
+ }
+ if (partKey.hasPartNumber()) {
+ matchedContainers.addAll(
+ getContainers(partEntry.getValue().getKeyLocationInfos(),
targetContainerIds));
+ }
+ }
+ } catch (Exception e) {
+ err().println("Exception occurred reading multipartPartsTable for upload
" + uploadId + ", " + e);
+ }
+ return matchedContainers;
+ }
+
+ private Set<Long> getContainers(List<OmKeyLocationInfoGroup>
locationVersions, Set<Long> targetContainerIds) {
Set<Long> keyContainers = new HashSet<>();
- keyInfo.getKeyLocationVersions().forEach(
+ locationVersions.forEach(
e -> e.getLocationList().forEach(
blk -> {
long cid = blk.getBlockID().getContainerID();
diff --git
a/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/om/TestContainerToKeyMapping.java
b/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/om/TestContainerToKeyMapping.java
index 4cad62cd719..3a2e396fd5f 100644
---
a/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/om/TestContainerToKeyMapping.java
+++
b/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/om/TestContainerToKeyMapping.java
@@ -31,6 +31,7 @@
import org.apache.hadoop.hdds.client.StandaloneReplicationConfig;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
+import org.apache.hadoop.ozone.OzoneConsts;
import org.apache.hadoop.ozone.debug.OzoneDebug;
import org.apache.hadoop.ozone.om.OMMetadataManager;
import org.apache.hadoop.ozone.om.OmMetadataManagerImpl;
@@ -41,6 +42,8 @@
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.OmMultipartPartInfo;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartPartKey;
import org.apache.hadoop.ozone.om.helpers.OmVolumeArgs;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.KeyInfo;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.PartKeyInfo;
@@ -76,6 +79,7 @@ public class TestContainerToKeyMapping {
private static final long CONTAINER_ID_2 = 2L;
private static final long CONTAINER_ID_3 = 3L;
private static final long CONTAINER_ID_4 = 4L;
+ private static final long CONTAINER_ID_5 = 5L;
private static final long UNREFERENCED_FILE_ID = 500L;
private static final long MISSING_DIR_ID = 999L; // Non-existent parent
private static final long OPEN_FILE_ID = 600L;
@@ -189,6 +193,18 @@ public void
testContainerToKeyMappingWithMPUOnlyFileNames() {
assertThat(output).contains("/vol1/obs-bucket/mpuKey/test-upload-id");
}
+ @Test
+ public void testContainerToKeyMappingWithSplitSchemaMPU() {
+ int exitCode = execute("--containers", String.valueOf(CONTAINER_ID_5),
"--in-progress");
+ assertEquals(0, exitCode);
+
+ String output = outWriter.toString();
+
+ assertThat(output).contains("\"" + CONTAINER_ID_5 + "\"");
+ assertThat(output).contains("\"openKeys\"");
+
assertThat(output).contains("/vol1/obs-bucket/splitMpuKey/split-upload-id");
+ }
+
@Test
public void testNonExistentContainer() {
@@ -292,6 +308,7 @@ private void createTestData() throws Exception {
// Create MPU (multipart upload) for OBS bucket with parts in container 5
createMultipartUpload();
+ createSplitSchemaMultipartUpload();
}
/**
@@ -339,6 +356,46 @@ private void createMultipartUpload() throws Exception {
omMetadataManager.getMultipartInfoTable().put(mpuKey, mpuInfo);
}
+ /**
+ * Helper method to create a split-schema multipart upload with parts in
+ * multipartPartsTable.
+ */
+ private void createSplitSchemaMultipartUpload() throws Exception {
+ String mpuKeyName = "splitMpuKey";
+ String uploadId = "split-upload-id";
+
+ OmKeyInfo part1Info = new OmKeyInfo.Builder(createOBSKeyInfo(
+ mpuKeyName + "/" + uploadId + "/part-1", MPU_PART1_ID + 10,
CONTAINER_ID_5))
+ .addMetadata(OzoneConsts.ETAG, "etag-1")
+ .build();
+ OmMultipartPartInfo partInfo1 = OmMultipartPartInfo.from(
+ mpuKeyName + "/" + uploadId + "/part-1", 1, part1Info);
+
+ OmKeyInfo part2Info = new OmKeyInfo.Builder(createOBSKeyInfo(
+ mpuKeyName + "/" + uploadId + "/part-2", MPU_PART2_ID + 10,
CONTAINER_ID_5))
+ .addMetadata(OzoneConsts.ETAG, "etag-2")
+ .build();
+ OmMultipartPartInfo partInfo2 = OmMultipartPartInfo.from(
+ mpuKeyName + "/" + uploadId + "/part-2", 2, part2Info);
+
+
omMetadataManager.getMultipartPartsTable().put(OmMultipartPartKey.of(uploadId,
1), partInfo1);
+
omMetadataManager.getMultipartPartsTable().put(OmMultipartPartKey.of(uploadId,
2), partInfo2);
+
+ OmMultipartKeyInfo mpuInfo = new OmMultipartKeyInfo.Builder()
+ .setUploadID(uploadId)
+ .setCreationTime(System.currentTimeMillis())
+
.setReplicationConfig(StandaloneReplicationConfig.getInstance(HddsProtos.ReplicationFactor.ONE))
+ .setSchemaVersion(OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION)
+ .setObjectID(MPU_KEY_ID + 10)
+ .setParentID(0)
+ .setUpdateID(1)
+ .build();
+
+ String mpuKey = omMetadataManager.getMultipartKey(
+ VOLUME_NAME, OBS_BUCKET_NAME, mpuKeyName, uploadId);
+ omMetadataManager.getMultipartInfoTable().put(mpuKey, mpuInfo);
+ }
+
/**
* Helper method to create OmKeyInfo with a block in specified container
(FSO).
*/
@@ -386,6 +443,8 @@ private OmKeyInfo createOBSKeyInfo(String keyName, long
objectId, long container
.setDataSize(1024)
.setObjectID(objectId)
.setUpdateID(1)
+ .setCreationTime(System.currentTimeMillis())
+ .setModificationTime(System.currentTimeMillis())
.addOmKeyLocationInfoGroup(locationGroup)
.build();
}
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OmMetadataManagerImpl.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OmMetadataManagerImpl.java
index f31815bc208..99b1c18d8ec 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OmMetadataManagerImpl.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OmMetadataManagerImpl.java
@@ -1559,8 +1559,13 @@ public List<ExpiredMultipartUploadsBucket>
getExpiredMultipartUploads(
expiredMPUs.get(mapKey)
.addMultipartUploads(builder.setName(dbMultipartInfoKey)
.build());
- numParts += omMultipartKeyInfo.getPartKeyInfoMap().size();
- // TODO: Add the expired part handling from the new table when the
complete flow is done
+
+ if (omMultipartKeyInfo.getSchemaVersion()
+ == OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION) {
+ numParts += OMMultipartUploadUtils.countParts(this,
expiredMultipartUpload.getUploadId());
+ } else {
+ numParts += omMultipartKeyInfo.getPartKeyInfoMap().size();
+ }
}
}
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/util/OMMultipartUploadUtils.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/util/OMMultipartUploadUtils.java
index 598f7dc9410..d6fd32af51d 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/util/OMMultipartUploadUtils.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/util/OMMultipartUploadUtils.java
@@ -165,6 +165,16 @@ public static SortedMap<Integer, OmMultipartPartInfo>
scanParts(
return parts;
}
+ /**
+ * Count the multipart parts belonging to a given upload in the split
+ * multipartPartsTable, honouring cache tombstones and pending commits. The
+ * count therefore matches the set of parts a subsequent abort/cleanup would
+ * process, which makes it suitable for batch sizing.
+ */
+ public static int countParts(OMMetadataManager omMetadataManager, String
uploadId) throws IOException {
+ return scanParts(omMetadataManager, uploadId).size();
+ }
+
public static List<OmMultipartPartKey> getPartKeys(String uploadId,
SortedMap<Integer, OmMultipartPartInfo> parts) {
List<OmMultipartPartKey> partKeys = new ArrayList<>(parts.size());
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/KeyLifecycleService.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/KeyLifecycleService.java
index ead8acc7cd5..5310f5b4837 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/KeyLifecycleService.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/KeyLifecycleService.java
@@ -1138,11 +1138,23 @@ private void processMultipartUploads(OmBucketInfo
bucketInfo, List<OmLCRule> rul
abortExpiredMultipartUploadsAndClear(bucketInfo,
expiredUploads);
}
- // Get part count for this MPU (at least 1 even if no parts
uploaded yet)
- int partCount = Math.max(1,
mpuKeyInfo.getPartKeyInfoMap().size());
- expiredUploads.add(upload, partCount);
+ // Split-schema MPUs keep parts in multipartPartsTable (the
embedded map
+ // is empty); legacy MPUs use the embedded map. An MPU with no
uploaded
+ // parts is valid (S3 allows aborting it with an empty parts
list).
+ int uploadedParts;
+ try {
+ uploadedParts = mpuKeyInfo.getSchemaVersion()
+ == OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION
+ ? OMMultipartUploadUtils.countParts(omMetadataManager,
upload.getUploadId())
+ : mpuKeyInfo.getPartKeyInfoMap().size();
+ } catch (IOException e) {
+ LOG.warn("Failed to count parts for MPU {}/{}/{} uploadId {},
skipping",
+ volumeName, bucketName, keyName, upload.getUploadId(), e);
+ break;
+ }
+ expiredUploads.add(upload, uploadedParts);
LOG.debug("Multipart upload {}/{}/{} with uploadId {} ({} parts)
will be aborted",
- volumeName, bucketName, keyName, upload.getUploadId(),
partCount);
+ volumeName, bucketName, keyName, upload.getUploadId(),
uploadedParts);
break;
}
}
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestOmMetadataManager.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestOmMetadataManager.java
index 754fbc1ed21..1ffcc32f0ac 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestOmMetadataManager.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestOmMetadataManager.java
@@ -61,6 +61,7 @@
import static org.junit.jupiter.params.provider.Arguments.arguments;
import java.io.File;
+import java.io.IOException;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Arrays;
@@ -75,6 +76,7 @@
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import java.util.stream.Stream;
+import org.apache.hadoop.hdds.client.BlockID;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.protocol.StorageType;
@@ -89,8 +91,11 @@
import org.apache.hadoop.ozone.om.helpers.ListOpenFilesResult;
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.OmMultipartPartInfo;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartPartKey;
import org.apache.hadoop.ozone.om.helpers.OmMultipartUpload;
import org.apache.hadoop.ozone.om.helpers.OmVolumeArgs;
import org.apache.hadoop.ozone.om.helpers.OpenKeySession;
@@ -1047,6 +1052,226 @@ public void testGetExpiredMPUs() throws Exception {
assertThat(expiredMPUs).containsAll(names);
}
+ @Test
+ public void testGetExpiredMPUsSplitSchema() throws Exception {
+ final String bucketName = UUID.randomUUID().toString();
+ final String volumeName = UUID.randomUUID().toString();
+ final int numExpiredMPUs = 4;
+ final int numUnexpiredMPUs = 1;
+ final int numPartsPerMPU = 5;
+ final long expireThresholdMillis = ozoneConfiguration.getTimeDuration(
+ OZONE_OM_MPU_EXPIRE_THRESHOLD,
+ OZONE_OM_MPU_EXPIRE_THRESHOLD_DEFAULT,
+ TimeUnit.MILLISECONDS);
+
+ final Duration expireThreshold = Duration.ofMillis(expireThresholdMillis);
+
+ final long expiredMPUCreationTime =
+ expireThreshold.negated().plusMillis(Time.now()).toMillis();
+
+ Set<String> expiredMPUs = new HashSet<>();
+ for (int i = 0; i < numExpiredMPUs + numUnexpiredMPUs; i++) {
+ final long creationTime = i < numExpiredMPUs ?
+ expiredMPUCreationTime : Time.now();
+
+ String uploadId = OMMultipartUploadUtils.getMultipartUploadId();
+ final OmMultipartKeyInfo mpuKeyInfo = new OmMultipartKeyInfo.Builder(
+ OMRequestTestUtils.createOmMultipartKeyInfo(uploadId, creationTime,
+ HddsProtos.ReplicationType.RATIS,
+ HddsProtos.ReplicationFactor.ONE, 0L))
+
.setSchemaVersion(OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION)
+ .build();
+
+ String keyName = "expired-split" + i;
+ final OmKeyInfo keyInfo = OMRequestTestUtils.createOmKeyInfo(volumeName,
+ bucketName, keyName, RatisReplicationConfig.getInstance(ONE))
+ .setCreationTime(creationTime)
+ .build();
+
+ for (int j = 1; j <= numPartsPerMPU; j++) {
+ addSplitSchemaPart(uploadId, j);
+ }
+
+ final String mpuDbKey = OMRequestTestUtils.addMultipartInfoToTable(
+ false, keyInfo, mpuKeyInfo, 0L, omMetadataManager);
+
+ expiredMPUs.add(mpuDbKey);
+ }
+
+ List<ExpiredMultipartUploadsBucket> someExpiredMPUs =
+ omMetadataManager.getExpiredMultipartUploads(
+ expireThreshold,
+ (numExpiredMPUs * numPartsPerMPU) - (numPartsPerMPU));
+ List<String> names = getMultipartKeyNames(someExpiredMPUs);
+ assertEquals(numExpiredMPUs - 1, names.size());
+ assertThat(expiredMPUs).containsAll(names);
+
+ List<ExpiredMultipartUploadsBucket> allExpiredMPUs =
+ omMetadataManager.getExpiredMultipartUploads(expireThreshold,
+ (numExpiredMPUs * numPartsPerMPU));
+ names = getMultipartKeyNames(allExpiredMPUs);
+ assertEquals(numExpiredMPUs, names.size());
+ assertThat(expiredMPUs).containsAll(names);
+ }
+
+ @Test
+ public void testGetExpiredMPUsMixedSchema() throws Exception {
+ final String bucketName = UUID.randomUUID().toString();
+ final String volumeName = UUID.randomUUID().toString();
+ final int numPartsPerMPU = 3;
+ final long expireThresholdMillis = ozoneConfiguration.getTimeDuration(
+ OZONE_OM_MPU_EXPIRE_THRESHOLD,
+ OZONE_OM_MPU_EXPIRE_THRESHOLD_DEFAULT,
+ TimeUnit.MILLISECONDS);
+ final Duration expireThreshold = Duration.ofMillis(expireThresholdMillis);
+ final long expiredCreationTime =
+ expireThreshold.negated().plusMillis(Time.now()).toMillis();
+
+ Set<String> expiredLegacyKeys = new HashSet<>();
+ Set<String> expiredSplitKeys = new HashSet<>();
+
+ // Create 2 legacy-schema expired MPUs with embedded parts
+ for (int i = 0; i < 2; i++) {
+ String uploadId = OMMultipartUploadUtils.getMultipartUploadId();
+ OmMultipartKeyInfo mpuKeyInfo = OMRequestTestUtils
+ .createOmMultipartKeyInfo(uploadId, expiredCreationTime,
+ HddsProtos.ReplicationType.RATIS,
+ HddsProtos.ReplicationFactor.ONE, 0L);
+ String keyName = "legacy" + i;
+ OmKeyInfo keyInfo = OMRequestTestUtils.createOmKeyInfo(volumeName,
+ bucketName, keyName, RatisReplicationConfig.getInstance(ONE))
+ .setCreationTime(expiredCreationTime)
+ .build();
+ for (int j = 1; j <= numPartsPerMPU; j++) {
+ PartKeyInfo partKeyInfo = OMRequestTestUtils
+ .createPartKeyInfo(volumeName, bucketName, keyName, uploadId, j);
+ OMRequestTestUtils.addPart(partKeyInfo, mpuKeyInfo);
+ }
+ String mpuDbKey = OMRequestTestUtils.addMultipartInfoToTable(
+ false, keyInfo, mpuKeyInfo, 0L, omMetadataManager);
+ expiredLegacyKeys.add(mpuDbKey);
+ }
+
+ // Create 2 split-schema expired MPUs with parts in multipartPartsTable
+ for (int i = 0; i < 2; i++) {
+ String uploadId = OMMultipartUploadUtils.getMultipartUploadId();
+ OmMultipartKeyInfo mpuKeyInfo = new OmMultipartKeyInfo.Builder(
+ OMRequestTestUtils.createOmMultipartKeyInfo(uploadId,
expiredCreationTime,
+ HddsProtos.ReplicationType.RATIS,
+ HddsProtos.ReplicationFactor.ONE, 0L))
+
.setSchemaVersion(OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION)
+ .build();
+ String keyName = "split" + i;
+ OmKeyInfo keyInfo = OMRequestTestUtils.createOmKeyInfo(volumeName,
+ bucketName, keyName, RatisReplicationConfig.getInstance(ONE))
+ .setCreationTime(expiredCreationTime)
+ .build();
+ for (int j = 1; j <= numPartsPerMPU; j++) {
+ addSplitSchemaPart(uploadId, j);
+ }
+ String mpuDbKey = OMRequestTestUtils.addMultipartInfoToTable(
+ false, keyInfo, mpuKeyInfo, 0L, omMetadataManager);
+ expiredSplitKeys.add(mpuDbKey);
+ }
+
+ // Budget of 9 parts fits exactly 3 MPUs (3 parts each)
+ List<ExpiredMultipartUploadsBucket> someExpiredMPUs =
+ omMetadataManager.getExpiredMultipartUploads(expireThreshold,
+ numPartsPerMPU * 3);
+ List<String> names = getMultipartKeyNames(someExpiredMPUs);
+ assertEquals(3, names.size());
+
+ // Budget of 12 parts fits all 4 expired MPUs
+ List<ExpiredMultipartUploadsBucket> allExpiredMPUs =
+ omMetadataManager.getExpiredMultipartUploads(expireThreshold,
+ numPartsPerMPU * 4);
+ names = getMultipartKeyNames(allExpiredMPUs);
+ assertEquals(4, names.size());
+ assertThat(names).containsAll(expiredLegacyKeys);
+ assertThat(names).containsAll(expiredSplitKeys);
+ }
+
+ @Test
+ public void testGetExpiredMPUsZeroPartsSplitSchema() throws Exception {
+ final String bucketName = UUID.randomUUID().toString();
+ final String volumeName = UUID.randomUUID().toString();
+ final long expireThresholdMillis = ozoneConfiguration.getTimeDuration(
+ OZONE_OM_MPU_EXPIRE_THRESHOLD,
+ OZONE_OM_MPU_EXPIRE_THRESHOLD_DEFAULT,
+ TimeUnit.MILLISECONDS);
+ final Duration expireThreshold = Duration.ofMillis(expireThresholdMillis);
+ final long expiredCreationTime =
+ expireThreshold.negated().plusMillis(Time.now()).toMillis();
+
+ // Zero-parts split-schema MPU (freshly initiated, no parts uploaded)
+ String uploadIdZero = OMMultipartUploadUtils.getMultipartUploadId();
+ OmMultipartKeyInfo mpuZero = new OmMultipartKeyInfo.Builder(
+ OMRequestTestUtils.createOmMultipartKeyInfo(uploadIdZero,
expiredCreationTime,
+ HddsProtos.ReplicationType.RATIS,
+ HddsProtos.ReplicationFactor.ONE, 0L))
+ .setSchemaVersion(OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION)
+ .build();
+ OmKeyInfo keyInfoZero = OMRequestTestUtils.createOmKeyInfo(volumeName,
+ bucketName, "aaa-zero-parts",
RatisReplicationConfig.getInstance(ONE))
+ .setCreationTime(expiredCreationTime)
+ .build();
+ String zeroKey = OMRequestTestUtils.addMultipartInfoToTable(
+ false, keyInfoZero, mpuZero, 0L, omMetadataManager);
+
+ // Split-schema MPU with 3 parts
+ String uploadIdThree = OMMultipartUploadUtils.getMultipartUploadId();
+ OmMultipartKeyInfo mpuThree = new OmMultipartKeyInfo.Builder(
+ OMRequestTestUtils.createOmMultipartKeyInfo(uploadIdThree,
expiredCreationTime,
+ HddsProtos.ReplicationType.RATIS,
+ HddsProtos.ReplicationFactor.ONE, 0L))
+ .setSchemaVersion(OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION)
+ .build();
+ OmKeyInfo keyInfoThree = OMRequestTestUtils.createOmKeyInfo(volumeName,
+ bucketName, "zzz-three-parts",
RatisReplicationConfig.getInstance(ONE))
+ .setCreationTime(expiredCreationTime)
+ .build();
+ for (int j = 1; j <= 3; j++) {
+ addSplitSchemaPart(uploadIdThree, j);
+ }
+ String threeKey = OMRequestTestUtils.addMultipartInfoToTable(
+ false, keyInfoThree, mpuThree, 0L, omMetadataManager);
+
+ // maxParts=0 returns nothing (loop pre-check 0 < 0 is false)
+ List<ExpiredMultipartUploadsBucket> noneExpired =
+ omMetadataManager.getExpiredMultipartUploads(expireThreshold, 0);
+ assertTrue(getMultipartKeyNames(noneExpired).isEmpty());
+
+ // maxParts=3 should return both: zero-parts contributes 0, three-parts
+ // contributes 3, total = 3 which satisfies the loop exit condition
+ List<ExpiredMultipartUploadsBucket> allExpired =
+ omMetadataManager.getExpiredMultipartUploads(expireThreshold, 3);
+ List<String> names = getMultipartKeyNames(allExpired);
+ assertEquals(2, names.size());
+ assertThat(names).contains(zeroKey);
+ assertThat(names).contains(threeKey);
+ }
+
+ private void addSplitSchemaPart(String uploadId, int partNumber) throws
IOException {
+ OmKeyLocationInfo locationInfo = new OmKeyLocationInfo.Builder()
+ .setBlockID(new BlockID(1L, partNumber))
+ .setLength(100)
+ .build();
+ OmKeyLocationInfoGroup locationGroup = new OmKeyLocationInfoGroup(0,
+ Collections.singletonList(locationInfo));
+ String partName = "part-" + partNumber;
+ OmKeyInfo keyInfo = OMRequestTestUtils.createOmKeyInfo("vol", "bucket",
"key",
+ RatisReplicationConfig.getInstance(ONE))
+ .setDataSize(100L)
+ .setObjectID(partNumber)
+ .setUpdateID(partNumber)
+ .addOmKeyLocationInfoGroup(locationGroup)
+ .addMetadata(org.apache.hadoop.ozone.OzoneConsts.ETAG, "etag-" +
partNumber)
+ .build();
+ OmMultipartPartInfo partInfo = OmMultipartPartInfo.from(partName,
partNumber, keyInfo);
+ omMetadataManager.getMultipartPartsTable().put(
+ OmMultipartPartKey.of(uploadId, partNumber), partInfo);
+ }
+
private List<String> getOpenKeyNames(
Collection<OpenKeyBucket.Builder> openKeyBuckets) {
return openKeyBuckets.stream()
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestKeyLifecycleService.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestKeyLifecycleService.java
index 6b8f2e6e12e..89746fb6c82 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestKeyLifecycleService.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestKeyLifecycleService.java
@@ -23,6 +23,7 @@
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.OzoneAcl.AclScope.ACCESS;
+import static org.apache.hadoop.ozone.OzoneConsts.ETAG;
import static org.apache.hadoop.ozone.OzoneConsts.OM_KEY_PREFIX;
import static
org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_KEY_LIFECYCLE_SERVICE_DELETE_BATCH_SIZE;
import static
org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_KEY_LIFECYCLE_SERVICE_ENABLED;
@@ -31,6 +32,7 @@
import static
org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_KEY_LIFECYCLE_SERVICE_STATE_SAVE_INTERVAL_MS_DEFAULT;
import static
org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_KEY_LIFECYCLE_SERVICE_STATE_SAVE_KEYS_PROCESSED;
import static
org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_KEY_LIFECYCLE_SERVICE_STATE_SAVE_KEYS_PROCESSED_DEFAULT;
+import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_DB_DIRS;
import static org.apache.hadoop.ozone.om.OmConfig.Keys.ENABLE_FILESYSTEM_PATHS;
import static
org.apache.hadoop.ozone.om.exceptions.OMException.ResultCodes.INVALID_REQUEST;
import static
org.apache.hadoop.ozone.om.helpers.BucketLayout.FILE_SYSTEM_OPTIMIZED;
@@ -50,10 +52,12 @@
import static org.junit.jupiter.params.provider.Arguments.arguments;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.argThat;
+import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.atLeast;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
import com.google.common.collect.ImmutableMap;
import java.io.File;
@@ -77,6 +81,7 @@
import org.apache.commons.lang3.RandomStringUtils;
import org.apache.commons.lang3.tuple.Pair;
import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.hdds.client.BlockID;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList;
@@ -115,12 +120,16 @@
import org.apache.hadoop.ozone.om.helpers.OmLifecycleScanState;
import org.apache.hadoop.ozone.om.helpers.OmMultipartInfo;
import org.apache.hadoop.ozone.om.helpers.OmMultipartKeyInfo;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartPartInfo;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartPartKey;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartUpload;
import org.apache.hadoop.ozone.om.helpers.OmVolumeArgs;
import org.apache.hadoop.ozone.om.helpers.OpenKeySession;
import org.apache.hadoop.ozone.om.helpers.RepeatedOmKeyInfo;
import org.apache.hadoop.ozone.om.protocol.OzoneManagerProtocol;
import org.apache.hadoop.ozone.om.request.OMRequestTestUtils;
import org.apache.hadoop.ozone.om.request.key.OMKeysDeleteRequest;
+import org.apache.hadoop.ozone.om.request.util.OMMultipartUploadUtils;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.LifecycleConfiguration;
import org.apache.hadoop.ozone.security.acl.IAccessAuthorizer;
@@ -146,6 +155,7 @@
import org.junit.jupiter.params.provider.EnumSource;
import org.junit.jupiter.params.provider.MethodSource;
import org.junit.jupiter.params.provider.ValueSource;
+import org.mockito.Mockito;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.event.Level;
@@ -3340,4 +3350,120 @@ private OmMultipartInfo
createTestMultipartUpload(String volumeName, String buck
public static String uniqueObjectName(String prefix) {
return prefix + OBJECT_COUNTER.getAndIncrement();
}
+
+ @Test
+ public void testPartCountIoExceptionSkipsUploadContinuesBucket(@TempDir File
tempDir)
+ throws Exception {
+ OzoneConfiguration omConf = new OzoneConfiguration();
+ omConf.set(OZONE_OM_DB_DIRS, tempDir.getAbsolutePath());
+ OMMetadataManager realMetadataManager = new OmMetadataManagerImpl(omConf,
null);
+ try {
+ String uploadA = OMMultipartUploadUtils.getMultipartUploadId();
+ String uploadB = OMMultipartUploadUtils.getMultipartUploadId();
+ String uploadC = OMMultipartUploadUtils.getMultipartUploadId();
+
+ addSplitSchemaPart(realMetadataManager, uploadA, 1);
+ addSplitSchemaPart(realMetadataManager, uploadA, 2);
+ addSplitSchemaPart(realMetadataManager, uploadC, 1);
+
+ // Spy only the lightweight parts table and inject an IOException for
+ // uploadB's prefix scan to simulate a corrupt read; uploadA and uploadC
+ // fall through to the real table with real data. doThrow(...).when(...)
+ // is used (not when(...).thenThrow(...)) so the real iterator() is not
+ // invoked during stubbing, which would leak a native RocksDB iterator
+ // and block getStore().close().
+ Table<OmMultipartPartKey, OmMultipartPartInfo> spyPartsTable =
+ Mockito.spy(realMetadataManager.getMultipartPartsTable());
+ Mockito.doThrow(new RocksDatabaseException("simulated corruption"))
+
.when(spyPartsTable).iterator(eq(OmMultipartPartKey.prefix(uploadB)));
+
+ OMMetadataManager mockMM = Mockito.mock(OMMetadataManager.class);
+ when(mockMM.getMultipartPartsTable()).thenReturn(spyPartsTable);
+
+ assertEquals(2, OMMultipartUploadUtils.countParts(mockMM, uploadA));
+ assertEquals(1, OMMultipartUploadUtils.countParts(mockMM, uploadC));
+ assertThrows(IOException.class,
+ () -> OMMultipartUploadUtils.countParts(mockMM, uploadB));
+
+ KeyLifecycleService.PartCountLimitedList list =
+ new KeyLifecycleService.PartCountLimitedList(10);
+ for (String uploadId : Arrays.asList(uploadA, uploadB, uploadC)) {
+ try {
+ int partCount = OMMultipartUploadUtils.countParts(mockMM, uploadId);
+ list.add(new OmMultipartUpload("v", "b", "k", uploadId), partCount);
+ } catch (IOException e) {
+ // per-MPU skip — bucket loop continues
+ }
+ }
+ // A (2 parts) and C (1 part) are added; B is skipped due to IOException
+ assertEquals(2, list.size()); // uploadA and uploadC only
+ assertEquals(3, list.getPartCount()); // 2 (uploadA) + 1 (uploadC)
+ } finally {
+ realMetadataManager.getStore().close();
+ }
+ }
+
+ @Test
+ public void testPartCountLimitedListBoundaryBehavior() {
+ // Zero-parts upload is addable and does not fill the list
+ KeyLifecycleService.PartCountLimitedList list =
+ new KeyLifecycleService.PartCountLimitedList(5);
+ list.add(new OmMultipartUpload("v", "b", "k", "id1"), 0);
+ assertEquals(1, list.size());
+ assertEquals(0, list.getPartCount());
+ assertFalse(list.isFull());
+ assertFalse(list.isEmpty());
+
+ // Exact boundary: partCount == maxPartCount triggers isFull
+ KeyLifecycleService.PartCountLimitedList exactList =
+ new KeyLifecycleService.PartCountLimitedList(5);
+ exactList.add(new OmMultipartUpload("v", "b", "k", "id2"), 5);
+ assertEquals(1, exactList.size());
+ assertEquals(5, exactList.getPartCount());
+ assertTrue(exactList.isFull());
+
+ // Over boundary: cumulative parts exceed max
+ KeyLifecycleService.PartCountLimitedList overList =
+ new KeyLifecycleService.PartCountLimitedList(5);
+ overList.add(new OmMultipartUpload("v", "b", "k", "id3"), 3);
+ assertFalse(overList.isFull());
+ overList.add(new OmMultipartUpload("v", "b", "k", "id4"), 3);
+ assertTrue(overList.isFull());
+ assertEquals(2, overList.size());
+ assertEquals(6, overList.getPartCount());
+
+ // clear() resets all state
+ overList.clear();
+ assertTrue(overList.isEmpty());
+ assertFalse(overList.isFull());
+ assertEquals(0, overList.size());
+ assertEquals(0, overList.getPartCount());
+ }
+
+ private static void addSplitSchemaPart(OMMetadataManager omMetadataManager,
+ String uploadId, int partNumber) throws IOException {
+ OmKeyLocationInfo locationInfo = new OmKeyLocationInfo.Builder()
+ .setBlockID(new BlockID(1L, partNumber))
+ .setLength(100)
+ .build();
+ OmKeyLocationInfoGroup locationGroup = new OmKeyLocationInfoGroup(0,
+ Collections.singletonList(locationInfo));
+ String partName = "part-" + partNumber;
+ OmKeyInfo keyInfo = new OmKeyInfo.Builder()
+ .setVolumeName("v")
+ .setBucketName("b")
+ .setKeyName("k")
+ .setReplicationConfig(RatisReplicationConfig.getInstance(THREE))
+ .setDataSize(100L)
+ .setCreationTime(System.currentTimeMillis())
+ .setModificationTime(System.currentTimeMillis())
+ .setObjectID(partNumber)
+ .setUpdateID(partNumber)
+ .addOmKeyLocationInfoGroup(locationGroup)
+ .addMetadata(ETAG, "etag-" + partNumber)
+ .build();
+ OmMultipartPartInfo partInfo = OmMultipartPartInfo.from(partName,
partNumber, keyInfo);
+ omMetadataManager.getMultipartPartsTable().put(
+ OmMultipartPartKey.of(uploadId, partNumber), partInfo);
+ }
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]