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 589bd98aaceb..d075c87ecd9f 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,10 @@ 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.OmVolumeArgs; import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.PartKeyInfo; import picocli.CommandLine; @@ -100,6 +103,7 @@ public class ContainerToKeyMapping extends AbstractSubcommand implements Callabl private Table openFileTable; private Table openKeyTable; private Table multipartInfoTable; + private Table multipartPartsTable; private DBStore dirTreeDbStore; private Table dirTreeTable; // Cache volume IDs to avoid repeated lookups @@ -138,6 +142,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 +318,18 @@ private void processMultipartUpload(Set containerIds, Map 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 +344,43 @@ private void processMultipartUpload(Set containerIds, Map getKeyContainers(OmKeyInfo keyInfo, Set 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 getSplitPartContainers(String multipartInfoDbKey, Set targetContainerIds) { + Set matchedContainers = new HashSet<>(); + String uploadId = multipartInfoDbKey.substring( + multipartInfoDbKey.lastIndexOf(OM_KEY_PREFIX) + OM_KEY_PREFIX.length()); + OmMultipartPartKey prefix = OmMultipartPartKey.prefix(uploadId); + try (TableIterator> + partIterator = multipartPartsTable.iterator(prefix)) { + while (partIterator.hasNext()) { + Table.KeyValue 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 getContainers(List locationVersions, Set targetContainerIds) { Set 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 4cad62cd719c..7d1d651d2488 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; @@ -189,6 +192,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_4), "--in-progress"); + assertEquals(0, exitCode); + + String output = outWriter.toString(); + + assertThat(output).contains("\"" + CONTAINER_ID_4 + "\""); + assertThat(output).contains("\"openKeys\""); + assertThat(output).contains("/vol1/obs-bucket/splitMpuKey/split-upload-id"); + } + @Test public void testNonExistentContainer() { @@ -292,6 +307,7 @@ private void createTestData() throws Exception { // Create MPU (multipart upload) for OBS bucket with parts in container 5 createMultipartUpload(); + createSplitSchemaMultipartUpload(); } /** @@ -339,6 +355,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_4)) + .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_4)) + .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 +442,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 f31815bc208c..99b1c18d8ec2 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 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 598f7dc94102..d6fd32af51df 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 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 getPartKeys(String uploadId, SortedMap parts) { List 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 ead8acc7cd51..5310f5b4837e 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 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 754fbc1ed21e..1ffcc32f0acb 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 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 someExpiredMPUs = + omMetadataManager.getExpiredMultipartUploads( + expireThreshold, + (numExpiredMPUs * numPartsPerMPU) - (numPartsPerMPU)); + List names = getMultipartKeyNames(someExpiredMPUs); + assertEquals(numExpiredMPUs - 1, names.size()); + assertThat(expiredMPUs).containsAll(names); + + List 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 expiredLegacyKeys = new HashSet<>(); + Set 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 someExpiredMPUs = + omMetadataManager.getExpiredMultipartUploads(expireThreshold, + numPartsPerMPU * 3); + List names = getMultipartKeyNames(someExpiredMPUs); + assertEquals(3, names.size()); + + // Budget of 12 parts fits all 4 expired MPUs + List 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 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 allExpired = + omMetadataManager.getExpiredMultipartUploads(expireThreshold, 3); + List 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 getOpenKeyNames( Collection 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 6b8f2e6e12ed..09f274daba45 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,111 @@ 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 omMetadataManager = new OmMetadataManagerImpl(omConf, null); + try { + String uploadA = OMMultipartUploadUtils.getMultipartUploadId(); + String uploadB = OMMultipartUploadUtils.getMultipartUploadId(); + String uploadC = OMMultipartUploadUtils.getMultipartUploadId(); + + addSplitSchemaPart(omMetadataManager, uploadA, 1); + addSplitSchemaPart(omMetadataManager, uploadA, 2); + addSplitSchemaPart(omMetadataManager, uploadC, 1); + + OMMetadataManager mockMM = Mockito.spy(omMetadataManager); + Table mockPartsTable = + Mockito.spy(omMetadataManager.getMultipartPartsTable()); + when(mockMM.getMultipartPartsTable()).thenReturn(mockPartsTable); + when(mockPartsTable.iterator(eq(OmMultipartPartKey.prefix(uploadB)))) + .thenThrow(new IOException("simulated corruption")); + + assertEquals(2, OMMultipartUploadUtils.countParts(mockMM, uploadA)); + 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 + } + } + assertEquals(2, list.size()); + assertEquals(3, list.getPartCount()); + } finally { + omMetadataManager.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); + } }