This is an automated email from the ASF dual-hosted git repository.
ChenSammi 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 711476cb0c4 HDDS-16089. Improve MPU table evaluation efficiency in
KeyLifeCycleService (#10949)
711476cb0c4 is described below
commit 711476cb0c4973a6f045a875d97de66916a6bc5b
Author: Priyesh Karatha <[email protected]>
AuthorDate: Fri Aug 14 13:45:42 2026 +0530
HDDS-16089. Improve MPU table evaluation efficiency in KeyLifeCycleService
(#10949)
---
.../ozone/om/service/KeyLifecycleService.java | 133 +++++++++++----------
.../ozone/om/service/TestKeyLifecycleService.java | 119 ++++++++++++++++++
2 files changed, 190 insertions(+), 62 deletions(-)
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 e551b9242e4..c36097fd4d3 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
@@ -1108,55 +1108,81 @@ private void processMultipartUploads(OmBucketInfo
bucketInfo, List<OmLCRule> rul
upload.setCreationTime(Instant.ofEpochMilli(mpuKeyInfo.getCreationTime()));
String keyName = upload.getKeyName();
- String multipartOpenKey;
- try {
- multipartOpenKey = OMMultipartUploadUtils.getMultipartOpenKey(
- volumeName, bucketName, keyName, upload.getUploadId(),
- omMetadataManager, bucketInfo.getBucketLayout());
- } catch (OMException e) {
- LOG.warn("Failed to get multipart open key for {}/{}/{}, skipping",
- volumeName, bucketName, keyName, e);
- continue;
- }
-
- OmKeyInfo openKeyInfo =
omMetadataManager.getOpenKeyTable(bucketInfo.getBucketLayout())
- .get(multipartOpenKey);
- if (openKeyInfo == null) {
- LOG.warn("Open key not found for multipart upload {}/{}/{},
skipping",
- volumeName, bucketName, keyName);
- continue;
- }
-
+ OmLCRule matchingRule = null;
+ OmKeyInfo openKeyInfo = null;
+ boolean openKeyFetchAttempted = false;
+ boolean skipUpload = false;
for (OmLCRule rule : ruleList) {
- if (shouldAbortUpload(openKeyInfo, upload, keyName, rule)) {
- if (expiredUploads.isFull()) {
- LOG.info("Multipart upload batch reached part count limit {},
aborting current batch " +
- "({} uploads, {} parts) for bucket {}/{}",
- mpuAbortLimitPerTask, expiredUploads.size(),
expiredUploads.getPartCount(),
- volumeName, bucketName);
- abortExpiredMultipartUploadsAndClear(bucketInfo,
expiredUploads);
- }
-
- // 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;
+ if (!passesAgeAndPrefix(upload, keyName, rule)) {
+ continue;
+ }
+ if (!rule.isTagEnable()) {
+ matchingRule = rule;
+ break;
+ }
+ if (!openKeyFetchAttempted) {
+ openKeyFetchAttempted = true;
try {
- uploadedParts = mpuKeyInfo.getSchemaVersion()
- == OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION
- ? OMMultipartUploadUtils.countParts(omMetadataManager,
upload.getUploadId())
- : mpuKeyInfo.getPartKeyInfoMap().size();
+ String multipartOpenKey =
OMMultipartUploadUtils.getMultipartOpenKey(
+ volumeName, bucketName, keyName, upload.getUploadId(),
+ omMetadataManager, bucketInfo.getBucketLayout());
+ openKeyInfo = omMetadataManager.getOpenKeyTable(
+ bucketInfo.getBucketLayout()).get(multipartOpenKey);
+ } catch (OMException e) {
+ LOG.warn("Failed to get multipart open key for {}/{}/{},
skipping",
+ volumeName, bucketName, keyName, e);
+ skipUpload = true;
+ break;
} catch (IOException e) {
- LOG.warn("Failed to count parts for MPU {}/{}/{} uploadId {},
skipping",
- volumeName, bucketName, keyName, upload.getUploadId(), e);
+ LOG.warn("Failed to read open key table for {}/{}/{},
skipping",
+ volumeName, bucketName, keyName, e);
+ skipUpload = true;
break;
}
- expiredUploads.add(upload, uploadedParts);
- LOG.debug("Multipart upload {}/{}/{} with uploadId {} ({} parts)
will be aborted",
- volumeName, bucketName, keyName, upload.getUploadId(),
uploadedParts);
+ if (openKeyInfo == null) {
+ LOG.debug("Orphan multipart upload {}/{}/{} has no open key
entry, skipping tag-requiring rules",
+ volumeName, bucketName, keyName);
+ }
+ }
+ if (openKeyInfo == null) {
+ continue;
+ }
+ OmLCFilter filter = rule.getFilter();
+ if (filter == null || filter.match(openKeyInfo, keyName)) {
+ matchingRule = rule;
break;
}
}
+
+ if (skipUpload || matchingRule == null) {
+ continue;
+ }
+
+ if (expiredUploads.isFull()) {
+ LOG.info("Multipart upload batch reached part count limit {},
aborting current batch " +
+ "({} uploads, {} parts) for bucket {}/{}",
+ mpuAbortLimitPerTask, expiredUploads.size(),
expiredUploads.getPartCount(),
+ volumeName, bucketName);
+ abortExpiredMultipartUploadsAndClear(bucketInfo, expiredUploads);
+ }
+
+ // 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);
+ continue;
+ }
+ expiredUploads.add(upload, uploadedParts);
+ LOG.debug("Multipart upload {}/{}/{} with uploadId {} ({} parts)
will be aborted",
+ volumeName, bucketName, keyName, upload.getUploadId(),
uploadedParts);
}
} catch (IOException e) {
LOG.warn("Failed to iterate multipartInfoTable for bucket {}/{}",
volumeName, bucketName, e);
@@ -1171,33 +1197,16 @@ private void processMultipartUploads(OmBucketInfo
bucketInfo, List<OmLCRule> rul
}
/**
- * Check if a multipart upload should be aborted based on the lifecycle
rule.
- *
- * @param openKeyInfo the open key information with tags
- * @param upload the multipart upload information
- * @param keyName the key name of the upload
- * @param rule the lifecycle rule to evaluate against
- * @return true if the upload should be aborted, false otherwise
+ * Returns true if the upload passes the age and prefix checks for the
given rule,
+ * without consulting the open key table (no tag evaluation).
*/
- private boolean shouldAbortUpload(OmKeyInfo openKeyInfo, OmMultipartUpload
upload,
- String keyName, OmLCRule rule) {
-
+ private boolean passesAgeAndPrefix(OmMultipartUpload upload, String
keyName, OmLCRule rule) {
if (!rule.getAbortIncompleteMultipartUpload().shouldAbort(
upload.getCreationTime().toEpochMilli())) {
return false;
}
-
String effectivePrefix = rule.getEffectivePrefix();
- if (effectivePrefix != null && !keyName.startsWith(effectivePrefix)) {
- return false;
- }
-
- OmLCFilter filter = rule.getFilter();
- if (filter != null && !filter.match(openKeyInfo, keyName)) {
- return false;
- }
-
- return true;
+ return effectivePrefix == null || keyName.startsWith(effectivePrefix);
}
/**
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 c80ae4e5881..012123680f6 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
@@ -2840,6 +2840,125 @@ void testAbortIncompleteMultipartUploadWithTagFilter()
throws Exception {
deleteLifecyclePolicy(volumeName, bucketName);
}
+ /**
+ * An MPU whose open key entry is missing (orphan) should still be aborted
when a lifecycle
+ * rule matches by age and prefix alone (no tag filter). The abort request
handler already
+ * tolerates a missing open key.
+ */
+ @Test
+ void testOrphanMpuAbortedByAgeAndPrefixRule() throws Exception {
+ final String volumeName = getTestName();
+ final String bucketName = uniqueObjectName("bucket");
+
+ createVolumeAndBucket(volumeName, bucketName, OBJECT_STORE,
+ UserGroupInformation.getCurrentUser().getShortUserName());
+
+ String owner = UserGroupInformation.getCurrentUser().getShortUserName();
+ long initialMpuCount = getMultipartUploadCount(volumeName, bucketName);
+
+ // Create two MPUs: one normal (has open key), one will become an orphan.
+ OmMultipartInfo normalMpu = createTestMultipartUpload(volumeName,
bucketName, "data/normal", owner);
+ OmMultipartInfo orphanMpu = createTestMultipartUpload(volumeName,
bucketName, "data/orphan", owner);
+
+ // Age both MPUs past the threshold.
+ long oldCreationTime = System.currentTimeMillis() -
TimeUnit.DAYS.toMillis(2);
+ updateMultipartUploadCreationTime(volumeName, bucketName, "data/normal",
+ normalMpu.getUploadID(), oldCreationTime);
+ updateMultipartUploadCreationTime(volumeName, bucketName, "data/orphan",
+ orphanMpu.getUploadID(), oldCreationTime);
+
+ // Simulate orphan: delete the open key entry for "data/orphan" directly
from the table.
+ // Use the same helper the service uses so the key format matches
exactly.
+ String resolvedOrphanOpenKey = OMMultipartUploadUtils
+ .getMultipartOpenKey(volumeName, bucketName, "data/orphan",
orphanMpu.getUploadID(),
+ metadataManager, OBJECT_STORE);
+
metadataManager.getOpenKeyTable(OBJECT_STORE).delete(resolvedOrphanOpenKey);
+
+ // Rule: abort all MPUs under "data/" after 1 day — no tag filter.
+ OmLCRule rule = new OmLCRule.Builder()
+ .setId("abort-data-prefix")
+ .setEnabled(true)
+ .setFilter(new OmLCFilter.Builder().setPrefix("data/").build())
+ .setAction(new OmLCAbortIncompleteMultipartUpload.Builder()
+ .setDaysAfterInitiation(1)
+ .build())
+ .build();
+
+ createLifecyclePolicy(volumeName, bucketName, OBJECT_STORE,
Collections.singletonList(rule));
+
+ // Both MPUs should be aborted: normal one (open key present) and orphan
(open key missing).
+ GenericTestUtils.waitFor(() ->
+ getMultipartUploadCount(volumeName, bucketName) - initialMpuCount ==
0,
+ WAIT_CHECK_INTERVAL, 10000);
+
+ String normalKey = metadataManager.getMultipartKey(volumeName,
bucketName,
+ "data/normal", normalMpu.getUploadID());
+ assertNull(metadataManager.getMultipartInfoTable().get(normalKey),
+ "Normal MPU should be aborted");
+
+ String orphanKey = metadataManager.getMultipartKey(volumeName,
bucketName,
+ "data/orphan", orphanMpu.getUploadID());
+ assertNull(metadataManager.getMultipartInfoTable().get(orphanKey),
+ "Orphan MPU (no open key) should be aborted when a tag-free rule
matches");
+
+ deleteLifecyclePolicy(volumeName, bucketName);
+ }
+
+ /**
+ * An MPU whose open key is missing (orphan) must NOT be aborted when the
only matching
+ * lifecycle rule requires a tag filter, because the tag metadata is
unavailable.
+ */
+ @Test
+ void testOrphanMpuNotAbortedByTagOnlyRule() throws Exception {
+ final String volumeName = getTestName();
+ final String bucketName = uniqueObjectName("bucket");
+
+ createVolumeAndBucket(volumeName, bucketName, OBJECT_STORE,
+ UserGroupInformation.getCurrentUser().getShortUserName());
+
+ String owner = UserGroupInformation.getCurrentUser().getShortUserName();
+ long initialMpuCount = getMultipartUploadCount(volumeName, bucketName);
+
+ OmMultipartInfo orphanMpu = createTestMultipartUpload(volumeName,
bucketName, "file.txt", owner);
+
+ // Age past the threshold.
+ long oldCreationTime = System.currentTimeMillis() -
TimeUnit.DAYS.toMillis(2);
+ updateMultipartUploadCreationTime(volumeName, bucketName, "file.txt",
+ orphanMpu.getUploadID(), oldCreationTime);
+
+ // Simulate orphan by removing its open key.
+ String resolvedOrphanOpenKey = OMMultipartUploadUtils
+ .getMultipartOpenKey(volumeName, bucketName, "file.txt",
orphanMpu.getUploadID(),
+ metadataManager, OBJECT_STORE);
+
metadataManager.getOpenKeyTable(OBJECT_STORE).delete(resolvedOrphanOpenKey);
+
+ // Rule requires tag match — cannot evaluate without open key.
+ OmLCRule tagRule = new OmLCRule.Builder()
+ .setId("abort-by-tag")
+ .setEnabled(true)
+ .setFilter(new OmLCFilter.Builder().setTag("env", "test").build())
+ .setAction(new OmLCAbortIncompleteMultipartUpload.Builder()
+ .setDaysAfterInitiation(1)
+ .build())
+ .build();
+
+ createLifecyclePolicy(volumeName, bucketName, OBJECT_STORE,
Collections.singletonList(tagRule));
+
+ // Wait long enough for the service to run at least once.
+ Thread.sleep(SERVICE_INTERVAL * 2);
+
+ // Orphan MPU must remain: the tag rule cannot evaluate without open key
metadata.
+ String orphanKey = metadataManager.getMultipartKey(volumeName,
bucketName,
+ "file.txt", orphanMpu.getUploadID());
+ assertNotNull(metadataManager.getMultipartInfoTable().get(orphanKey),
+ "Orphan MPU should NOT be aborted when only tag-requiring rules
exist");
+
+ assertEquals(initialMpuCount + 1, getMultipartUploadCount(volumeName,
bucketName),
+ "MPU count should not change");
+
+ deleteLifecyclePolicy(volumeName, bucketName);
+ }
+
}
/**
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]