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]


Reply via email to