roryqi commented on code in PR #11700: URL: https://github.com/apache/gravitino/pull/11700#discussion_r3429860005
########## design-docs/iceberg-remove-orphan-files-maintenance-job.md: ########## @@ -0,0 +1,534 @@ +<!-- + 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} │ +│ Evaluates: time since last cleanup ≥ threshold │ +│ 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 + private final long cleanupIntervalDays; // trigger threshold, default: 7 + + // Trigger / score expressions + public static final String TRIGGER_EXPR = + "days-since-last-orphan-cleanup >= cleanupIntervalDays"; + public static final String SCORE_EXPR = + "days-since-last-orphan-cleanup"; + + // Defaults + public static final long DEFAULT_OLDER_THAN_DAYS = 3; + public static final boolean DEFAULT_DRY_RUN = false; + public static final long DEFAULT_CLEANUP_INTERVAL_DAYS = 7; +} +``` + +#### 5.2.3 Policy Content Fields + +| Field | Type | Default | Description | +| --------------------- | --------- | ------- | ------------------------------------------------------------------------ | +| `olderThanDays` | `long` | 3 | Only remove orphan files older than this many days | +| `location` | `String` | null | Custom location to scan (null = table's default location) | +| `dryRun` | `boolean` | false | Preview-only mode — list orphan files without deleting | +| `cleanupIntervalDays` | `long` | 7 | Trigger threshold — only run when days since last cleanup exceeds this | + +#### 5.2.4 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, run weekly", + "enabled": true, + "content": { + "olderThanDays": 3, + "dryRun": false, + "cleanupIntervalDays": 7 + } +} +``` + +--- + +### 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); + // No TABLE_STATISTICS or PARTITION_STATISTICS needed — + // orphan file removal is table-level, time-driven + } + + @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()` only includes `TABLE_METADATA`. +- No partition scoring / selection logic is needed. +- The trigger expression evaluates time since last cleanup, not data + statistics. + +#### 5.3.2 Job Execution Context + +```java +// NEW: maintenance/optimizer/src/main/java/…/handler/orphan/ +// OrphanFileRemovalJobContext.java +public class OrphanFileRemovalJobContext implements JobExecutionContext { + private final NameIdentifier nameIdentifier; + private final Map<String, String> jobOptions; + private final String jobTemplateName; + + // jobOptions keys: + // older_than, location, dry_run, Review Comment: Why do we need location here? -- 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]
