lasdf1234 commented on code in PR #13531: URL: https://github.com/apache/gravitino/pull/13531#discussion_r4144570398
########## api/src/main/java/org/apache/gravitino/policy/IcebergOrphanFileRemovalContent.java: ########## @@ -0,0 +1,128 @@ +/* + * 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 { + /** 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", STRATEGY_TYPE_VALUE, "job.template-name", 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", "true"); + rules.put("score-expr", "1"); + rules.put("job.options.olderThanDays", olderThanDays); + rules.put("job.options.dryRun", dryRun); + if (location != null) { + rules.put("job.options.location", location); + } + return Collections.unmodifiableMap(rules); + } + + @Override + public void validate() { + PolicyContent.super.validate(); + Preconditions.checkArgument(olderThanDays >= 1, "olderThanDays must be at least 1"); Review Comment: Please also bound the upper end of `olderThanDays` here (not only `>= 1`). `GravitinoOrphanFileRemovalJobAdapter` does `clock.instant().minus(Duration.ofDays(days))`, and `Duration.ofDays(Long.MAX_VALUE)` overflows — the unit test already expects that path to throw, but a policy with a huge `olderThanDays` can still be created/stored and only fails at submit time. Align with `IcebergDataCompactionContent.validate()`: fail loud at content validation so REST create/update rejects bad values early. ########## api/src/main/java/org/apache/gravitino/policy/IcebergOrphanFileRemovalContent.java: ########## @@ -0,0 +1,128 @@ +/* + * 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 { + /** 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", STRATEGY_TYPE_VALUE, "job.template-name", JOB_TEMPLATE_NAME_VALUE); Review Comment: Convention nit vs `IcebergDataCompactionContent`: please introduce named constants for the property/rule keys (`STRATEGY_TYPE_KEY`, `JOB_TEMPLATE_NAME_KEY`, `TRIGGER_EXPR_KEY`, `SCORE_EXPR_KEY`, and the `job.options.*` keys) instead of string literals in `properties()`/`rules()`. Keeps the two built-in Iceberg policy contents consistent and avoids drift when adapters parse options. ########## common/src/main/java/org/apache/gravitino/dto/policy/PolicyContentDTO.java: ########## @@ -204,4 +206,67 @@ private PolicyContent toDomainContent() { rewriteOptions()); } } + /** Typed orphan cleanup policy content for REST requests and responses. */ + @EqualsAndHashCode + @ToString + @Builder(setterPrefix = "with") + @AllArgsConstructor(access = lombok.AccessLevel.PRIVATE) + class IcebergOrphanFileRemovalContentDTO implements PolicyContentDTO { + @JsonProperty("olderThanDays") + private Long olderThanDays; + + @JsonProperty("location") + private String location; + + @JsonProperty("dryRun") + private Boolean dryRun; + + private IcebergOrphanFileRemovalContentDTO() {} + + /** + * @return minimum age in days, defaulting to three + */ + public long olderThanDays() { + return olderThanDays == null + ? IcebergOrphanFileRemovalContent.DEFAULT_OLDER_THAN_DAYS + : olderThanDays; + } + /** + * @return optional scan location + */ + @Nullable + public String location() { + return location; + } + /** + * @return whether to list candidates without deleting them + */ + public boolean dryRun() { + return dryRun == null ? false : dryRun; + } + + @Override + public Set<MetadataObject.Type> supportedObjectTypes() { + return toDomainContent().supportedObjectTypes(); + } + + @Override + public Map<String, String> properties() { + return toDomainContent().properties(); + } + + @Override + public Map<String, Object> rules() { + return toDomainContent().rules(); + } + + @Override + public void validate() { Review Comment: Please call `PolicyContentDTO.super.validate()` before `toDomainContent().validate()`, same as `IcebergCompactionContentDTO.validate()`. Right now this DTO skips the interface-level checks the compaction content still runs. ########## maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/job/GravitinoOrphanFileRemovalJobAdapter.java: ########## @@ -0,0 +1,82 @@ +/* + * 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.Duration; +import java.time.ZoneOffset; +import java.time.format.DateTimeFormatter; +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 = + Long.parseLong( + options.getOrDefault( + "olderThanDays", + String.valueOf(IcebergOrphanFileRemovalContent.DEFAULT_OLDER_THAN_DAYS))); + Preconditions.checkArgument(days >= 1, "olderThanDays must be at least 1"); + String dryRun = options.getOrDefault("dryRun", "false"); + Preconditions.checkArgument( + "true".equals(dryRun) || "false".equals(dryRun), "dryRun must be true or false"); + String location = options.getOrDefault("location", ""); + if (options.containsKey("location")) { + 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( + "table_identifier", + IdentifierUtils.removeCatalogFromIdentifier(orphan.nameIdentifier()).toString(), + "older_than", + TIMESTAMP.format(clock.instant().minus(Duration.ofDays(days))), Review Comment: Related to the content-validation comment: even with adapter-side `days >= 1`, large values still blow up here via `Duration.ofDays`. Once content validates a sane upper bound, consider catching `ArithmeticException` / using `minus(days, ChronoUnit.DAYS)` with an explicit check so the error message stays about `olderThanDays` rather than an opaque overflow. ########## common/src/main/java/org/apache/gravitino/dto/policy/PolicyContentDTO.java: ########## @@ -204,4 +206,67 @@ private PolicyContent toDomainContent() { rewriteOptions()); } } + /** Typed orphan cleanup policy content for REST requests and responses. */ + @EqualsAndHashCode + @ToString + @Builder(setterPrefix = "with") + @AllArgsConstructor(access = lombok.AccessLevel.PRIVATE) + class IcebergOrphanFileRemovalContentDTO implements PolicyContentDTO { + @JsonProperty("olderThanDays") + private Long olderThanDays; + + @JsonProperty("location") + private String location; + + @JsonProperty("dryRun") + private Boolean dryRun; + + private IcebergOrphanFileRemovalContentDTO() {} + + /** + * @return minimum age in days, defaulting to three + */ + public long olderThanDays() { + return olderThanDays == null + ? IcebergOrphanFileRemovalContent.DEFAULT_OLDER_THAN_DAYS + : olderThanDays; + } + /** + * @return optional scan location + */ + @Nullable + public String location() { + return location; + } + /** + * @return whether to list candidates without deleting them + */ + public boolean dryRun() { + return dryRun == null ? false : dryRun; Review Comment: Prefer `IcebergOrphanFileRemovalContent.DEFAULT_DRY_RUN` (and the same pattern as `olderThanDays()`) instead of a bare `false`, so DTO defaults stay tied to the domain constants. ########## maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/handler/BaseExpressionStrategyHandler.java: ########## @@ -127,7 +127,8 @@ private int maxPartitionNum() { } private boolean isPartitionTable() { - return tableMetadata.partitioning().length > 0; + return dataRequirements().contains(DataRequirement.PARTITION_STATISTICS) Review Comment: This change is the right fix so table-level strategies (orphan cleanup) can run on partitioned tables without partition stats — thanks. Please add a short comment here explaining **why** `PARTITION_STATISTICS` gates `isPartitionTable()`, so a future handler doesn’t accidentally reintroduce “partitioned ⇒ need partition stats” for strategies that intentionally declare only `TABLE_METADATA`. ########## docs/open-api/policies.yaml: ########## @@ -417,10 +454,12 @@ components: description: A basic policy without content oneOf: - $ref: '#/components/schemas/CustomPolicy' + - $ref: '#/components/schemas/IcebergOrphanFileRemovalPolicy' discriminator: propertyName: policyType mapping: custom: '#/components/schemas/CustomPolicy' + system_iceberg_orphan_file_removal: '#/components/schemas/IcebergOrphanFileRemovalPolicy' Review Comment: OpenAPI now discriminates `system_iceberg_orphan_file_removal`, which is good, but `system_iceberg_compaction` is still missing from these `oneOf` / mapping blocks while Jackson `@JsonSubTypes` registers both. Either add compaction to the same OpenAPI discriminators in this PR (parity with Java), or call out the gap explicitly so we don’t ship one built-in type documented in OpenAPI and another only in code. ########## maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/handler/orphan/OrphanFileRemovalStrategyHandler.java: ########## @@ -0,0 +1,60 @@ +/* + * 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.List; +import java.util.Map; +import java.util.Set; +import org.apache.gravitino.NameIdentifier; +import org.apache.gravitino.maintenance.optimizer.api.common.PartitionPath; +import org.apache.gravitino.maintenance.optimizer.api.common.Strategy; +import org.apache.gravitino.maintenance.optimizer.api.recommender.JobExecutionContext; +import org.apache.gravitino.maintenance.optimizer.recommender.handler.BaseExpressionStrategyHandler; +import org.apache.gravitino.policy.IcebergOrphanFileRemovalContent; +import org.apache.gravitino.rel.Table; + +/** Evaluates orphan cleanup at table level, including for partitioned tables. */ +public class OrphanFileRemovalStrategyHandler extends BaseExpressionStrategyHandler { + /** Registered strategy type. */ + public static final String NAME = IcebergOrphanFileRemovalContent.STRATEGY_TYPE_VALUE; + + @Override + public Set<DataRequirement> dataRequirements() { + return Set.of(DataRequirement.TABLE_METADATA); + } + + @Override + public String strategyType() { + return NAME; + } + + @Override + protected JobExecutionContext buildJobExecutionContext( + NameIdentifier nameIdentifier, + Strategy strategy, + Table tableMetadata, + List<PartitionPath> partitions, + Map<String, String> jobOptions) { + return new OrphanFileRemovalJobContext( + nameIdentifier, + jobOptions, + strategy.jobTemplateName(), + tableMetadata.properties() == null ? null : tableMetadata.properties().get("location")); Review Comment: Please avoid the magic `"location"` string. Elsewhere Iceberg uses `IcebergTable.PROP_LOCATION` / `IcebergConstants.LOCATION` (`"location"`). If depending on `catalog-lakehouse-iceberg` from optimizer is undesirable, at least put a named constant next to the handler / shared util so this doesn’t silently diverge from the catalog property key. -- 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]
