laserninja commented on code in PR #13531: URL: https://github.com/apache/gravitino/pull/13531#discussion_r4151295755
########## api/src/main/java/org/apache/gravitino/policy/IcebergOrphanFileRemovalContent.java: ########## @@ -0,0 +1,159 @@ +/* + * 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.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"); + rules.put(SCORE_EXPR_KEY, "1"); + rules.put(JOB_OPTIONS_PREFIX + OLDER_THAN_DAYS_KEY, olderThanDays); + rules.put(JOB_OPTIONS_PREFIX + DRY_RUN_KEY, dryRun); + if (location != null) { + rules.put(JOB_OPTIONS_PREFIX + LOCATION_KEY, location); + } + return Collections.unmodifiableMap(rules); + } + + @Override + public void validate() { + PolicyContent.super.validate(); + validateOlderThanDays(olderThanDays); + Preconditions.checkArgument( + location == null || !location.trim().isEmpty(), "location must not be blank"); Review Comment: Fixed in d91caf8a1. Blank locations and leading/trailing whitespace (including Unicode whitespace) are now rejected during content validation. Reproduced the original bug through HTTP create returning 200, then added unit and raw HTTP create/update tests asserting 400 and preservation of the existing policy after rejected updates. All affected suites and the final regression tests passed. ########## api/src/main/java/org/apache/gravitino/policy/IcebergOrphanFileRemovalContent.java: ########## @@ -0,0 +1,159 @@ +/* + * 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.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"); + rules.put(SCORE_EXPR_KEY, "1"); + rules.put(JOB_OPTIONS_PREFIX + OLDER_THAN_DAYS_KEY, olderThanDays); + rules.put(JOB_OPTIONS_PREFIX + DRY_RUN_KEY, dryRun); + if (location != null) { + rules.put(JOB_OPTIONS_PREFIX + LOCATION_KEY, location); + } + return Collections.unmodifiableMap(rules); + } + + @Override + public void validate() { + PolicyContent.super.validate(); + validateOlderThanDays(olderThanDays); + Preconditions.checkArgument( + location == null || !location.trim().isEmpty(), "location must not be blank"); + } + + /** + * Validates retention before storing a policy or submitting a job. + * + * @param days minimum file age in days + * @throws IllegalArgumentException if days is outside the supported range + */ + public static void validateOlderThanDays(long days) { Review Comment: Fixed in d91caf8a1. Moved validateOlderThanDays immediately above validate(), as suggested. ########## docs/table-maintenance-service/optimizer-configuration.md: ########## @@ -160,5 +160,51 @@ not supported. For a secured Iceberg REST catalog, supply its authentication settings explicitly in `spark_conf`, as described above. See [Remove Orphan Files](./optimizer-cli-reference.md#remove-orphan-files) for a -complete submission example. Orphan cleanup has no built-in scheduling policy -in this release. +complete submission example. For policy-driven submission, configure the handler +and policy below. Periodic scheduling remains separate. + +### Orphan cleanup policy integration + +Register the handler alongside the existing optimizer providers: + +```properties +gravitino.optimizer.strategyHandler.iceberg-orphan-file-removal.className = org.apache.gravitino.maintenance.optimizer.recommender.handler.orphan.OrphanFileRemovalStrategyHandler +``` + +The built-in job adapter is registered automatically. Keep the same Spark and +catalog submission configuration as for direct orphan cleanup, including +`catalog_name` (the catalog alias configured in Spark) and `spark_conf`. +The adapter supplies `table_identifier`, `older_than`, `location`, and `dry_run`. Review Comment: Fixed in d91caf8a1. The policy section now explicitly states that the adapter emits a UTC cutoff ending in Z, independent of spark.sql.session.timeZone. The direct-job parameter description also clarifies that the session time zone applies when no offset is supplied. ########## common/src/main/java/org/apache/gravitino/dto/policy/PolicyContentDTO.java: ########## @@ -204,4 +206,68 @@ private PolicyContent toDomainContent() { rewriteOptions()); } } + /** Typed orphan cleanup policy content for REST requests and responses. */ Review Comment: Fixed in d91caf8a1. Added the missing blank lines before the nested DTO and between its accessor methods and subsequent Javadocs. Ran Spotless. ########## maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/job/GravitinoOrphanFileRemovalJobAdapter.java: ########## @@ -0,0 +1,100 @@ +/* + * 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 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, ""); + if (options.containsKey(IcebergOrphanFileRemovalContent.LOCATION_KEY)) { + Preconditions.checkArgument( + orphan.tableLocation() != null, + "Table location is required to validate a custom scan location"); + OrphanFileLocationUtils.validateLocation(orphan.tableLocation(), location); + location = OrphanFileLocationUtils.normalizeLocation(location).toString(); + } + return Map.of( Review Comment: Fixed in d91caf8a1. Added named constants for the four job-template keys. Also addressed your review-summary request with a TestGravitinoJobSubmitter case verifying the built-in orphan template resolves to GravitinoOrphanFileRemovalJobAdapter. The submitter and adapter 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]
