lasdf1234 commented on code in PR #13531: URL: https://github.com/apache/gravitino/pull/13531#discussion_r4151002041
########## 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: **Convention:** `validateOlderThanDays` is a public static helper sitting between instance `validate()` and `equals()`. Per this repo’s class member order (statics with other statics, then instance methods), please move it up next to the constants / other statics (or immediately above `validate()`), matching how other API types keep validation helpers grouped. ########## 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: **Bug (validation gap):** `location` is only rejected when `trim()` is empty, but the raw value is stored as-is. So a create/update with e.g. `"location": " s3://bucket/table/data"` passes `validate()`, is persisted, then fails later in `OrphanFileLocationUtils.normalizeLocation` / adapter submit with a less actionable error. Please either: 1. reject locations where `location != location.trim()` (and prefer `StringUtils.isNotBlank` like compaction rewrite options), or 2. normalize to `location.trim()` in the constructor / factory before storing. Add a unit test for leading/trailing whitespace (the existing case only covers `" "`). ########## 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: **Style nit vs `IcebergCompactionContentDTO`:** missing blank line before this nested type, and between `olderThanDays()` / `location()` / `dryRun()` and their following Javadoc blocks. Please match the compaction DTO formatting (Google Java Style member spacing) so the two built-in content DTOs stay consistent for future edits. ########## 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: **Convention nit:** job template keys (`table_identifier`, `older_than`, `location`, `dry_run`) are still magic strings, while policy option keys were correctly centralized on `IcebergOrphanFileRemovalContent`. Not a functional bug (they match `IcebergRemoveOrphanFilesJob` placeholders), but named constants next to the adapter (or shared with the job template) would prevent a placeholder rename from silently desyncing submit vs execution. Optional if you want to keep scope tight — compaction still hardcodes too. ########## 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: **Docs consistency:** policy-driven submission converts `olderThanDays` to an **explicit UTC** cutoff (`GravitinoOrphanFileRemovalJobAdapter` + the non-UTC session Spark test), but the orphan job section above still describes `older_than` as interpreted in the **Spark session time zone**. Please state here that the adapter always emits a UTC timestamp (`…Z` / offset), so operators don’t mix “N days” policy semantics with a non-UTC `spark.sql.session.timeZone` the way direct job submission does. -- 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]
