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]

Reply via email to