laserninja commented on code in PR #13531: URL: https://github.com/apache/gravitino/pull/13531#discussion_r4147358520
########## 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: Fixed in 782715170. Content validation now accepts olderThanDays only in the range 1..36500 (approximately 100 years), with the same bound documented in OpenAPI. Added boundary tests and raw HTTP create/update tests confirming invalid values return HTTP 400; a rejected update preserves the existing policy. ########## 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: Fixed in 782715170. The DTO now calls PolicyContentDTO.super.validate() before toDomainContent().validate(), matching the compaction DTO. ########## 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: Fixed in 782715170. Added named property/rule keys and job option constants, and reused the option constants in the adapter. ########## 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: Fixed in 782715170. The DTO now uses IcebergOrphanFileRemovalContent.DEFAULT_DRY_RUN; the policy factory and adapter also use the shared default. ########## 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: Fixed in 782715170. The adapter reuses the content retention bound and uses minus(days, ChronoUnit.DAYS). Parsing and timestamp-range failures now produce an IllegalArgumentException identifying olderThanDays. Added tests for invalid values, the maximum accepted value, and timestamp underflow. ########## 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: Fixed in 782715170. Added TABLE_LOCATION_PROPERTY next to the handler, with a comment identifying it as the Iceberg catalog table-location property. -- 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]
