This is an automated email from the ASF dual-hosted git repository.
roryqi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new b86048be92 [#11195] feat(maintenance): add Iceberg orphan file cleanup
job (#13153)
b86048be92 is described below
commit b86048be927f69708ee2fc8ed0265f5c8b5225ea
Author: Akshay Thorat <[email protected]>
AuthorDate: Wed Sep 23 18:45:34 2026 -0700
[#11195] feat(maintenance): add Iceberg orphan file cleanup job (#13153)
### What changes were proposed in this pull request?
Add and register `builtin-iceberg-remove-orphan-files`, with cutoff,
scan location, dry-run, and Spark configuration options. Validate that
the scan stays within the target table, reject symlinks within the
table, and log dry-run candidates. Bound remote ancestor inspection to
the table root so warehouse parent symlinks and inaccessible parents do
not block cleanup. Retain the shared runtime check, CLI error reporting,
and Spark cleanup. Spark renders procedure string literals so quoted
names work without changing shared utilities. Defer Spark-specific
exception handling until execution so the server can discover the
template without Spark installed.
Update the optimizer overview, configuration guide, CLI reference, and
design document. This is the job-layer implementation from #11700;
policy, strategy, and adapter integration will follow separately.
### Why are the changes needed?
Failed or incomplete writes leave unreferenced files in table storage.
This job provides direct orphan cleanup through the existing job
submission API.
Related to #11195.
### Does this PR introduce _any_ user-facing change?
Adds a v1 Spark job template. The cutoff defaults to three days ago;
explicit cutoffs must be at least 24 hours old, including for dry runs.
Dry-run defaults to false. Custom scan locations must remain within the
table's storage location.
The submission example supplies all template placeholders, using empty
strings for optional cutoff and location defaults. Secured REST catalogs
use explicit Spark configuration.
### How was this patch tested?
With JDK 17:
```bash
./gradlew :maintenance:jobs:spotlessApply :maintenance:jobs:build
:maintenance:jobs:javadoc -PskipWeb=true
```
All 165 jobs tests passed, including real local Spark tests for template
arguments, dry-run preservation, orphan deletion, retention of
referenced and recent files, explicit cutoffs, quoted identifiers, and
unsafe input/location rejection. Subprocess tests verify CLI failure
status and the missing-runtime diagnostic. Remote filesystem validation
uses an in-memory Hadoop filesystem stub, with regressions for warehouse
parent symlinks, inaccessible parents, normalized paths, and rejection
of links inside the table. An isolated classloader test reproduces
server-side template discovery without Spark, Iceberg, or Hadoop.
Docker-dependent tests were not enabled. Javadoc generation passed with
seven existing warnings in unchanged files.
---
.../iceberg-remove-orphan-files-maintenance-job.md | 30 +--
.../optimizer-cli-reference.md | 66 ++++-
.../optimizer-configuration.md | 27 ++
docs/table-maintenance-service/optimizer.md | 15 +-
.../jobs/BuiltInJobTemplateProvider.java | 4 +-
.../jobs/iceberg/IcebergRemoveOrphanFilesJob.java | 298 +++++++++++++++++++++
.../jobs/iceberg/RemoteLocationValidator.java | 82 ++++++
.../iceberg/TestIcebergRemoveOrphanFilesJob.java | 133 +++++++++
.../TestIcebergRemoveOrphanFilesJobMain.java | 144 ++++++++++
.../TestIcebergRemoveOrphanFilesJobWithSpark.java | 192 +++++++++++++
.../jobs/iceberg/TestRemoteLocationValidator.java | 197 ++++++++++++++
11 files changed, 1169 insertions(+), 19 deletions(-)
diff --git a/design-docs/iceberg-remove-orphan-files-maintenance-job.md
b/design-docs/iceberg-remove-orphan-files-maintenance-job.md
index 885f78cb62..74421e3bc6 100644
--- a/design-docs/iceberg-remove-orphan-files-maintenance-job.md
+++ b/design-docs/iceberg-remove-orphan-files-maintenance-job.md
@@ -209,7 +209,7 @@ public class IcebergOrphanFileRemovalContent implements
PolicyContent {
| 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.
|
+| `olderThanDays` | `long` | 3 | Only remove orphan files older than
this many days; must be at least 1. 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
|
@@ -232,13 +232,11 @@ A 3-day window is long enough to cover:
- 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.
+**Explicit cutoffs must be at least 24 hours old.** Iceberg's Spark SQL
+procedure enforces this minimum, including for dry runs. The planned policy
+must require `olderThanDays >= 1`; `olderThanDays: 0` is not supported.
+The three-day default remains the recommended starting point for protecting
+in-flight writes. Use a longer interval if writers can run longer than that.
Run with `dryRun: true` first to review the file list.
#### 5.2.5 Example Policy Creation
@@ -362,10 +360,11 @@ public class GravitinoOrphanFileRemovalJobAdapter
+ "." + ctx.nameIdentifier().name());
// Convert olderThanDays → absolute timestamp
- // olderThanDays == 0 means "now", i.e. remove all orphan files
+ // Iceberg requires explicit cutoffs to be at least 24 hours old.
Map<String, String> opts = ctx.jobOptions();
long days = Long.parseLong(
opts.getOrDefault("olderThanDays", "3"));
+ Preconditions.checkArgument(days >= 1, "olderThanDays must be at least
1");
String ts = Instant.now()
.minus(Duration.ofDays(days))
.toString()
@@ -529,9 +528,9 @@ The `older_than` threshold should be set conservatively.
Files from in-flight
writes or concurrent operations may not yet be referenced by a committed
snapshot. A minimum of 3 days is recommended.
-Setting `olderThanDays` to `0` removes all orphan files regardless of age.
-This is only safe when no writer is active against the table - see
-[Section 5.2.4](#524-why-olderthandays-defaults-to-3).
+The Spark procedure rejects explicit cutoffs less than 24 hours old, even
+for dry runs. `olderThanDays` must be at least `1`; zero-day cleanup is not
+supported. See [Section 5.2.4](#524-why-olderthandays-defaults-to-3).
---
@@ -591,10 +590,9 @@ combined into a single PR.
somewhere (e.g., job output metadata) for review before actual deletion?
4. ~~**PR granularity**~~ - Resolved: single PR for policy + strategy +
adapter layers since total code is expected to be under 1000 lines.
-5. ~~**`older_than` minimum**~~ - Resolved: no hard minimum is enforced.
- `olderThanDays: 0` is a deliberate escape hatch for reclaiming storage
- when no writer is active. See
- [Section 5.2.4](#524-why-olderthandays-defaults-to-3).
+5. ~~**`older_than` minimum**~~ - Resolved: the Spark SQL procedure enforces
+ a 24-hour minimum for explicit cutoffs. The policy must reject
+ `olderThanDays < 1`. See [Section
5.2.4](#524-why-olderthandays-defaults-to-3).
6. **Minimum run interval** - Out of scope. A uniform minimum-interval
mechanism will be defined across all four system built-in policies.
diff --git a/docs/table-maintenance-service/optimizer-cli-reference.md
b/docs/table-maintenance-service/optimizer-cli-reference.md
index c93b2bde9b..cba7e1cad0 100644
--- a/docs/table-maintenance-service/optimizer-cli-reference.md
+++ b/docs/table-maintenance-service/optimizer-cli-reference.md
@@ -237,13 +237,14 @@ EvaluationResult{scopeType=TABLE,
identifier=rest_catalog.db.t1, partitionPath=<
## Built-in Job Templates
-Three job templates ship with the service, and they are complementary rather
than alternatives. A full maintenance pass collects statistics, compacts data
files, and then expires the snapshot history that compaction just created.
+Four job templates ship with the service, and they are complementary rather
than alternatives. A full maintenance pass collects statistics, compacts data
files, expires the snapshot history that compaction just created, and removes
old orphan files.
| Job template | What it does
|
|---------------------------------------|-------------------------------------------|
| `builtin-iceberg-update-stats` | Collects file statistics and metrics
|
| `builtin-iceberg-rewrite-data-files` | Compacts small data files
|
| `builtin-iceberg-expire-snapshots` | Removes old snapshot metadata
|
+| `builtin-iceberg-remove-orphan-files` | Removes unreferenced files from
storage |
Each can be submitted directly over REST, and the first two are also what the
policy-driven workflow submits on your behalf. See [Quick
Start](./optimizer.md#walkthrough) for the policy-driven path.
@@ -365,3 +366,66 @@ Expire Snapshots Results:
- [Configuration](./optimizer-configuration.md) for the three configuration
layers
- [Iceberg Compaction Policy](../iceberg-compaction-policy.md) for tuning the
built-in strategy
- [Manage Jobs](../manage-jobs-in-gravitino.md) for job status and templates
+
+## Remove Orphan Files
+
+`builtin-iceberg-remove-orphan-files` runs Iceberg's `remove_orphan_files`
Spark
+procedure. It removes files in the scan location that are no longer referenced
+by table metadata. This job is available for direct submission; policy-driven
+scheduling is a separate feature.
+
+| Key | Description
| Default |
+| ------------------ |
----------------------------------------------------------------------------------------------
| -------------------------------- |
+| `catalog_name` | Iceberg catalog registered in Spark
| Required |
+| `table_identifier` | Table identifier within that catalog, such as
`db.sample` | Required
|
+| `older_than` | Cutoff timestamp in the Spark session time zone;
explicit values must be at least 24 hours old | Three days ago (Iceberg
default) |
+| `location` | Scan only this directory within the table's storage
location | Table location |
+| `dry_run` | `true` logs candidate paths without deleting; `false`
deletes | `false` |
+| `spark_conf` | JSON map of custom Spark configuration
| None |
+
+The template uses the same Spark and catalog connection settings as the other
+Iceberg jobs. Supply every template placeholder in `jobConf`: use empty strings
+for `older_than` and `location` to keep their defaults, an explicit boolean
string
+for `dry_run`, and `{}` for `spark_conf` when no overrides are needed.
+For example, submit a preview using:
+
+```json
+{
+ "jobTemplateName": "builtin-iceberg-remove-orphan-files",
+ "jobConf": {
+ "catalog_name": "rest_catalog",
+ "table_identifier": "db.t1",
+ "older_than": "",
+ "location": "",
+ "dry_run": "true",
+ "spark_conf": "{}",
+ "spark_master": "local[2]",
+ "spark_executor_instances": "1",
+ "spark_executor_cores": "1",
+ "spark_executor_memory": "1g",
+ "spark_driver_memory": "1g",
+ "catalog_type": "rest",
+ "catalog_uri": "http://localhost:9001/iceberg",
+ "warehouse_location": ""
+ }
+}
+```
+
+POST this body to `/api/metalakes/{metalake}/jobs`. Review the candidate paths
in
+the job logs before resubmitting with `dry_run: "false"`. A successful run also
+logs the candidate count. CLI equivalents are `--catalog`, `--table`,
+`--older-than`, `--location`, `--dry-run true|false`, and `--spark-conf`.
+
+The job validates the location before executing the procedure. A custom
+location must be the table's own location or a descendant, on the same storage
+scheme and authority. Relative paths, ambiguous percent-encoded paths, query
+strings, fragments, and symlinks in the scan directory are rejected. Filesystem
+validation errors fail the job before deletion. Keep the scan directory free of
+concurrent location or symlink changes during cleanup.
+
+Retain the three-day default unless a longer interval is needed for your
writers.
+Files staged by active writers can appear to be orphaned. Iceberg 1.11's SQL
+procedure rejects explicit cutoffs less than 24 hours old, including a cutoff
of
+"now", even for dry runs. This job preserves that safeguard and does not enable
+Iceberg's testing override. Spark's `spark.sql.parser.escapedStringLiterals`
must
+remain `false` for procedure arguments to be interpreted correctly.
diff --git a/docs/table-maintenance-service/optimizer-configuration.md
b/docs/table-maintenance-service/optimizer-configuration.md
index ae6350e9b2..f356f97cf7 100644
--- a/docs/table-maintenance-service/optimizer-configuration.md
+++ b/docs/table-maintenance-service/optimizer-configuration.md
@@ -135,3 +135,30 @@ Four things are worth confirming before assuming a
configuration problem is a co
- [CLI Reference](./optimizer-cli-reference.md) for every command and the
built-in job templates
- [Troubleshooting](./optimizer-troubleshooting.md) when a command or job fails
- [Extension Guide](./optimizer-extension-guide.md) for custom strategies and
providers
+
+## Orphan File Cleanup Job Configuration
+
+Submit `builtin-iceberg-remove-orphan-files` with the same Spark and catalog
+settings described above. Its job-specific `jobConf` keys are:
+
+| Key | Meaning
| Default |
+| ------------------ |
-------------------------------------------------------------------------------------
| -------------------------------- |
+| `catalog_name` | Iceberg catalog registered in Spark
| Required |
+| `table_identifier` | Table identifier within the catalog, such as
`db.sample` | Required |
+| `older_than` | Timestamp in the Spark session time zone; must be at
least 24 hours old | Three days ago (Iceberg default) |
+| `location` | Scan only this directory within the table's storage
location | Table location |
+| `dry_run` | `true` logs candidate paths without deleting; `false`
deletes | `false` |
+| `spark_conf` | JSON string containing custom Spark settings, including
the Iceberg runtime if needed | None |
+
+For direct template submission, supply `older_than: ""` and `location: ""` to
+use the defaults, `dry_run: "false"` (or `"true"` to preview), and
+`spark_conf: "{}"` when no overrides are needed.
+
+Keep the three-day default unless your workload needs a longer retention
window.
+The 24-hour minimum also applies to dry runs; passing the current timestamp is
+not supported. For a secured Iceberg REST catalog, supply its authentication
+settings explicitly in `spark_conf`, as described above.
+
+See [Remove Orphan Files](./optimizer-cli-reference.md#remove-orphan-files)
for a
+complete submission example. Orphan cleanup has no built-in scheduling policy
+in this release.
diff --git a/docs/table-maintenance-service/optimizer.md
b/docs/table-maintenance-service/optimizer.md
index b1c66b557c..774591d55a 100644
--- a/docs/table-maintenance-service/optimizer.md
+++ b/docs/table-maintenance-service/optimizer.md
@@ -21,7 +21,7 @@ The CLI binary, its configuration file, and its configuration
keys carry the old
Confirm your environment matches this list before starting an evaluation
against the built-ins. Anything outside it needs a custom extension, which is
covered in the [Extension Guide](./optimizer-extension-guide.md).
-- Compaction is the only built-in strategy. There is no built-in snapshot
expiration, orphan file cleanup, or sort and cluster maintenance.
+- Compaction is the only built-in strategy. Snapshot expiration and orphan
file cleanup are available as directly submitted built-in jobs, but do not yet
have built-in scheduling strategies. Sort and cluster maintenance also require
custom strategies.
- Compaction applies to Iceberg tables only, and only where every partition
uses an identity transform.
- The service is driven through the CLI workflow rather than running on a
schedule of its own.
@@ -57,6 +57,19 @@ Three identifiers look interchangeable and are not.
`--strategy-name` takes the **policy name**, despite what it is called.
Passing either of the other two reports no matching identifiers rather than
naming the mistake.
+## Direct Orphan File Cleanup
+
+Submit `builtin-iceberg-remove-orphan-files` through the jobs REST API to
reclaim
+unreferenced files. Start with `dry_run: "true"` and review the candidate paths
+in the job logs before allowing deletion. The default cutoff is three days ago;
+explicit cutoffs must be at least 24 hours old, including for dry runs. A
custom
+scan location must remain within the target table's storage location.
+
+This is a directly submitted job, not a new scheduling strategy. See
+[Remove Orphan Files](./optimizer-cli-reference.md#remove-orphan-files) for the
+submission example and
[Configuration](./optimizer-configuration.md#orphan-file-cleanup-job-configuration)
+for the per-job options.
+
## Walkthrough
This takes one Iceberg table through the whole workflow: create it, fill it
with small files, attach a compaction policy, collect statistics, and let the
service decide to compact it. It runs against a local Spark and takes about
fifteen minutes.
diff --git
a/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/BuiltInJobTemplateProvider.java
b/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/BuiltInJobTemplateProvider.java
index 3053c02d96..471824af43 100644
---
a/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/BuiltInJobTemplateProvider.java
+++
b/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/BuiltInJobTemplateProvider.java
@@ -26,6 +26,7 @@ import java.util.stream.Collectors;
import org.apache.gravitino.job.JobTemplate;
import org.apache.gravitino.job.JobTemplateProvider;
import org.apache.gravitino.maintenance.jobs.iceberg.IcebergExpireSnapshotsJob;
+import
org.apache.gravitino.maintenance.jobs.iceberg.IcebergRemoveOrphanFilesJob;
import
org.apache.gravitino.maintenance.jobs.iceberg.IcebergRewriteDataFilesJob;
import
org.apache.gravitino.maintenance.jobs.iceberg.IcebergUpdateStatsAndMetricsJob;
import org.apache.gravitino.maintenance.jobs.spark.SparkPiJob;
@@ -47,7 +48,8 @@ public class BuiltInJobTemplateProvider implements
JobTemplateProvider {
new SparkPiJob(),
new IcebergRewriteDataFilesJob(),
new IcebergUpdateStatsAndMetricsJob(),
- new IcebergExpireSnapshotsJob());
+ new IcebergExpireSnapshotsJob(),
+ new IcebergRemoveOrphanFilesJob());
@Override
public List<? extends JobTemplate> jobTemplates() {
diff --git
a/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergRemoveOrphanFilesJob.java
b/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergRemoveOrphanFilesJob.java
new file mode 100644
index 0000000000..e8569a966d
--- /dev/null
+++
b/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergRemoveOrphanFilesJob.java
@@ -0,0 +1,298 @@
+/*
+ * 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.jobs.iceberg;
+
+import static org.apache.spark.sql.functions.lit;
+
+import com.google.common.base.Preconditions;
+import java.io.IOException;
+import java.net.URI;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.Paths;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.stream.Stream;
+import javax.annotation.Nullable;
+import org.apache.gravitino.job.JobTemplateProvider;
+import org.apache.gravitino.job.SparkJobTemplate;
+import org.apache.gravitino.maintenance.jobs.BuiltInJob;
+import
org.apache.gravitino.maintenance.optimizer.common.util.IcebergSparkConfigUtils;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.spark.Spark3Util;
+import org.apache.spark.sql.AnalysisException;
+import org.apache.spark.sql.Row;
+import org.apache.spark.sql.SparkSession;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/** Removes unreferenced Iceberg files after validating the requested scan
location. */
+public class IcebergRemoveOrphanFilesJob implements BuiltInJob {
+ private static final Logger LOG =
LoggerFactory.getLogger(IcebergRemoveOrphanFilesJob.class);
+ private static final String NAME =
+ JobTemplateProvider.BUILTIN_NAME_PREFIX + "iceberg-remove-orphan-files";
+ private static final String VERSION = "v1";
+
+ @Override
+ public SparkJobTemplate jobTemplate() {
+ return SparkJobTemplate.builder()
+ .withName(NAME)
+ .withComment("Built-in Iceberg orphan file cleanup job template")
+ .withExecutable(resolveExecutable(IcebergRemoveOrphanFilesJob.class))
+ .withClassName(IcebergRemoveOrphanFilesJob.class.getName())
+ .withArguments(buildArguments())
+ .withConfigs(buildSparkConfigs())
+ .withCustomFields(
+ Collections.singletonMap(JobTemplateProvider.PROPERTY_VERSION_KEY,
VERSION))
+ .build();
+ }
+
+ /**
+ * Runs orphan file cleanup using named arguments.
+ *
+ * <p>Required: {@code --catalog name --table db.table}. Optional: {@code
--older-than 'yyyy-MM-dd
+ * HH:mm:ss'}, {@code --location path}, {@code --dry-run true|false}, and
{@code --spark-conf
+ * json}. The cutoff defaults to three days ago and dry-run defaults to
false. Iceberg's minimum
+ * retention interval is preserved. A custom location must be within the
table.
+ *
+ * @param args named command-line arguments
+ */
+ public static void main(String[] args) {
+ int exitCode = run(args);
+ if (exitCode != 0) {
+ System.exit(exitCode);
+ }
+ }
+
+ static int run(String[] args) {
+ return Runner.run(args);
+ }
+
+ static long execute(SparkSession spark, Map<String, String> options)
+ throws IOException, AnalysisException {
+ String catalog = requireOption(options, "catalog");
+ String identifier = requireOption(options, "table");
+ boolean dryRun = parseDryRun(options.get("dry-run"));
+ // Backslash escapes in string literals must retain Spark's default
interpretation.
+ Preconditions.checkArgument(
+
!Boolean.parseBoolean(spark.conf().get("spark.sql.parser.escapedStringLiterals",
"false")),
+ "spark.sql.parser.escapedStringLiterals must be false");
+ Table table =
+ Spark3Util.loadIcebergTable(
+ spark, IcebergJobUtils.escapeSqlIdentifier(catalog) + "." +
identifier);
+ String location = options.getOrDefault("location", table.location());
+ validateLocation(table.location(), location);
+ validateLocalLocation(table.location(), location);
+ validateRemoteLocation(spark, table.location(), location);
+ String sql =
+ buildProcedureCall(
+ catalog,
+ identifier,
+ options.get("older-than"),
+ normalizeLocation(location).toString(),
+ dryRun);
+ Iterator<Row> results = spark.sql(sql).toLocalIterator();
+ long count = 0;
+ while (results.hasNext()) {
+ Row row = results.next();
+ if (dryRun) {
+ LOG.info("Orphan file (dry-run): {}", row.getString(0));
+ }
+ count++;
+ }
+ LOG.info(
+ "Orphan file cleanup completed: {} files {}, table {}.{}",
+ count,
+ dryRun ? "found (dry-run)" : "removed",
+ catalog,
+ identifier);
+ return count;
+ }
+
+ static String buildProcedureCall(
+ String catalog,
+ String table,
+ @Nullable String olderThan,
+ @Nullable String location,
+ boolean dryRun) {
+ // Render string literals using Spark so quotes and backslashes round-trip
correctly.
+ StringBuilder sql =
+ new StringBuilder("CALL ")
+ .append(IcebergJobUtils.escapeSqlIdentifier(catalog))
+ .append(".system.remove_orphan_files(table => ")
+ .append(lit(table).expr().sql());
+ if (olderThan != null && !olderThan.isEmpty()) {
+ sql.append(", older_than => TIMESTAMP
").append(lit(olderThan).expr().sql());
+ }
+ if (location != null && !location.isEmpty()) {
+ sql.append(", location => ").append(lit(location).expr().sql());
+ }
+ return sql.append(", dry_run => ").append(dryRun).append(")").toString();
+ }
+
+ static boolean parseDryRun(@Nullable String value) {
+ Preconditions.checkArgument(
+ value == null || "false".equals(value) || "true".equals(value),
+ "--dry-run must be true or false");
+ return "true".equals(value);
+ }
+
+ static void validateLocation(String tableLocation, String location) {
+ URI root = normalizeLocation(tableLocation);
+ URI requested = normalizeLocation(location);
+ String rootPath = root.getPath().replaceAll("/+$", "");
+ String childPath = requested.getPath().replaceAll("/+$", "");
+ boolean sameStorage =
+ Objects.equals(root.getScheme(), requested.getScheme())
+ && Objects.equals(root.getAuthority(), requested.getAuthority());
+ Preconditions.checkArgument(
+ sameStorage
+ && (childPath.equals(rootPath)
+ || childPath.startsWith(rootPath.endsWith("/") ? rootPath :
rootPath + "/")),
+ "location must be within the table's storage location: %s",
+ tableLocation);
+ }
+
+ private static List<String> buildArguments() {
+ return Arrays.asList(
+ "--catalog",
+ "{{catalog_name}}",
+ "--table",
+ "{{table_identifier}}",
+ "--older-than",
+ "{{older_than}}",
+ "--location",
+ "{{location}}",
+ "--dry-run",
+ "{{dry_run}}",
+ "--spark-conf",
+ "{{spark_conf}}");
+ }
+
+ private static Map<String, String> buildSparkConfigs() {
+ return IcebergSparkConfigUtils.buildTemplateSparkConfigs();
+ }
+
+ private static void printUsage() {
+ LOG.error(
+ "Usage: IcebergRemoveOrphanFilesJob --catalog <name> --table
<db.table> "
+ + "[--older-than 'yyyy-MM-dd HH:mm:ss'] [--location <path>] "
+ + "[--dry-run true|false] [--spark-conf <json>]");
+ }
+
+ private static URI normalizeLocation(String value) {
+ // Reject ambiguous encoded paths rather than allowing different
filesystem decoders to
+ // interpret the containment check and the subsequent listing differently.
+ Preconditions.checkArgument(
+ !value.isEmpty() && !value.contains("%") && !value.contains("\\"),
+ "Invalid scan location: %s",
+ value);
+ URI uri = URI.create(value);
+ Preconditions.checkArgument(
+ uri.getQuery() == null
+ && uri.getFragment() == null
+ && uri.getPath() != null
+ && uri.getPath().startsWith("/"),
+ "Scan location must be an absolute path without query or fragment: %s",
+ value);
+ if (uri.getScheme() == null || "file".equals(uri.getScheme())) {
+ return (uri.getScheme() == null ? Paths.get(value) :
Paths.get(uri)).normalize().toUri();
+ }
+ return uri.normalize();
+ }
+
+ private static void validateLocalLocation(String tableLocation, String
location)
+ throws IOException {
+ URI requested = normalizeLocation(location);
+ if (!"file".equals(requested.getScheme())) {
+ return;
+ }
+ Path lexicalRoot = Paths.get(normalizeLocation(tableLocation));
+ Path root = lexicalRoot.toRealPath();
+ Path scan = Paths.get(requested);
+ Preconditions.checkArgument(
+ scan.toRealPath().startsWith(root),
+ "Scan location resolves outside the table's storage location");
+ for (Path ancestor = scan;
+ ancestor != null && ancestor.startsWith(lexicalRoot);
+ ancestor = ancestor.getParent()) {
+ Preconditions.checkArgument(
+ !Files.isSymbolicLink(ancestor), "Symlinks are not allowed in the
scan location");
+ }
+ // Do not follow symlinks during validation. Iceberg must never list
another table through one.
+ try (Stream<Path> paths = Files.walk(scan)) {
+ Preconditions.checkArgument(
+ paths.noneMatch(Files::isSymbolicLink), "Symlinks are not allowed in
the scan location");
+ }
+ }
+
+ private static void validateRemoteLocation(
+ SparkSession spark, String tableLocation, String location) throws
IOException {
+ if ("file".equals(normalizeLocation(location).getScheme())) {
+ return;
+ }
+ RemoteLocationValidator.validate(
+ spark.sparkContext().hadoopConfiguration(), tableLocation, location);
+ }
+
+ private static String requireOption(Map<String, String> options, String key)
{
+ String value = options.get(key);
+ Preconditions.checkArgument(value != null && !value.trim().isEmpty(),
"--%s is required", key);
+ return value;
+ }
+
+ // Defer verification of Spark-specific exception handlers until job
execution. The server
+ // loads this job's template without Spark or Iceberg on its classpath.
+ private static final class Runner {
+ private static int run(String[] args) {
+ Map<String, String> options = IcebergJobUtils.parseArguments(args);
+ SparkSession.Builder builder =
+ SparkSession.builder().appName("Gravitino Built-in Iceberg Remove
Orphan Files");
+ try {
+ requireOption(options, "catalog");
+ requireOption(options, "table");
+ parseDryRun(options.get("dry-run"));
+
IcebergJobUtils.parseCustomSparkConfigs(options.get("spark-conf")).forEach(builder::config);
+ } catch (IllegalArgumentException e) {
+ LOG.error("Invalid remove orphan files job arguments: {}",
e.getMessage());
+ printUsage();
+ return 1;
+ }
+
+ SparkSession spark = null;
+ try {
+ spark = builder.getOrCreate();
+ IcebergJobUtils.requireIcebergSparkRuntime();
+ execute(spark, options);
+ return 0;
+ } catch (IOException | AnalysisException | RuntimeException e) {
+ LOG.error("Error executing remove orphan files job", e);
+ return 1;
+ } finally {
+ if (spark != null) {
+ spark.stop();
+ }
+ }
+ }
+ }
+}
diff --git
a/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/RemoteLocationValidator.java
b/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/RemoteLocationValidator.java
new file mode 100644
index 0000000000..6d1a93a722
--- /dev/null
+++
b/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/RemoteLocationValidator.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.jobs.iceberg;
+
+import com.google.common.base.Preconditions;
+import java.io.IOException;
+import java.util.ArrayDeque;
+import java.util.Deque;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.FileStatus;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.fs.RemoteIterator;
+
+/** Validates scan paths on Hadoop filesystems that can resolve symbolic
links. */
+final class RemoteLocationValidator {
+ private RemoteLocationValidator() {}
+
+ static void validate(Configuration conf, String tableLocation, String
location)
+ throws IOException {
+ Path scan = new Path(location);
+ FileSystem fs = scan.getFileSystem(conf);
+ validate(fs, new Path(tableLocation), scan);
+ }
+
+ static void validate(FileSystem fs, Path tableLocation, Path scan) throws
IOException {
+ // Object stores do not support symbolic links. The URI containment check
is sufficient there.
+ if (!fs.supportsSymlinks()) {
+ return;
+ }
+ Path lexicalRoot = new Path(tableLocation.toUri().normalize());
+ scan = new Path(scan.toUri().normalize());
+ IcebergRemoveOrphanFilesJob.validateLocation(lexicalRoot.toString(),
scan.toString());
+ Path root = fs.resolvePath(lexicalRoot);
+ IcebergRemoveOrphanFilesJob.validateLocation(root.toString(),
fs.resolvePath(scan).toString());
+ // Inspect the table root itself, but not warehouse parents outside this
table. Use the
+ // lexical boundary because a warehouse parent symlink can change the
resolved root.
+ for (Path ancestor = scan; ancestor != null; ancestor =
ancestor.getParent()) {
+ Preconditions.checkArgument(
+ !fs.getFileLinkStatus(ancestor).isSymlink(),
+ "Symlinks are not allowed in the scan location");
+ if (ancestor.equals(lexicalRoot)) {
+ break;
+ }
+ }
+ Deque<Path> pending = new ArrayDeque<>();
+ pending.add(scan);
+ while (!pending.isEmpty()) {
+ Path path = pending.removeFirst();
+ FileStatus status = fs.getFileLinkStatus(path);
+ Preconditions.checkArgument(
+ !status.isSymlink(), "Symlinks are not allowed in the scan
location");
+ if (status.isDirectory()) {
+ RemoteIterator<FileStatus> children = fs.listStatusIterator(path);
+ while (children.hasNext()) {
+ FileStatus child = children.next();
+ Preconditions.checkArgument(
+ !child.isSymlink(), "Symlinks are not allowed in the scan
location");
+ if (child.isDirectory()) {
+ pending.addLast(child.getPath());
+ }
+ }
+ }
+ }
+ }
+}
diff --git
a/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRemoveOrphanFilesJob.java
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRemoveOrphanFilesJob.java
new file mode 100644
index 0000000000..d51f2a6838
--- /dev/null
+++
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRemoveOrphanFilesJob.java
@@ -0,0 +1,133 @@
+/*
+ * 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.jobs.iceberg;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.io.File;
+import java.net.URL;
+import java.net.URLClassLoader;
+import java.util.ArrayList;
+import java.util.List;
+import org.apache.gravitino.maintenance.jobs.BuiltInJobTemplateProvider;
+import org.junit.jupiter.api.Test;
+
+/** Tests orphan cleanup templates, SQL arguments, and location boundaries. */
+public class TestIcebergRemoveOrphanFilesJob {
+ /** Verifies template registration. */
+ @Test
+ public void testTemplateRegistration() {
+ assertTrue(new
IcebergRemoveOrphanFilesJob().jobTemplate().environments().isEmpty());
+ assertTrue(
+ new BuiltInJobTemplateProvider()
+ .jobTemplates().stream()
+ .anyMatch(t ->
t.name().equals("builtin-iceberg-remove-orphan-files")));
+ assertTrue(new
IcebergRemoveOrphanFilesJob().jobTemplate().arguments().contains("{{dry_run}}"));
+ }
+
+ /** Verifies server-side template discovery without Spark, Iceberg, or
Hadoop libraries. */
+ @Test
+ public void testTemplateLoadsWithoutSparkRuntime() throws Exception {
+ List<URL> urls = new ArrayList<>();
+ for (String entry :
System.getProperty("java.class.path").split(File.pathSeparator)) {
+ File file = new File(entry);
+ String name = file.getName();
+ if (!name.startsWith("spark-")
+ && !name.startsWith("iceberg-")
+ && !name.startsWith("hadoop-")) {
+ urls.add(file.toURI().toURL());
+ }
+ }
+ try (URLClassLoader loader =
+ new URLClassLoader(
+ urls.toArray(new URL[0]),
ClassLoader.getSystemClassLoader().getParent())) {
+ assertThrows(
+ ClassNotFoundException.class,
+ () -> loader.loadClass("org.apache.spark.sql.SparkSession"));
+ Class<?> providerClass =
loader.loadClass(BuiltInJobTemplateProvider.class.getName());
+ Object provider = providerClass.getDeclaredConstructor().newInstance();
+ List<?> templates = (List<?>)
providerClass.getMethod("jobTemplates").invoke(provider);
+ List<String> names = new ArrayList<>();
+ for (Object template : templates) {
+ names.add((String)
template.getClass().getMethod("name").invoke(template));
+ }
+ assertTrue(names.contains("builtin-iceberg-remove-orphan-files"));
+ }
+ }
+
+ /** Verifies defaults and explicit dry run. */
+ @Test
+ public void testDefaultsAndExplicitDryRun() {
+ assertEquals(
+ "CALL `cat`.system.remove_orphan_files(table => 'db.tbl', dry_run =>
false)",
+ IcebergRemoveOrphanFilesJob.buildProcedureCall("cat", "db.tbl", null,
null, false));
+ assertEquals(
+ "CALL `cat`.system.remove_orphan_files(table => 'db.tbl', older_than
=> TIMESTAMP '2024-01-01 00:00:00', location => '/warehouse/tbl/data', dry_run
=> true)",
+ IcebergRemoveOrphanFilesJob.buildProcedureCall(
+ "cat", "db.tbl", "2024-01-01 00:00:00", "/warehouse/tbl/data",
true));
+ assertFalse(IcebergRemoveOrphanFilesJob.parseDryRun(null));
+ assertFalse(IcebergRemoveOrphanFilesJob.parseDryRun("false"));
+ assertTrue(IcebergRemoveOrphanFilesJob.parseDryRun("true"));
+ assertThrows(
+ IllegalArgumentException.class, () ->
IcebergRemoveOrphanFilesJob.parseDryRun("yes"));
+ }
+
+ /** Verifies sql escaping. */
+ @Test
+ public void testSqlEscaping() {
+ String sql =
+ IcebergRemoveOrphanFilesJob.buildProcedureCall("ca`t", "db.ta'ble",
null, "/a/b'c", true);
+ assertTrue(sql.contains("`ca``t`"));
+ assertTrue(sql.contains("db.ta\\'ble"));
+ assertTrue(sql.contains("/a/b\\'c"));
+ }
+
+ /** Verifies location boundaries. */
+ @Test
+ public void testLocationBoundaries() {
+ assertDoesNotThrow(
+ () ->
+ IcebergRemoveOrphanFilesJob.validateLocation(
+ "s3://bucket/db/table", "s3://bucket/db/table/data/"));
+ assertDoesNotThrow(
+ () ->
+ IcebergRemoveOrphanFilesJob.validateLocation(
+ "s3://bucket/db/table/", "s3://bucket/db/table"));
+ for (String location :
+ new String[] {
+ "s3://bucket/db/table2",
+ "s3://other/db/table",
+ "s3://bucket/db",
+ "s3://bucket/db/table/../other",
+ "s3://bucket/db/table/%2e%2e/other",
+ "s3://bucket/db/table?x=1",
+ "relative/path",
+ "s3a://bucket/db/table"
+ }) {
+ assertThrows(
+ IllegalArgumentException.class,
+ () ->
IcebergRemoveOrphanFilesJob.validateLocation("s3://bucket/db/table", location),
+ location);
+ }
+ }
+}
diff --git
a/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRemoveOrphanFilesJobMain.java
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRemoveOrphanFilesJobMain.java
new file mode 100644
index 0000000000..5d6f293944
--- /dev/null
+++
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRemoveOrphanFilesJobMain.java
@@ -0,0 +1,144 @@
+/*
+ * 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.jobs.iceberg;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import java.io.File;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.attribute.FileTime;
+import java.time.Instant;
+import java.time.temporal.ChronoUnit;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.hadoop.HadoopTables;
+import org.apache.iceberg.types.Types;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+/** Tests orphan cleanup CLI behavior and process exit status. */
+public class TestIcebergRemoveOrphanFilesJobMain {
+ @TempDir Path tempDir;
+
+ /** Verifies template arguments run preview and cleanup. */
+ @Test
+ public void testTemplateArgumentsRunPreviewAndCleanup() throws Exception {
+ Path table = tempDir.resolve("db/cli");
+ new HadoopTables(new Configuration())
+ .create(
+ new Schema(Types.NestedField.optional(1, "id",
Types.IntegerType.get())),
+ table.toString());
+ Path data = Files.createDirectories(table.resolve("data"));
+ Path orphan = Files.write(data.resolve("orphan"), new byte[] {1});
+ Files.setLastModifiedTime(orphan, FileTime.from(Instant.now().minus(5,
ChronoUnit.DAYS)));
+ Map<String, String> config = new HashMap<>();
+ config.put("spark.master", "local[2]");
+ config.put("spark.ui.enabled", "false");
+ config.put("spark.sql.shuffle.partitions", "2");
+ config.put(
+ "spark.sql.extensions",
+ "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions");
+ config.put("spark.sql.catalog.cli",
"org.apache.iceberg.spark.SparkCatalog");
+ config.put("spark.sql.catalog.cli.type", "hadoop");
+ config.put("spark.sql.catalog.cli.warehouse", tempDir.toString());
+ String json = new ObjectMapper().writeValueAsString(config);
+ Map<String, String> jobConf = new HashMap<>();
+ jobConf.put("catalog_name", "cli");
+ jobConf.put("table_identifier", "db.cli");
+ jobConf.put("older_than", "");
+ jobConf.put("location", "");
+ jobConf.put("spark_conf", json);
+ jobConf.put("dry_run", "true");
+ assertEquals(0,
IcebergRemoveOrphanFilesJob.run(templateArguments(jobConf)));
+ assertTrue(Files.exists(orphan));
+ jobConf.put("dry_run", "false");
+ assertEquals(0,
IcebergRemoveOrphanFilesJob.run(templateArguments(jobConf)));
+ assertFalse(Files.exists(orphan));
+ }
+
+ /** Verifies cli failure exit and usage. */
+ @Test
+ public void testCliFailureExitAndUsage() throws Exception {
+ String output = runFailure(false, new String[] {"--catalog", "cli"});
+ assertTrue(output.contains("Usage: IcebergRemoveOrphanFilesJob"), output);
+ assertTrue(output.contains("--table is required"), output);
+ }
+
+ /** Verifies missing runtime has actionable error. */
+ @Test
+ public void testMissingRuntimeHasActionableError() throws Exception {
+ String output =
+ runFailure(
+ true,
+ new String[] {
+ "--catalog",
+ "cli",
+ "--table",
+ "db.cli",
+ "--spark-conf",
+ "{\"spark.master\":\"local[1]\",\"spark.ui.enabled\":\"false\"}"
+ });
+ assertTrue(output.contains("iceberg-spark-runtime"), output);
+ }
+
+ private String[] templateArguments(Map<String, String> jobConf) {
+ return new IcebergRemoveOrphanFilesJob()
+ .jobTemplate().arguments().stream()
+ .map(
+ arg -> arg.startsWith("{{") ? jobConf.get(arg.substring(2,
arg.length() - 2)) : arg)
+ .toArray(String[]::new);
+ }
+
+ private String runFailure(boolean omitIcebergRuntime, String[] args) throws
Exception {
+ String classpath =
+
Arrays.stream(System.getProperty("java.class.path").split(File.pathSeparator))
+ .filter(entry -> !omitIcebergRuntime ||
!entry.contains("iceberg-spark-runtime"))
+ .collect(Collectors.joining(File.pathSeparator));
+ List<String> command = new ArrayList<>();
+ command.add(new File(System.getProperty("java.home"),
"bin/java").toString());
+ command.add("--add-opens=java.base/sun.nio.ch=ALL-UNNAMED");
+ command.add("-cp");
+ command.add(classpath);
+ command.add(IcebergRemoveOrphanFilesJob.class.getName());
+ command.addAll(Arrays.asList(args));
+ Path log = Files.createTempFile(tempDir, "cli-", ".log");
+ Process process =
+ new
ProcessBuilder(command).redirectErrorStream(true).redirectOutput(log.toFile()).start();
+ try {
+ assertTrue(process.waitFor(45, TimeUnit.SECONDS), "CLI did not exit");
+ String output = new String(Files.readAllBytes(log),
StandardCharsets.UTF_8);
+ assertEquals(1, process.exitValue(), output);
+ return output;
+ } finally {
+ process.destroyForcibly();
+ }
+ }
+}
diff --git
a/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRemoveOrphanFilesJobWithSpark.java
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRemoveOrphanFilesJobWithSpark.java
new file mode 100644
index 0000000000..1d3a8d821e
--- /dev/null
+++
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRemoveOrphanFilesJobWithSpark.java
@@ -0,0 +1,192 @@
+/*
+ * 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.jobs.iceberg;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.attribute.FileTime;
+import java.time.Instant;
+import java.time.temporal.ChronoUnit;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.stream.Collectors;
+import java.util.stream.Stream;
+import org.apache.spark.sql.SparkSession;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+/** Tests orphan cleanup against real Iceberg tables in local Spark. */
+public class TestIcebergRemoveOrphanFilesJobWithSpark {
+ @TempDir static Path tempDir;
+ private static SparkSession spark;
+
+ @BeforeAll
+ static void setUp() {
+ spark =
+ SparkSession.builder()
+ .master("local[2]")
+ .appName("TestRemoveOrphanFiles")
+ .config("spark.ui.enabled", "false")
+ .config("spark.sql.shuffle.partitions", "2")
+ .config(
+ "spark.sql.extensions",
+
"org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
+ .config("spark.sql.catalog.test_catalog",
"org.apache.iceberg.spark.SparkCatalog")
+ .config("spark.sql.catalog.test_catalog.type", "hadoop")
+ .config("spark.sql.catalog.test_catalog.warehouse",
tempDir.toString())
+ .getOrCreate();
+ spark.sql("CREATE NAMESPACE test_catalog.db");
+ }
+
+ @AfterAll
+ static void tearDown() {
+ if (spark != null) {
+ spark.stop();
+ }
+ }
+
+ /** Verifies dry run and deletion preserve live and recent files. */
+ @Test
+ public void testDryRunAndDeletionPreserveLiveAndRecentFiles() throws
Exception {
+ spark.sql("CREATE TABLE test_catalog.db.cleanup (id INT) USING iceberg");
+ spark.sql("INSERT INTO test_catalog.db.cleanup VALUES (1), (2)");
+ Path data = tempDir.resolve("db/cleanup/data");
+ // Age referenced files too, so survival depends on metadata references,
not the cutoff.
+ try (Stream<Path> files = Files.list(data)) {
+ for (Path file : files.collect(Collectors.toList())) {
+ Files.setLastModifiedTime(file, FileTime.from(Instant.now().minus(5,
ChronoUnit.DAYS)));
+ }
+ }
+ Path old = Files.write(data.resolve("old-orphan.parquet"), new byte[] {1});
+ Files.setLastModifiedTime(old, FileTime.from(Instant.now().minus(5,
ChronoUnit.DAYS)));
+ Path recent = Files.write(data.resolve("recent-orphan.parquet"), new
byte[] {1});
+ Map<String, String> args = args("cleanup");
+ args.put("dry-run", "true");
+ assertEquals(1, IcebergRemoveOrphanFilesJob.execute(spark, args));
+ assertTrue(Files.exists(old));
+ args.put("dry-run", "false");
+ assertEquals(1, IcebergRemoveOrphanFilesJob.execute(spark, args));
+ assertFalse(Files.exists(old));
+ assertTrue(Files.exists(recent));
+ assertEquals(2, spark.table("test_catalog.db.cleanup").count());
+ assertEquals(0, IcebergRemoveOrphanFilesJob.execute(spark, args));
+ }
+
+ /** Verifies custom location and outside rejection. */
+ @Test
+ public void testCustomLocationAndOutsideRejection() throws Exception {
+ spark.sql("CREATE TABLE test_catalog.db.scoped (id INT) USING iceberg");
+ Path root = tempDir.resolve("db/scoped");
+ Path sub = Files.createDirectories(root.resolve("staged"));
+ Path inside = Files.write(sub.resolve("old"), new byte[] {1});
+ Path outside = Files.write(root.resolve("old"), new byte[] {1});
+ FileTime old = FileTime.from(Instant.now().minus(5, ChronoUnit.DAYS));
+ Files.setLastModifiedTime(inside, old);
+ Files.setLastModifiedTime(outside, old);
+ Map<String, String> args = args("scoped");
+ args.put("location", sub.toString());
+ assertEquals(1, IcebergRemoveOrphanFilesJob.execute(spark, args));
+ assertFalse(Files.exists(inside));
+ assertTrue(Files.exists(outside));
+ args.put("location", tempDir.toString());
+ assertThrows(
+ IllegalArgumentException.class, () ->
IcebergRemoveOrphanFilesJob.execute(spark, args));
+ assertTrue(Files.exists(outside));
+ Path link = root.resolve("escape");
+ Files.createSymbolicLink(link, tempDir);
+ args.put("location", link.toString());
+ assertThrows(
+ IllegalArgumentException.class, () ->
IcebergRemoveOrphanFilesJob.execute(spark, args));
+ }
+
+ /** Verifies invalid inputs fail before deletion. */
+ @Test
+ public void testInvalidInputsFailBeforeDeletion() throws Exception {
+ spark.sql("CREATE TABLE test_catalog.db.invalid (id INT) USING iceberg");
+ Map<String, String> options = args("invalid");
+ options.put("dry-run", "yes");
+ assertThrows(
+ IllegalArgumentException.class, () ->
IcebergRemoveOrphanFilesJob.execute(spark, options));
+ options.remove("dry-run");
+ options.put("location", tempDir.resolve("db/invalid/missing").toString());
+ assertThrows(IOException.class, () ->
IcebergRemoveOrphanFilesJob.execute(spark, options));
+ options.remove("location");
+ options.put("older-than", "2999-01-01 00:00:00");
+ assertThrows(
+ IllegalArgumentException.class, () ->
IcebergRemoveOrphanFilesJob.execute(spark, options));
+ options.remove("older-than");
+ spark.conf().set("spark.sql.parser.escapedStringLiterals", "true");
+ try {
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> IcebergRemoveOrphanFilesJob.execute(spark, options));
+ } finally {
+ spark.conf().set("spark.sql.parser.escapedStringLiterals", "false");
+ }
+ assertEquals(1, IcebergRemoveOrphanFilesJob.run(new String[] {}));
+ assertEquals(1, IcebergRemoveOrphanFilesJob.run(new String[] {"--catalog",
"test_catalog"}));
+ assertEquals(
+ 1,
+ IcebergRemoveOrphanFilesJob.run(
+ new String[] {
+ "--catalog", "test_catalog", "--table", "db.invalid",
"--dry-run", "yes"
+ }));
+ assertEquals(
+ 1,
+ IcebergRemoveOrphanFilesJob.run(
+ new String[] {
+ "--catalog", "test_catalog", "--table", "db.invalid",
"--spark-conf", "{bad}"
+ }));
+ }
+
+ /** Verifies explicit cutoff and quoted table name. */
+ @Test
+ public void testExplicitCutoffAndQuotedTableName() throws Exception {
+ Path location = tempDir.resolve("db/quo'te");
+ spark.sql("CREATE TABLE test_catalog.db.`quo'te` (id INT) USING iceberg");
+ Path old = Files.write(location.resolve("old"), new byte[] {1});
+ Path newer = Files.write(location.resolve("newer"), new byte[] {1});
+ Files.setLastModifiedTime(old,
FileTime.from(Instant.parse("2023-01-01T00:00:00Z")));
+ Files.setLastModifiedTime(newer,
FileTime.from(Instant.parse("2025-01-01T00:00:00Z")));
+ Map<String, String> options = args("`quo'te`");
+ options.put("older-than", "2024-01-01 00:00:00");
+ options.put("dry-run", "true");
+ assertEquals(1, IcebergRemoveOrphanFilesJob.execute(spark, options));
+ assertTrue(Files.exists(old));
+ options.put("dry-run", "false");
+ assertEquals(1, IcebergRemoveOrphanFilesJob.execute(spark, options));
+ assertFalse(Files.exists(old));
+ assertTrue(Files.exists(newer));
+ }
+
+ private static Map<String, String> args(String table) {
+ Map<String, String> args = new HashMap<>();
+ args.put("catalog", "test_catalog");
+ args.put("table", "db." + table);
+ return args;
+ }
+}
diff --git
a/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestRemoteLocationValidator.java
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestRemoteLocationValidator.java
new file mode 100644
index 0000000000..af342db6ab
--- /dev/null
+++
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestRemoteLocationValidator.java
@@ -0,0 +1,197 @@
+/*
+ * 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.jobs.iceberg;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import org.apache.hadoop.fs.FileStatus;
+import org.apache.hadoop.fs.FilterFileSystem;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.fs.RemoteIterator;
+import org.junit.jupiter.api.Test;
+
+/** Tests table-scoped validation on symlink-capable filesystems. */
+public class TestRemoteLocationValidator {
+ private final StubFileSystem fs = new StubFileSystem();
+ private final Path root = new Path("hdfs://host/table");
+ private final Path scan = new Path(root, "data");
+
+ /** Verifies object store does not require symlink inspection. */
+ @Test
+ public void testObjectStoreDoesNotRequireSymlinkInspection() throws
Exception {
+ fs.symlinksSupported = false;
+ RemoteLocationValidator.validate(fs, root, scan);
+ assertEquals(0, fs.resolveCalls);
+ }
+
+ /** Verifies resolved path outside table is rejected. */
+ @Test
+ public void testResolvedPathOutsideTableIsRejected() {
+ fs.resolvedPaths.put(scan, new Path("hdfs://host/other/data"));
+ assertThrows(
+ IllegalArgumentException.class, () ->
RemoteLocationValidator.validate(fs, root, scan));
+ }
+
+ /** Verifies inspection failure is not ignored. */
+ @Test
+ public void testInspectionFailureIsNotIgnored() {
+ fs.failResolution = true;
+ assertThrows(IOException.class, () -> RemoteLocationValidator.validate(fs,
root, scan));
+ }
+
+ /** Verifies symbolic link is rejected. */
+ @Test
+ public void testSymbolicLinkIsRejected() {
+ fs.statuses.put(scan, link(scan));
+ assertThrows(
+ IllegalArgumentException.class, () ->
RemoteLocationValidator.validate(fs, root, scan));
+ }
+
+ /** Verifies descendant symbolic link is rejected. */
+ @Test
+ public void testDescendantSymbolicLinkIsRejected() {
+ fs.listings.put(scan, new FileStatus[] {link(new Path(scan, "link"))});
+ assertThrows(
+ IllegalArgumentException.class, () ->
RemoteLocationValidator.validate(fs, root, scan));
+ }
+
+ /** Verifies directories are inspected recursively. */
+ @Test
+ public void testDirectoriesAreInspectedRecursively() throws Exception {
+ Path sub = new Path(scan, "sub");
+ fs.listings.put(scan, new FileStatus[] {directory(sub)});
+ RemoteLocationValidator.validate(fs, root, scan);
+ assertTrue(fs.listedPaths.contains(sub));
+ }
+
+ /** Verifies warehouse parent symlink is outside validation scope. */
+ @Test
+ public void testWarehouseParentSymlinkIsOutsideValidationScope() throws
Exception {
+ Path table = new Path("hdfs://host/warehouse/table");
+ Path data = new Path(table, "data");
+ fs.resolvedPaths.put(table, new Path("hdfs://host/storage/table"));
+ fs.resolvedPaths.put(data, new Path("hdfs://host/storage/table/data"));
+ fs.statuses.put(table.getParent(), link(table.getParent()));
+ RemoteLocationValidator.validate(fs, table, data);
+ assertTrue(fs.listedPaths.contains(data));
+ }
+
+ /** Verifies unreadable parent is outside validation scope. */
+ @Test
+ public void testUnreadableParentIsOutsideValidationScope() throws Exception {
+ fs.unreadablePath = root.getParent();
+ RemoteLocationValidator.validate(fs, root, scan);
+ assertTrue(fs.listedPaths.contains(scan));
+ }
+
+ /** Verifies table root and intermediate symlinks are rejected. */
+ @Test
+ public void testTableRootAndIntermediateSymlinksAreRejected() {
+ fs.statuses.put(root, link(root));
+ assertThrows(
+ IllegalArgumentException.class, () ->
RemoteLocationValidator.validate(fs, root, scan));
+ fs.statuses.clear();
+ fs.statuses.put(scan, link(scan));
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> RemoteLocationValidator.validate(fs, root, new Path(scan,
"nested")));
+ }
+
+ /** Verifies normalized paths stop at the table boundary, including a
whole-table scan. */
+ @Test
+ public void testNormalizedPathsAndWholeTableScan() throws Exception {
+ fs.unreadablePath = root.getParent();
+ RemoteLocationValidator.validate(fs, new Path(root, "."), new Path(root,
"staged/../data"));
+ assertTrue(fs.listedPaths.contains(scan));
+ RemoteLocationValidator.validate(fs, root, root);
+ assertTrue(fs.listedPaths.contains(root));
+ }
+
+ private static FileStatus directory(Path path) {
+ return new FileStatus(0, true, 1, 0, 0, path);
+ }
+
+ private static FileStatus link(Path path) {
+ FileStatus status = directory(path);
+ status.setSymlink(new Path("hdfs://host/other"));
+ return status;
+ }
+
+ /**
+ * In-memory filesystem responses for deterministic validation tests without
external services.
+ */
+ private static class StubFileSystem extends FilterFileSystem {
+ private final Map<Path, Path> resolvedPaths = new HashMap<>();
+ private final Map<Path, FileStatus> statuses = new HashMap<>();
+ private final Map<Path, FileStatus[]> listings = new HashMap<>();
+ private final List<Path> listedPaths = new ArrayList<>();
+ private boolean symlinksSupported = true;
+ private boolean failResolution;
+ private Path unreadablePath;
+ private int resolveCalls;
+
+ @Override
+ public boolean supportsSymlinks() {
+ return symlinksSupported;
+ }
+
+ @Override
+ public Path resolvePath(Path path) throws IOException {
+ resolveCalls++;
+ if (failResolution) {
+ throw new IOException("Permission denied");
+ }
+ return resolvedPaths.getOrDefault(path, path);
+ }
+
+ @Override
+ public FileStatus getFileLinkStatus(Path path) throws IOException {
+ if (path.equals(unreadablePath)) {
+ throw new IOException("Permission denied: " + path);
+ }
+ return statuses.getOrDefault(path, directory(path));
+ }
+
+ @Override
+ public RemoteIterator<FileStatus> listStatusIterator(Path path) {
+ listedPaths.add(path);
+ FileStatus[] children = listings.getOrDefault(path, new FileStatus[0]);
+ return new RemoteIterator<FileStatus>() {
+ private int index;
+
+ @Override
+ public boolean hasNext() {
+ return index < children.length;
+ }
+
+ @Override
+ public FileStatus next() {
+ return children[index++];
+ }
+ };
+ }
+ }
+}