lasdf1234 commented on code in PR #11700:
URL: https://github.com/apache/gravitino/pull/11700#discussion_r4005490913


##########
design-docs/iceberg-remove-orphan-files-maintenance-job.md:
##########
@@ -0,0 +1,644 @@
+<!--
+  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.
+-->
+
+# Design: Built-in Iceberg Remove Orphan Files Maintenance Job
+
+| Field   | Value                                                        |
+| ------- | ------------------------------------------------------------ |
+| Status  | Draft                                                        |
+| Authors | @laserninja                                                  |
+| Created | 2026-06-16                                                   |
+| Issue   | [#11195](https://github.com/apache/gravitino/issues/11195)   |
+| Module  | `api`, `maintenance/jobs`, `maintenance/optimizer`            |
+
+---
+
+## 1. Background
+
+Orphan files accumulate in Iceberg table storage locations from failed writes,
+incomplete transactions, schema evolution, or concurrent operations. These 
files
+are no longer referenced by any table snapshot but remain on disk, wasting
+significant storage — especially in high-write-volume environments.
+
+The existing built-in maintenance jobs (`builtin-iceberg-rewrite-data-files` 
for
+data compaction, `builtin-iceberg-update-stats` for metrics, and
+`builtin-iceberg-expire-snapshots` for metadata cleanup) address data file
+optimization and snapshot lifecycle but do not cover orphan file removal.
+
+PR [#10500](https://github.com/apache/gravitino/pull/10500) added Trino-side
+delegation for `remove_orphan_files` as a procedure, but there is no
+server-side built-in job that can be triggered automatically via the Table
+Maintenance Service (Optimizer) policies.
+
+This design proposes adding full end-to-end support for Iceberg orphan file
+removal: from policy definition through strategy evaluation to Spark job
+execution.
+
+---
+
+## 2. Goals
+
+1. Add a new built-in policy type `system_iceberg_orphan_file_removal` for
+   declarative orphan file cleanup configuration.
+2. Add a strategy handler that evaluates when orphan file removal should run
+   based on time since last cleanup or table statistics.
+3. Add a job adapter that converts strategy evaluation results into job
+   configurations.
+4. Add the Spark job that executes Iceberg's `remove_orphan_files` procedure.
+5. Ensure the full flow works end-to-end: policy → strategy → job
+   submission → Spark execution.
+
+---
+
+## 3. Non-Goals
+
+- Snapshot expiration (separate Iceberg procedure, separate issue 
[#11194](https://github.com/apache/gravitino/issues/11194)).
+- Automatic policy creation — users must explicitly create and attach
+  policies.
+- Changes to the Optimizer scheduling framework itself.
+- Custom file-level filtering beyond what Iceberg's procedure supports.
+
+---
+
+## 4. Existing Architecture Overview
+
+The Gravitino maintenance module follows a layered architecture for automated
+table maintenance. The existing Iceberg compaction flow establishes the
+pattern:
+
+```
+Policy Creation (REST API)
+    ↓
+GravitinoStrategyProvider (loads policies as strategies)
+    ↓
+CompactionStrategyHandler (evaluates trigger / score expressions)
+    ↓
+CompactionJobContext → GravitinoCompactionJobAdapter (converts to job config)
+    ↓
+GravitinoJobSubmitter (submits job via REST)
+    ↓
+IcebergRewriteDataFilesJob (Spark execution)
+```
+
+### 4.1 Layer Summary
+
+| Layer | Compaction Components | Purpose |
+| --- | --- | --- |
+| **Policy** | `Policy.BuiltInType.ICEBERG_COMPACTION`, 
`IcebergDataCompactionContent` | Define configuration, thresholds, expressions |
+| **Strategy** | `CompactionStrategyHandler` extends 
`BaseExpressionStrategyHandler` | Evaluate trigger conditions, score partitions 
|
+| **Adapter** | `GravitinoCompactionJobAdapter`, `CompactionJobContext` | 
Convert evaluation result to job configuration |
+| **Job** | `IcebergRewriteDataFilesJob`, registered in 
`BuiltInJobTemplateProvider` | Execute Spark procedure |
+
+---
+
+## 5. Proposed Design
+
+We add the same four layers for orphan file removal, following the compaction
+pattern.
+
+### 5.1 Architecture Diagram
+
+```
+┌──────────────────────────────────────────────────────────────┐
+│  REST API: POST /metalakes/{m}/policies                      │
+│  type: "system_iceberg_orphan_file_removal"                  │
+│  content: IcebergOrphanFileRemovalContent                    │
+│  { olderThanDays, location, dryRun }                         │
+└──────────────────────────┬───────────────────────────────────┘
+                           ↓
+┌──────────────────────────────────────────────────────────────┐
+│  GravitinoStrategyProvider                                   │
+│  Loads policy → GravitinoStrategy                            │
+│  strategyType: "iceberg-orphan-file-removal"                 │
+│  jobTemplateName: "builtin-iceberg-remove-orphan-files"      │
+└──────────────────────────┬───────────────────────────────────┘
+                           ↓
+┌──────────────────────────────────────────────────────────────┐
+│  OrphanFileRemovalStrategyHandler                            │
+│  extends BaseExpressionStrategyHandler                       │
+│  dataRequirements: {TABLE_METADATA, TABLE_STATISTICS}        │
+│  Evaluates: custom-last-orphan-cleanup-time (statistic_meta) │
+│  Returns: StrategyEvaluation with score + context            │
+└──────────────────────────┬───────────────────────────────────┘
+                           ↓
+┌──────────────────────────────────────────────────────────────┐
+│  OrphanFileRemovalJobContext → JobAdapter                     │
+│  Extracts: older_than, location, dry_run                     │
+│  Builds: job config map for template substitution            │
+└──────────────────────────┬───────────────────────────────────┘
+                           ↓
+┌──────────────────────────────────────────────────────────────┐
+│  GravitinoJobSubmitter                                       │
+│  Template: "builtin-iceberg-remove-orphan-files"             │
+│  Submits via REST: POST /metalakes/{m}/jobs                  │
+└──────────────────────────┬───────────────────────────────────┘
+                           ↓
+┌──────────────────────────────────────────────────────────────┐
+│  IcebergRemoveOrphanFilesJob (Spark)                         │
+│  CALL catalog.system.remove_orphan_files(                    │
+│      table => '…', older_than => TIMESTAMP '…',              │
+│      location => '…', dry_run => bool)                       │
+└──────────────────────────────────────────────────────────────┘
+```
+
+---
+
+### 5.2 Layer 1 — Policy Definition (`api/`)
+
+#### 5.2.1 New Policy Type
+
+Add `ICEBERG_ORPHAN_FILE_REMOVAL` to `Policy.BuiltInType`:
+
+```java
+// api/src/main/java/org/apache/gravitino/policy/Policy.java
+enum BuiltInType {
+    ICEBERG_COMPACTION("system_iceberg_compaction",
+        IcebergDataCompactionContent.class),
+    ICEBERG_ORPHAN_FILE_REMOVAL("system_iceberg_orphan_file_removal",
+        IcebergOrphanFileRemovalContent.class),  // NEW
+    CUSTOM("custom", CustomContent.class);
+}
+```
+
+#### 5.2.2 New Policy Content Class
+
+Create `IcebergOrphanFileRemovalContent` following the
+`IcebergDataCompactionContent` pattern:
+
+```java
+// NEW: api/src/main/java/org/apache/gravitino/policy/
+//      IcebergOrphanFileRemovalContent.java
+public class IcebergOrphanFileRemovalContent implements PolicyContent {
+
+    // Strategy metadata
+    public static final String STRATEGY_TYPE_VALUE =
+        "iceberg-orphan-file-removal";
+    public static final String JOB_TEMPLATE_NAME_VALUE =
+        "builtin-iceberg-remove-orphan-files";
+
+    // Configurable fields
+    private final long olderThanDays;     // default: 3
+    private final String location;        // default: null (table location)
+    private final boolean dryRun;         // default: false
+
+    // Trigger / score expressions.
+    // The interval threshold comes from the uniform minimum-interval
+    // mechanism shared by all system built-in policies, not this content.
+    public static final String TRIGGER_EXPR =
+        "custom-days-since-last-orphan-cleanup >= minIntervalDays";
+    public static final String SCORE_EXPR =
+        "custom-days-since-last-orphan-cleanup";
+
+    // Defaults
+    public static final long DEFAULT_OLDER_THAN_DAYS = 3;
+    public static final boolean DEFAULT_DRY_RUN = false;
+}
+```
+
+#### 5.2.3 Policy Content Fields
+
+| Field | Type | Default | Description |
+| --- | --- | --- | --- |
+| `olderThanDays` | `long` | 3 | Only remove orphan files older than this many 
days. See [5.2.4](#524-why-olderthandays-defaults-to-3) for the rationale. |
+| `location` | `String` | null | Custom location to scan. When specified, 
**only** this location is scanned instead of the table's default location. Must 
be validated against the table's own location — see [Section 
6.1](#61-location-validation). If null, the table's registered storage location 
is used. |
+| `dryRun` | `boolean` | false | Preview-only mode — list orphan files without 
deleting |
+
+> **Note:** The minimum interval between runs is intentionally **not** a field
+> of this policy. A uniform minimum-interval mechanism will be defined across
+> all four system built-in policies and is out of scope for this design.
+
+#### 5.2.4 Why `olderThanDays` Defaults to 3
+
+The 3-day default mirrors Iceberg's own default for
+`remove_orphan_files` and exists to protect **in-flight writes**.
+
+Orphan file detection compares files on storage against files referenced by
+table metadata. A file written by a transaction that has not yet committed
+looks identical to an orphan file. If such a file is deleted, the in-flight
+commit fails or, worse, produces a table that references a missing file.
+A 3-day window is long enough to cover:
+
+- Long-running Spark / Flink write jobs that stage files before committing
+- Retried or paused jobs that resume hours or days later
+- Clock skew between the storage system and the job runtime
+
+**To remove all orphan files regardless of age**, set `olderThanDays` to `0`.
+The adapter then passes the current timestamp as `older_than`, so every
+unreferenced file is eligible for deletion.
+
+**This is unsafe while any writer is active** and should only be used when
+all writes to the table are known to be stopped — for example, during a
+maintenance window or when reclaiming storage from a decommissioned table.
+Run with `dryRun: true` first to review the file list.
+
+#### 5.2.5 Example Policy Creation
+
+```json
+POST /metalakes/default/policies
+{
+    "name": "remove_orphans_weekly",
+    "type": "system_iceberg_orphan_file_removal",
+    "comment": "Remove orphan files older than 3 days",
+    "enabled": true,
+    "content": {
+        "olderThanDays": 3,
+        "dryRun": false
+    }
+}
+```
+
+---
+
+### 5.3 Layer 2 — Strategy Handler (`maintenance/optimizer/`)
+
+#### 5.3.1 Strategy Handler
+
+```java
+// NEW: maintenance/optimizer/src/main/java/…/handler/orphan/
+//      OrphanFileRemovalStrategyHandler.java
+public class OrphanFileRemovalStrategyHandler
+        extends BaseExpressionStrategyHandler {
+
+    public static final String NAME = "iceberg-orphan-file-removal";
+
+    @Override
+    public String strategyType() {
+        return NAME;
+    }
+
+    @Override
+    public Set<DataRequirement> dataRequirements() {
+        return ImmutableSet.of(
+                DataRequirement.TABLE_METADATA,
+                DataRequirement.TABLE_STATISTICS);
+        // TABLE_STATISTICS supplies custom-last-orphan-cleanup-time,
+        // read from statistic_meta the same way compaction reads its metrics
+    }
+
+    @Override
+    protected JobExecutionContext buildJobExecutionContext(
+            NameIdentifier nameIdentifier,
+            Strategy strategy,
+            Table table,
+            List<PartitionPath> partitions,
+            Map<String, String> jobOptions) {
+        return new OrphanFileRemovalJobContext(
+                nameIdentifier, jobOptions, strategy.jobTemplateName());
+    }
+}
+```
+
+**Key design decision:** Orphan file removal operates at the **table level**,
+not partition level. Unlike compaction, which scores and selects individual
+partitions, `remove_orphan_files` scans the entire table's storage location.
+Therefore:
+
+- `dataRequirements()` excludes `PARTITION_STATISTICS`.
+- No partition scoring / selection logic is needed.
+
+**Trigger modes:** Gravitino supports two trigger mechanisms, and this
+design does not limit users to one:
+
+1. **Event trigger** — The strategy handler evaluates table metadata
+   (e.g., snapshot count changes, write events) and triggers cleanup when
+   conditions are met.
+2. **Time trigger** — The Optimizer's scheduling framework can invoke the
+   strategy handler periodically, and the handler decides whether cleanup
+   is needed based on the last cleanup time.
+
+Both modes use the same strategy handler; the difference is in how
+often the handler is invoked.
+
+#### 5.3.2 Tracking the Last Cleanup Time
+

Review Comment:
   This section can be omitted. Later on, I will have a dedicated issue to 
determine the scheduling time. Additionally, no statistical information needs 
to be stored in the "statistic_meta" table. After some further consideration, I 
realized that the execution time of orphan files can actually be obtained from 
the "gravitino" job table. This is a general logic that can be applied to other 
built-in strategies as well. Therefore, the "statistic_meta" table does not 
need to store any information, and this section can be removed.



-- 
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