laserninja commented on code in PR #13531: URL: https://github.com/apache/gravitino/pull/13531#discussion_r4151815927
########## api/src/main/java/org/apache/gravitino/policy/IcebergOrphanFileRemovalContent.java: ########## @@ -0,0 +1,163 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.gravitino.policy; + +import com.google.common.base.Preconditions; +import com.google.common.collect.ImmutableMap; +import com.google.common.collect.ImmutableSet; +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import javax.annotation.Nullable; +import org.apache.commons.lang3.StringUtils; +import org.apache.gravitino.MetadataObject; + +/** Configuration for table-level Iceberg orphan cleanup on an optimizer invocation. */ +public class IcebergOrphanFileRemovalContent implements PolicyContent { + /** Property key for strategy type. */ + public static final String STRATEGY_TYPE_KEY = "strategy.type"; + /** Property key for the job template. */ + public static final String JOB_TEMPLATE_NAME_KEY = "job.template-name"; + /** Rule key for the trigger expression. */ + public static final String TRIGGER_EXPR_KEY = "trigger-expr"; + /** Rule key for the score expression. */ + public static final String SCORE_EXPR_KEY = "score-expr"; + /** Prefix for policy rules forwarded as job options. */ + public static final String JOB_OPTIONS_PREFIX = "job.options."; + /** Retention option key shared with the job adapter. */ + public static final String OLDER_THAN_DAYS_KEY = "olderThanDays"; + /** Scan location option key shared with the job adapter. */ + public static final String LOCATION_KEY = "location"; + /** Preview mode option key shared with the job adapter. */ + public static final String DRY_RUN_KEY = "dryRun"; + /** Maximum retention, approximately 100 years, to keep cleanup timestamps practical. */ + public static final long MAX_OLDER_THAN_DAYS = 36500; + /** Strategy type for orphan cleanup. */ + public static final String STRATEGY_TYPE_VALUE = "iceberg-orphan-file-removal"; + /** Existing Spark job template used for cleanup. */ + public static final String JOB_TEMPLATE_NAME_VALUE = "builtin-iceberg-remove-orphan-files"; + /** Default retention in days. */ + public static final long DEFAULT_OLDER_THAN_DAYS = 3; + /** Default execution mode. */ + public static final boolean DEFAULT_DRY_RUN = false; + + private final long olderThanDays; + @Nullable private final String location; + private final boolean dryRun; + + private IcebergOrphanFileRemovalContent() { + this(DEFAULT_OLDER_THAN_DAYS, null, DEFAULT_DRY_RUN); + } + + IcebergOrphanFileRemovalContent(long olderThanDays, @Nullable String location, boolean dryRun) { + this.olderThanDays = olderThanDays; + this.location = location; + this.dryRun = dryRun; + } + + /** + * @return minimum age of files eligible for removal, in days + */ + public long olderThanDays() { + return olderThanDays; + } + + /** + * @return optional scan location within the table, or null for the table location + */ + @Nullable + public String location() { + return location; + } + + /** + * @return whether candidates should only be listed + */ + public boolean dryRun() { + return dryRun; + } + + @Override + public Set<MetadataObject.Type> supportedObjectTypes() { + return ImmutableSet.of( + MetadataObject.Type.CATALOG, MetadataObject.Type.SCHEMA, MetadataObject.Type.TABLE); + } + + @Override + public Map<String, String> properties() { + return ImmutableMap.of( + STRATEGY_TYPE_KEY, STRATEGY_TYPE_VALUE, JOB_TEMPLATE_NAME_KEY, JOB_TEMPLATE_NAME_VALUE); + } + + @Override + public Map<String, Object> rules() { + Map<String, Object> rules = new LinkedHashMap<>(); + // Eligibility is evaluated when the optimizer is invoked. Scheduling is owned by the caller. + rules.put(TRIGGER_EXPR_KEY, "true"); Review Comment: Addressed in efb339668 and the PR description. Both now explicitly state that every enabled, selected submission launches a Spark job, with no metric gate, cooldown, or last-run check. Operators must rate-limit through external scheduling and start with dryRun: true; dry-run jobs also consume Spark resources. ########## maintenance/optimizer-api/src/main/java/org/apache/gravitino/maintenance/optimizer/common/util/OrphanFileLocationUtils.java: ########## @@ -0,0 +1,78 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.gravitino.maintenance.optimizer.common.util; + +import com.google.common.base.Preconditions; +import java.net.URI; +import java.nio.file.Paths; +import java.util.Objects; + +/** Lexical containment checks shared by orphan cleanup submission and execution. */ +public final class OrphanFileLocationUtils { + private OrphanFileLocationUtils() {} Review Comment: Fixed in efb339668. Added the blank line after the private constructor and ran Spotless. ########## maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/handler/orphan/OrphanFileRemovalJobContext.java: ########## @@ -0,0 +1,73 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.gravitino.maintenance.optimizer.recommender.handler.orphan; + +import java.util.Map; +import javax.annotation.Nullable; +import org.apache.gravitino.NameIdentifier; +import org.apache.gravitino.maintenance.optimizer.api.recommender.JobExecutionContext; + +/** Immutable target and options for table-level orphan cleanup. */ +public class OrphanFileRemovalJobContext implements JobExecutionContext { + private final NameIdentifier name; + private final Map<String, String> jobOptions; + private final String jobTemplateName; + @Nullable private final String tableLocation; + + /** + * Creates the context from table metadata and policy options. + * + * @param name normalized catalog.schema.table identifier + * @param jobOptions policy options + * @param jobTemplateName template to submit + * @param tableLocation table storage root, if available + */ + public OrphanFileRemovalJobContext( + NameIdentifier name, + Map<String, String> jobOptions, + String jobTemplateName, + @Nullable String tableLocation) { + this.name = name; + this.jobOptions = Map.copyOf(jobOptions); + this.jobTemplateName = jobTemplateName; + this.tableLocation = tableLocation; + } + + @Override + public NameIdentifier nameIdentifier() { + return name; + } + + @Override + public Map<String, String> jobOptions() { + return jobOptions; + } + + @Override + public String jobTemplateName() { Review Comment: Fixed in efb339668. Added the blank line before tableLocation()'s Javadoc and ran Spotless. ########## api/src/main/java/org/apache/gravitino/policy/Policy.java: ########## @@ -42,10 +42,12 @@ enum BuiltInType { ICEBERG_COMPACTION( BUILT_IN_TYPE_PREFIX + "iceberg_compaction", IcebergDataCompactionContent.class), - /** - * Custom policy type. "custom" is a fixed string that indicates the policy is a non-built-in - * type. - */ + /** Iceberg orphan file cleanup policy. */ + ICEBERG_ORPHAN_FILE_REMOVAL( + BUILT_IN_TYPE_PREFIX + "iceberg_orphan_file_removal", + IcebergOrphanFileRemovalContent.class), + + /** Custom policy type. */ Review Comment: Fixed in efb339668. Restored the note that non-built-in policies use the fixed wire value custom. ########## maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/job/GravitinoOrphanFileRemovalJobAdapter.java: ########## @@ -0,0 +1,104 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.gravitino.maintenance.optimizer.recommender.job; + +import com.google.common.base.Preconditions; +import java.time.Clock; +import java.time.DateTimeException; +import java.time.ZoneOffset; +import java.time.format.DateTimeFormatter; +import java.time.temporal.ChronoUnit; +import java.util.Map; +import org.apache.gravitino.maintenance.optimizer.api.recommender.JobExecutionContext; +import org.apache.gravitino.maintenance.optimizer.common.util.IdentifierUtils; +import org.apache.gravitino.maintenance.optimizer.common.util.OrphanFileLocationUtils; +import org.apache.gravitino.maintenance.optimizer.recommender.handler.orphan.OrphanFileRemovalJobContext; +import org.apache.gravitino.policy.IcebergOrphanFileRemovalContent; + +/** Converts orphan cleanup options to the existing Spark template's configuration. */ +public class GravitinoOrphanFileRemovalJobAdapter implements GravitinoJobAdapter { + private static final String TABLE_IDENTIFIER_KEY = "table_identifier"; + private static final String OLDER_THAN_KEY = "older_than"; + private static final String LOCATION_KEY = "location"; + private static final String DRY_RUN_KEY = "dry_run"; + private static final DateTimeFormatter TIMESTAMP = + DateTimeFormatter.ofPattern("uuuu-MM-dd HH:mm:ssXXX").withZone(ZoneOffset.UTC); + private final Clock clock; + + /** Creates an adapter using the UTC system clock. */ + public GravitinoOrphanFileRemovalJobAdapter() { + this(Clock.systemUTC()); + } + + GravitinoOrphanFileRemovalJobAdapter(Clock clock) { + this.clock = clock; + } + + @Override + public Map<String, String> jobConfig(JobExecutionContext context) { + Preconditions.checkArgument( + context instanceof OrphanFileRemovalJobContext, + "jobExecutionContext must be OrphanFileRemovalJobContext"); + OrphanFileRemovalJobContext orphan = (OrphanFileRemovalJobContext) context; + Map<String, String> options = orphan.jobOptions(); + long days; + try { + days = + Long.parseLong( + options.getOrDefault( + IcebergOrphanFileRemovalContent.OLDER_THAN_DAYS_KEY, + String.valueOf(IcebergOrphanFileRemovalContent.DEFAULT_OLDER_THAN_DAYS))); + } catch (NumberFormatException e) { + throw new IllegalArgumentException("olderThanDays must be an integer", e); + } + IcebergOrphanFileRemovalContent.validateOlderThanDays(days); + String dryRun = + options.getOrDefault( + IcebergOrphanFileRemovalContent.DRY_RUN_KEY, + String.valueOf(IcebergOrphanFileRemovalContent.DEFAULT_DRY_RUN)); + Preconditions.checkArgument( + "true".equals(dryRun) || "false".equals(dryRun), "dryRun must be true or false"); + String location = options.getOrDefault(IcebergOrphanFileRemovalContent.LOCATION_KEY, ""); Review Comment: Fixed in efb339668 using the job-side option you suggested. execute() now treats null or blank location as the table root before containment and filesystem validation, independently of CLI parsing. Reproduced the empty-location failure against a real Iceberg table in local Spark, then added a regression covering empty/whitespace locations, dry-run preservation, and deletion. All orphan cleanup job tests passed. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
