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 621c448d77 [#11196] feat(maintenance): Implement Iceberg rewrite
manifests job (#13329)
621c448d77 is described below
commit 621c448d7750293e6ad958167a5fc996c1f3d3c2
Author: Akshay Thorat <[email protected]>
AuthorDate: Sun Sep 27 23:19:06 2026 -0700
[#11196] feat(maintenance): Implement Iceberg rewrite manifests job (#13329)
### What changes were proposed in this pull request?
Implements the first phase of the merged design in #12942: the
`builtin-iceberg-rewrite-manifests` Spark job, registration, validated
arguments, spec selection, and operator documentation. Manifest
statistics and the Policy -> Strategy -> Adapter integration for
automatic triggering remain separate follow-ups in the design, as
requested in review.
Continues #12937 on current main and preserves the original commits by
@ibrahimErbilen from #11216, with co-author attribution on follow-up
work.
Adds Spark-backed tests invoking the job entry point through resolved
template arguments. They verify default and explicit partition-spec
selection, unchanged table records and other-spec manifests, repeat-run
no-op behavior, and cleanup without metadata changes after an
unknown-spec failure. The entry point checks the Iceberg runtime, closes
Spark on success or failure, prints user-visible results to stdout for
output.log, and reports concise errors to stderr with a nonzero exit.
Strict argument parsing is shared through IcebergJobUtils while
preserving existing callers. The overview and configuration guides
document direct submission and the remaining policy integration.
### Why are the changes needed?
Manifest rewriting maintains manifest layout independently of data-file
compaction. The previous implementation lacked Spark-backed verification
of partition-spec selection and failure cleanup.
Fix: #11196
### Does this PR introduce _any_ user-facing change?
Adds `builtin-iceberg-rewrite-manifests` with required `catalog_name`
and `table_identifier`, plus optional `use_caching`, `spec_id`, and
`spark_conf`. Existing jobConf/CLI naming conventions are retained.
Invalid, duplicate, or missing argument values fail before session
creation; omitted optional values retain Iceberg defaults. `spec_id`
selects an existing spec, without migrating files between specs.
### How was this patch tested?
- 183 jobs-module tests passed, including 33 manifest unit tests and 3
Spark-backed tests on Spark 3.5 / Iceberg 1.11.0.
- `./gradlew :maintenance:jobs:test :maintenance:jobs:spotlessCheck
:maintenance:jobs:javadoc :docs:build rat -PskipITs` passed with JDK 17.
- Process-level checks verified exit code 1, concise stderr, and no
stdout or stack trace for missing arguments and malformed Spark
configuration.
- `git diff --check` passed.
Spark tests use a local Hadoop catalog; remote catalog deployment and
the server REST submission path were not exercised end to end.
---------
Co-authored-by: İbrahim ERBİLEN <[email protected]>
---
.../optimizer-cli-reference.md | 97 ++++-
.../optimizer-configuration.md | 23 ++
docs/table-maintenance-service/optimizer.md | 17 +
.../jobs/BuiltInJobTemplateProvider.java | 2 +
.../maintenance/jobs/iceberg/IcebergJobUtils.java | 66 +++
.../jobs/iceberg/IcebergRewriteManifestsJob.java | 266 ++++++++++++
.../jobs/TestBuiltInJobTemplateProvider.java | 12 +
.../jobs/iceberg/TestIcebergJobUtils.java | 45 +++
.../iceberg/TestIcebergRewriteManifestsJob.java | 446 +++++++++++++++++++++
.../TestIcebergRewriteManifestsJobWithSpark.java | 258 ++++++++++++
10 files changed, 1231 insertions(+), 1 deletion(-)
diff --git a/docs/table-maintenance-service/optimizer-cli-reference.md
b/docs/table-maintenance-service/optimizer-cli-reference.md
index cba7e1cad0..5fe7e04a7e 100644
--- a/docs/table-maintenance-service/optimizer-cli-reference.md
+++ b/docs/table-maintenance-service/optimizer-cli-reference.md
@@ -237,7 +237,7 @@ EvaluationResult{scopeType=TABLE,
identifier=rest_catalog.db.t1, partitionPath=<
## Built-in Job Templates
-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.
+Five 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, consolidates
manifests, and removes old orphan files.
| Job template | What it does
|
|---------------------------------------|-------------------------------------------|
@@ -245,6 +245,7 @@ Four job templates ship with the service, and they are
complementary rather than
| `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 |
+| `builtin-iceberg-rewrite-manifests` | Consolidates small manifest files
|
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.
@@ -360,6 +361,100 @@ Expire Snapshots Results:
Deleted manifest lists: 3
```
+## Rewrite Manifests
+
+`builtin-iceberg-rewrite-manifests` consolidates a table's manifest files.
Frequent writes can leave many small manifests. Scan planning uses
manifest-list summaries to prune them, then reads the remaining manifests.
Rewriting consolidates and clusters entries for one existing partition spec
without repartitioning data files.
+
+This complements `builtin-iceberg-rewrite-data-files`: that job improves the
data file layout, this one improves the metadata that points at it.
+
+The job calls Iceberg's `rewrite_manifests` stored procedure through Spark SQL.
+
+| Property | Value
|
+| ------------- |
------------------------------------------------------------------------------ |
+| Name | `builtin-iceberg-rewrite-manifests`
|
+| Type | Spark
|
+| Version | `v1`
|
+| Main class |
`org.apache.gravitino.maintenance.jobs.iceberg.IcebergRewriteManifestsJob` |
+
+### Parameters
+
+`catalog_name` and `table_identifier` are required. For this job, omitted or
blank optional argument values use the defaults below; unresolved optional
template placeholders are also treated as absent.
+
+| Key | Description
| Default |
+| -------------------- |
------------------------------------------------------------------------ |
----------------------------------- |
+| `catalog_name` | Iceberg catalog name as registered in Spark
| Required |
+| `table_identifier` | Fully qualified table name, such as `db.sample`
| Required |
+| `use_caching` | Caches table metadata in Spark while rewriting;
`true` or `false` | Installed Iceberg version's default |
+| `spec_id` | Existing partition spec whose manifests to rewrite;
non-negative integer | The table's current spec |
+| `spark_conf` | JSON map of Spark configuration
| None |
+
+Leave `use_caching` unset to use the installed Iceberg version's default
(`false` in Iceberg 1.11.0). Set it explicitly when consistent behavior across
versions is required.
+
+### How to find `spec_id`
+
+Omit `spec_id` for routine maintenance of the current partition spec. Do not
guess IDs such as `0` or `1`. Read `default-spec-id` and the `partition-specs`
array from the table's current Iceberg metadata JSON. The array maps each
`spec-id` to its partition fields and transforms; `default-spec-id` identifies
the current spec. These are Iceberg table metadata fields, not Gravitino
catalog properties.
+
+To discover which specs have manifests in the current snapshot, run:
+
+```sql
+SELECT DISTINCT partition_spec_id
+FROM rest_catalog.db.t1.manifests;
+```
+
+This query lists represented specs, not which one is current. A defined spec
with no manifests may be absent; use the metadata JSON to identify the default
and interpret the transforms.
+
+For example, suppose metadata shows spec `0` uses `day(event_time)` and the
current spec `1` uses `hour(event_time)`. Omitting `spec_id` rewrites eligible
manifests for spec `1`. Passing `"spec_id": "0"` consolidates the old day-spec
manifests. Neither run converts day-partitioned data files to hour partitioning
or rewrites manifests belonging to the other spec.
+
+`spec_id` selects the existing spec whose manifests are eligible for
rewriting, and replacement manifests use that same spec. Iceberg validates that
the ID exists. This follows the [Iceberg 1.11.0 action
implementation](https://github.com/apache/iceberg/blob/apache-iceberg-1.11.0/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteManifestsSparkAction.java),
where `findMatchingManifests` compares each manifest's `partitionSpecId()` to
the selected spec. It is not a par [...]
+
+### Submitting the Job
+
+```bash
+curl -X POST -H "Accept: application/vnd.gravitino.v1+json" \
+ -H "Content-Type: application/json" \
+ -d '{
+ "jobTemplateName": "builtin-iceberg-rewrite-manifests",
+ "jobConf": {
+ "catalog_name": "rest_catalog",
+ "table_identifier": "db.t1",
+ "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": ""
+ }
+ }' \
+ http://localhost:8090/api/metalakes/test/jobs/runs
+```
+
+The request above uses Iceberg defaults. Adding `"use_caching": "false"` and
`"spec_id": "2"` to `jobConf` produces:
+
+```sql
+CALL `rest_catalog`.system.rewrite_manifests(
+ table => 'db.t1',
+ use_caching => false,
+ spec_id => 2
+)
+```
+
+### Verifying the Result
+
+```bash
+curl -sS "http://localhost:8090/api/metalakes/test/jobs/runs/{job_id}" | jq
'.job.status'
+cat
/tmp/gravitino/jobs/staging/test/builtin-iceberg-rewrite-manifests/{job_id}/output.log
+```
+
+A successful run reports its state as `SUCCEEDED` and logs how many manifests
it replaced:
+
+```text
+Rewrite Manifests Results: Rewritten manifests: 24, Added manifests: 2
+```
+
+Both counts at zero indicate a successful no-op, for example when no manifests
for the selected spec need rewriting. Other specs remain unchanged.
+
## Related
- [Table Maintenance Service](./optimizer.md) for the concepts and the
walkthrough
diff --git a/docs/table-maintenance-service/optimizer-configuration.md
b/docs/table-maintenance-service/optimizer-configuration.md
index 5ce59539e6..465afb3e06 100644
--- a/docs/table-maintenance-service/optimizer-configuration.md
+++ b/docs/table-maintenance-service/optimizer-configuration.md
@@ -112,6 +112,29 @@ your Spark, Scala, and Iceberg versions. Details are under
`warehouse_location` may be empty for local filesystem testing. Set it to the
warehouse URI for HDFS or cloud object storage.
+### Rewrite Manifests Job
+
+Submit `builtin-iceberg-rewrite-manifests` directly through
+`POST /api/metalakes/{metalake}/jobs/runs`. Its `jobConf` uses these
job-specific keys:
+
+| Key | Meaning |
Default |
+| ------------------ | ----------------------------------------------------- |
--------------------------------------------- |
+| `catalog_name` | Iceberg catalog registered in Spark |
Required |
+| `table_identifier` | Table identifier, for example `db.t1` |
Required |
+| `spec_id` | Existing partition spec whose manifests to rewrite |
Current table spec |
+| `use_caching` | `true` or `false` to control caching during rewriting |
Installed Iceberg default (`false` in 1.11.0) |
+| `spark_conf` | JSON string containing additional Spark settings |
None |
+
+Include the Spark and catalog template settings shown in the
+[submission example](./optimizer-cli-reference.md#submitting-the-job), and
make the matching
+Iceberg Spark runtime available as described above. For this job, omitted,
blank, or
+unresolved optional argument placeholders are treated as absent. Required
catalog/table
+arguments cannot be blank or unresolved. `spec_id` must be a non-negative
integer identifying
+an existing spec; it does not repartition data or migrate manifests between
specs.
+
+This job currently supports direct submission. Policy -> Strategy -> Adapter
integration
+for automatic manifest maintenance will follow separately.
+
## Running Against a Local Filesystem
On a machine with no HDFS, Spark still defaults to `hdfs://localhost:9000` and
fails. Set the default filesystem explicitly, in `spark_conf` for job
submissions and in the CLI `spark_conf` value:
diff --git a/docs/table-maintenance-service/optimizer.md
b/docs/table-maintenance-service/optimizer.md
index 193f7f8c9b..affd3f0b5b 100644
--- a/docs/table-maintenance-service/optimizer.md
+++ b/docs/table-maintenance-service/optimizer.md
@@ -313,6 +313,23 @@ grep -E "Rewritten data files|Added data files|completed
successfully" "${log_di
REST job status is polled rather than pushed, so it lags the real Spark
process by up to one poll interval. That is why the prerequisites lower it to
ten seconds.
+## Rewriting Iceberg Manifests
+
+`builtin-iceberg-rewrite-manifests` consolidates manifest metadata for one
existing partition
+spec to improve scan planning. It leaves data files and other specs' manifests
unchanged.
+Submit it directly with `POST /api/metalakes/{metalake}/jobs/runs`, setting
+`jobTemplateName` to `builtin-iceberg-rewrite-manifests` and providing
`catalog_name`,
+`table_identifier`, and the Spark/catalog settings in `jobConf`.
+
+Omit `spec_id` to maintain the current partition spec, or supply an existing
spec ID to
+maintain that spec. Optional `use_caching` controls Spark caching during the
rewrite.
+See [Rewrite Manifests](./optimizer-cli-reference.md#rewrite-manifests) for
the complete
+submission example and result checks, and
[configuration](./optimizer-configuration.md#rewrite-manifests-job)
+for the job keys. Results appear in the job's `output.log`.
+
+Policy-driven manifest maintenance through Policy -> Strategy -> Adapter is a
follow-up;
+`submit-strategy-jobs` does not yet select this job automatically.
+
## Related
- [Configuration](./optimizer-configuration.md) for the three configuration
layers
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 471824af43..df7cae0888 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
@@ -28,6 +28,7 @@ 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.IcebergRewriteManifestsJob;
import
org.apache.gravitino.maintenance.jobs.iceberg.IcebergUpdateStatsAndMetricsJob;
import org.apache.gravitino.maintenance.jobs.spark.SparkPiJob;
import org.slf4j.Logger;
@@ -47,6 +48,7 @@ public class BuiltInJobTemplateProvider implements
JobTemplateProvider {
ImmutableList.of(
new SparkPiJob(),
new IcebergRewriteDataFilesJob(),
+ new IcebergRewriteManifestsJob(),
new IcebergUpdateStatsAndMetricsJob(),
new IcebergExpireSnapshotsJob(),
new IcebergRemoveOrphanFilesJob());
diff --git
a/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergJobUtils.java
b/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergJobUtils.java
index 9f95e92405..0340be3c33 100644
---
a/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergJobUtils.java
+++
b/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergJobUtils.java
@@ -20,6 +20,9 @@ package org.apache.gravitino.maintenance.jobs.iceberg;
import java.util.HashMap;
import java.util.Map;
+import java.util.Set;
+import java.util.regex.Pattern;
+import javax.annotation.Nullable;
import
org.apache.gravitino.maintenance.optimizer.common.util.IcebergSparkConfigUtils;
import org.apache.spark.sql.SparkSession;
@@ -32,9 +35,33 @@ import org.apache.spark.sql.SparkSession;
public final class IcebergJobUtils {
private static final String ICEBERG_SPARK_CATALOG =
"org.apache.iceberg.spark.SparkCatalog";
+ /**
+ * Matches a job template placeholder that no job configuration value
replaced, e.g. {@code
+ * {{use_caching}}}.
+ */
+ private static final Pattern UNRESOLVED_PLACEHOLDER_PATTERN =
Pattern.compile("^\\{\\{[^{}]*}}$");
private IcebergJobUtils() {}
+ /**
+ * Return the given argument value unless it is an unresolved job template
placeholder.
+ *
+ * <p>Job templates declare optional parameters as {@code {{name}}}
placeholders. When the caller
+ * omits a parameter, the server leaves the placeholder untouched and it
reaches the job as a
+ * literal {@code "{{name}}"} argument. Treating that as a real value would
forward nonsense to
+ * Iceberg, so callers use this method to map it back to "not supplied".
+ *
+ * @param value the argument value to inspect
+ * @return the value, or null if it is null or an unresolved placeholder
+ */
+ @Nullable
+ public static String nullIfUnresolvedPlaceholder(@Nullable String value) {
+ if (value == null ||
UNRESOLVED_PLACEHOLDER_PATTERN.matcher(value.trim()).matches()) {
+ return null;
+ }
+ return value;
+ }
+
/**
* Escape single quotes in SQL string literals by replacing ' with ''.
*
@@ -99,6 +126,45 @@ public final class IcebergJobUtils {
return argMap;
}
+ /**
+ * Parse named arguments with explicit values and validate supported and
required names.
+ *
+ * <p>Blank values and unresolved template placeholders are normalized to
null. Duplicate names,
+ * unknown names, missing values, and absent required arguments are
rejected. Unlike the
+ * single-argument overload, this parser does not accept bare boolean flags.
+ *
+ * @param args command line arguments
+ * @param supported allowed argument names without the leading dashes
+ * @param requiredArguments required argument names without the leading
dashes
+ * @return map of argument names to normalized values
+ * @throws IllegalArgumentException if arguments are invalid
+ */
+ public static Map<String, String> parseArguments(
+ String[] args, Set<String> supported, Set<String> requiredArguments) {
+ Map<String, String> parsed = new HashMap<>();
+ for (int i = 0; i < args.length; i += 2) {
+ String flag = args[i];
+ if (flag == null || !flag.startsWith("--") ||
!supported.contains(flag.substring(2))) {
+ throw new IllegalArgumentException("Unknown argument: " + flag);
+ }
+ if (i + 1 == args.length || args[i + 1] == null || args[i +
1].startsWith("--")) {
+ throw new IllegalArgumentException("Missing value for " + flag);
+ }
+ String key = flag.substring(2);
+ if (parsed.containsKey(key)) {
+ throw new IllegalArgumentException("Duplicate argument: " + flag);
+ }
+ String value = nullIfUnresolvedPlaceholder(args[i + 1].trim());
+ parsed.put(key, value == null || value.isEmpty() ? null : value);
+ }
+ for (String required : requiredArguments) {
+ if (parsed.get(required) == null) {
+ throw new IllegalArgumentException("--" + required + " is required");
+ }
+ }
+ return parsed;
+ }
+
/**
* Parse custom Spark configurations from JSON string.
*
diff --git
a/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergRewriteManifestsJob.java
b/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergRewriteManifestsJob.java
new file mode 100644
index 0000000000..4044379c7e
--- /dev/null
+++
b/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergRewriteManifestsJob.java
@@ -0,0 +1,266 @@
+/*
+ * 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 java.util.Arrays;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+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.spark.sql.Row;
+import org.apache.spark.sql.SparkSession;
+
+/**
+ * Built-in job for rewriting Iceberg table manifest files.
+ *
+ * <p>This job leverages Iceberg's RewriteManifestsProcedure to consolidate
small manifest files and
+ * cluster manifest entries within an existing partition spec to improve scan
planning.
+ */
+public class IcebergRewriteManifestsJob implements BuiltInJob {
+
+ private static final String NAME =
+ JobTemplateProvider.BUILTIN_NAME_PREFIX + "iceberg-rewrite-manifests";
+ private static final String VERSION = "v1";
+
+ @Override
+ public SparkJobTemplate jobTemplate() {
+ return SparkJobTemplate.builder()
+ .withName(NAME)
+ .withComment(
+ "Built-in Iceberg rewrite manifests job template for scan planning
optimization")
+ .withExecutable(resolveExecutable(IcebergRewriteManifestsJob.class))
+ .withClassName(IcebergRewriteManifestsJob.class.getName())
+ .withArguments(buildArguments())
+ .withConfigs(buildSparkConfigs())
+ .withCustomFields(
+ Collections.singletonMap(JobTemplateProvider.PROPERTY_VERSION_KEY,
VERSION))
+ .build();
+ }
+
+ /**
+ * Main entry point for the rewrite manifests job.
+ *
+ * <p>Uses named arguments for flexibility:
+ *
+ * <ul>
+ * <li>--catalog <catalog_name> Required. Iceberg catalog name.
+ * <li>--table <table_identifier> Required. Table name (db.table)
+ * <li>--use-caching <boolean> Optional. Whether to cache the table
metadata in Spark
+ * while rewriting (default: Iceberg's own default)
+ * <li>--spec-id <int> Optional. Rewrite manifests belonging to this
partition spec ID
+ * (default: the table's current spec)
+ * <li>--spark-conf <spark_conf_json> Optional. JSON map of custom
Spark configurations
+ * </ul>
+ *
+ * <p><b>Important Notes on Special Characters:</b>
+ *
+ * <ul>
+ * <li><b>Via Gravitino API:</b> Pass values as-is without shell escaping.
Gravitino handles
+ * escaping internally via ProcessBuilder.
+ * <li><b>Via Command Line:</b> Use shell quoting for values containing
whitespace.
+ * </ul>
+ *
+ * <p>Example via command line: --catalog iceberg_catalog --table db.sample
--use-caching false
+ *
+ * <p>Example via Gravitino API:
+ *
+ * <pre>{@code
+ * Map<String, String> jobConf = new HashMap<>();
+ * jobConf.put("catalog_name", "iceberg_catalog");
+ * jobConf.put("table_identifier", "db.sample");
+ * jobConf.put("use_caching", "false");
+ * metalake.runJob("builtin-iceberg-rewrite-manifests", jobConf);
+ * }</pre>
+ *
+ * @param args named job arguments
+ */
+ public static void main(String[] args) {
+ int exitCode = run(args);
+ if (exitCode != 0) {
+ System.exit(exitCode);
+ }
+ }
+
+ static int run(String[] args) {
+ try {
+ execute(args);
+ return 0;
+ } catch (Exception e) {
+ System.err.println("Error rewriting manifests: " + e.getMessage());
+ printUsage();
+ return 1;
+ }
+ }
+
+ static Map<String, String> parseArguments(String[] args) {
+ Map<String, String> parsed =
+ IcebergJobUtils.parseArguments(
+ args,
+ new HashSet<>(
+ Arrays.asList("catalog", "table", "use-caching", "spec-id",
"spark-conf")),
+ new HashSet<>(Arrays.asList("catalog", "table")));
+ validateUseCaching(parsed.get("use-caching"));
+ validateSpecId(parsed.get("spec-id"));
+ return parsed;
+ }
+
+ /**
+ * Build the SQL CALL statement for the rewrite_manifests procedure.
+ *
+ * @param catalogName Iceberg catalog name
+ * @param tableIdentifier Fully qualified table name
+ * @param useCaching Whether to cache table metadata during the rewrite
+ * @param specId Existing partition spec ID whose manifests to rewrite
+ * @return SQL CALL statement
+ */
+ static String buildProcedureCall(
+ String catalogName,
+ String tableIdentifier,
+ @Nullable String useCaching,
+ @Nullable String specId) {
+ StringBuilder sql = new StringBuilder();
+ sql.append("CALL ")
+ .append(IcebergJobUtils.escapeSqlIdentifier(catalogName))
+ .append(".system.rewrite_manifests(");
+ sql.append("table =>
'").append(IcebergJobUtils.escapeSqlString(tableIdentifier)).append("'");
+
+ if (useCaching != null && !useCaching.isEmpty()) {
+ sql.append(", use_caching => ").append(Boolean.parseBoolean(useCaching));
+ }
+
+ if (specId != null && !specId.isEmpty()) {
+ sql.append(", spec_id => ").append(Integer.parseInt(specId));
+ }
+
+ sql.append(")");
+ return sql.toString();
+ }
+
+ /**
+ * Validate the use-caching parameter value.
+ *
+ * <p>{@link Boolean#parseBoolean(String)} maps anything that is not {@code
"true"} to {@code
+ * false}, so a typo would silently disable caching. Reject such values
instead.
+ *
+ * @param useCaching the use-caching value to validate
+ * @throws IllegalArgumentException if the value is neither "true" nor
"false"
+ */
+ static void validateUseCaching(@Nullable String useCaching) {
+ if (useCaching == null || useCaching.isEmpty()) {
+ return; // use-caching is optional
+ }
+
+ if (!"true".equalsIgnoreCase(useCaching) &&
!"false".equalsIgnoreCase(useCaching)) {
+ throw new IllegalArgumentException(
+ "Invalid use-caching value '" + useCaching + "'. Must be either
'true' or 'false'");
+ }
+ }
+
+ /**
+ * Validate the spec-id parameter value.
+ *
+ * <p>Iceberg partition spec IDs start at 0 and the procedure rejects an
unknown ID, but failing
+ * here keeps a malformed value from reaching Spark as an unparseable SQL
literal.
+ *
+ * @param specId the spec-id value to validate
+ * @throws IllegalArgumentException if the value is not a non-negative
integer
+ */
+ static void validateSpecId(@Nullable String specId) {
+ if (specId == null || specId.isEmpty()) {
+ return; // spec-id is optional
+ }
+
+ try {
+ if (Integer.parseInt(specId) < 0) {
+ throw new IllegalArgumentException(
+ "Invalid spec-id value '" + specId + "'. Must be a non-negative
integer");
+ }
+ } catch (NumberFormatException e) {
+ throw new IllegalArgumentException(
+ "Invalid spec-id value '" + specId + "'. Must be a non-negative
integer");
+ }
+ }
+
+ private static void execute(String[] args) {
+ Map<String, String> argMap = parseArguments(args);
+ String sql =
+ buildProcedureCall(
+ argMap.get("catalog"),
+ argMap.get("table"),
+ argMap.get("use-caching"),
+ argMap.get("spec-id"));
+ Map<String, String> configs =
IcebergJobUtils.parseCustomSparkConfigs(argMap.get("spark-conf"));
+ SparkSession.Builder builder =
+ SparkSession.builder().appName("Gravitino Built-in Iceberg Rewrite
Manifests");
+ configs.forEach(builder::config);
+
+ try (SparkSession spark = builder.getOrCreate()) {
+ IcebergJobUtils.requireIcebergSparkRuntime();
+ System.out.println("Executing Iceberg rewrite_manifests procedure: " +
sql);
+ List<Row> results = spark.sql(sql).collectAsList();
+ if (!results.isEmpty()) {
+ Row result = results.get(0);
+ System.out.printf(
+ "Rewrite Manifests Results: Rewritten manifests: %d, Added
manifests: %d%n",
+ ((Number) result.get(0)).longValue(), ((Number)
result.get(1)).longValue());
+ }
+ System.out.println("Rewrite manifests job completed successfully");
+ }
+ }
+
+ private static void printUsage() {
+ System.err.println(
+ "Usage: IcebergRewriteManifestsJob --catalog <catalog_name> --table
<table_identifier>"
+ + " [--use-caching <true|false>] [--spec-id <non-negative
integer>]"
+ + " [--spark-conf <json>]");
+ }
+
+ /**
+ * Build template arguments list with named argument format.
+ *
+ * @return list of template arguments
+ */
+ private static List<String> buildArguments() {
+ return Arrays.asList(
+ "--catalog",
+ "{{catalog_name}}",
+ "--table",
+ "{{table_identifier}}",
+ "--use-caching",
+ "{{use_caching}}",
+ "--spec-id",
+ "{{spec_id}}",
+ "--spark-conf",
+ "{{spark_conf}}");
+ }
+
+ /**
+ * Build Spark configuration template.
+ *
+ * @return map of Spark configuration keys to template values
+ */
+ private static Map<String, String> buildSparkConfigs() {
+ return IcebergSparkConfigUtils.buildTemplateSparkConfigs();
+ }
+}
diff --git
a/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/TestBuiltInJobTemplateProvider.java
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/TestBuiltInJobTemplateProvider.java
index f2cf0f6937..4470dc10df 100644
---
a/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/TestBuiltInJobTemplateProvider.java
+++
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/TestBuiltInJobTemplateProvider.java
@@ -18,6 +18,7 @@
*/
package org.apache.gravitino.maintenance.jobs;
+import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -29,6 +30,17 @@ import org.junit.jupiter.api.Test;
public class TestBuiltInJobTemplateProvider {
+ /** Ensures manifest rewriting is discoverable exactly once through the
built-in provider. */
+ @Test
+ public void testRewriteManifestsTemplateRegistered() {
+ assertEquals(
+ 1L,
+ new BuiltInJobTemplateProvider()
+ .jobTemplates().stream()
+ .filter(template ->
template.name().equals("builtin-iceberg-rewrite-manifests"))
+ .count());
+ }
+
@Test
public void testJobTemplatesReturnsNonEmptyList() {
BuiltInJobTemplateProvider provider = new BuiltInJobTemplateProvider();
diff --git
a/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergJobUtils.java
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergJobUtils.java
index 5f03f58e3e..89fb606e3d 100644
---
a/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergJobUtils.java
+++
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergJobUtils.java
@@ -19,13 +19,58 @@
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.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.Map;
+import java.util.Set;
import org.junit.jupiter.api.Test;
public class TestIcebergJobUtils {
+ @Test
+ public void testStrictArgumentsNormalizeTemplateValues() {
+ Map<String, String> parsed =
+ IcebergJobUtils.parseArguments(
+ new String[] {"--required", " value ", "--optional", "
{{optional}} "},
+ new HashSet<>(Arrays.asList("required", "optional")),
+ Collections.singleton("required"));
+ assertEquals("value", parsed.get("required"));
+ assertNull(parsed.get("optional"));
+ }
+
+ @Test
+ public void testStrictArgumentsRejectMalformedInput() {
+ Set<String> supported = Collections.singleton("key");
+ String[][] invalid = {
+ {"--unknown", "value"},
+ {"key", "value"},
+ {null, "value"},
+ {"--key"},
+ {"--key", null},
+ {"--key", "--key"},
+ {"--key", "", "--key", "value"},
+ {"--key", "{{key}}"},
+ {"--key", " "},
+ {}
+ };
+ for (String[] args : invalid) {
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> IcebergJobUtils.parseArguments(args, supported, supported));
+ }
+ }
+
+ @Test
+ public void testLegacyArgumentsStillAcceptBareFlags() {
+ assertEquals("true", IcebergJobUtils.parseArguments(new String[]
{"--flag"}).get("flag"));
+ }
+
@Test
public void testRequireIcebergSparkRuntimeSucceedsWhenPresent() {
// Test classpath includes iceberg-spark-runtime.
diff --git
a/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRewriteManifestsJob.java
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRewriteManifestsJob.java
new file mode 100644
index 0000000000..e4b5fa322b
--- /dev/null
+++
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRewriteManifestsJob.java
@@ -0,0 +1,446 @@
+/*
+ * 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.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.io.ByteArrayOutputStream;
+import java.io.PrintStream;
+import java.nio.charset.StandardCharsets;
+import java.util.Map;
+import org.apache.gravitino.job.JobTemplateProvider;
+import org.apache.gravitino.job.SparkJobTemplate;
+import org.junit.jupiter.api.Test;
+
+public class TestIcebergRewriteManifestsJob {
+
+ @Test
+ public void testInvalidEntryArgumentsFailBeforeSparkStarts() {
+ String[][] invalid = {
+ {},
+ {"--catalog", "cat", "--table"},
+ {"--catalog", "{{catalog_name}}", "--table", "db.t"},
+ {"--catalog", " ", "--table", "db.t"},
+ {"--catalog", "cat", "--table", "db.t", "--use-caching", "yes"},
+ {"--catalog", "cat", "--table", "db.t", "--spec-id", "-1"},
+ {"--catalog", "cat", "--table", "db.t", "--spark-conf", "not-json"},
+ {"--catalog", "cat", "--table", "db.t", "--unknown", "value"},
+ {"--catalog", "cat", "--table", "db.t", "--catalog", "other"}
+ };
+ for (String[] args : invalid) {
+ PrintStream originalErr = System.err;
+ ByteArrayOutputStream errors = new ByteArrayOutputStream();
+ try (PrintStream captured = new PrintStream(errors, true,
StandardCharsets.UTF_8)) {
+ System.setErr(captured);
+ assertEquals(1, IcebergRewriteManifestsJob.run(args));
+ } finally {
+ System.setErr(originalErr);
+ }
+ String message = errors.toString(StandardCharsets.UTF_8);
+ assertTrue(message.startsWith("Error rewriting manifests:"));
+ assertTrue(message.contains("Usage: IcebergRewriteManifestsJob"));
+ assertFalse(message.contains("\tat "));
+ }
+ }
+
+ @Test
+ public void testParseArgumentsRejectsInvalidNamesAndMissingValues() {
+ String[][] invalid = {
+ {"--catalog", "cat", "--table", "db.t", "--unknown", "value"},
+ {"--catalog", "cat", "--table", "db.t", "--spec-id", "", "--spec-id",
"0"},
+ {"--catalog", "cat", "--table", "{{table_identifier}}"},
+ {"--catalog", "cat", "--table", "db.t", "--use-caching"},
+ {"--catalog", "cat", "--table", "db.t", "--spec-id", "--use-caching",
"true"},
+ {"--catalog", "cat", "--table", null},
+ {null, "cat", "--table", "db.t"}
+ };
+ for (String[] args : invalid) {
+ assertThrows(
+ IllegalArgumentException.class, () ->
IcebergRewriteManifestsJob.parseArguments(args));
+ }
+ }
+
+ @Test
+ public void testEntryArgumentsOmitEmptyAndUnresolvedOptionals() {
+ Map<String, String> args =
+ IcebergRewriteManifestsJob.parseArguments(
+ new String[] {
+ "--catalog",
+ "cat",
+ "--table",
+ "db.t",
+ "--use-caching",
+ "{{use_caching}}",
+ "--spec-id",
+ "",
+ "--spark-conf",
+ " "
+ });
+ assertNull(args.get("use-caching"));
+ assertNull(args.get("spec-id"));
+ assertNull(args.get("spark-conf"));
+ }
+
+ @Test
+ public void testJobTemplateHasCorrectName() {
+ IcebergRewriteManifestsJob job = new IcebergRewriteManifestsJob();
+ SparkJobTemplate template = job.jobTemplate();
+
+ assertNotNull(template);
+ assertEquals("builtin-iceberg-rewrite-manifests", template.name());
+ }
+
+ @Test
+ public void testJobTemplateHasComment() {
+ IcebergRewriteManifestsJob job = new IcebergRewriteManifestsJob();
+ SparkJobTemplate template = job.jobTemplate();
+
+ assertNotNull(template.comment());
+ assertFalse(template.comment().trim().isEmpty());
+ assertTrue(template.comment().contains("Iceberg"));
+ assertTrue(template.comment().contains("manifests"));
+ }
+
+ @Test
+ public void testJobTemplateHasExecutable() {
+ IcebergRewriteManifestsJob job = new IcebergRewriteManifestsJob();
+ SparkJobTemplate template = job.jobTemplate();
+
+ assertNotNull(template.executable());
+ assertFalse(template.executable().trim().isEmpty());
+ }
+
+ @Test
+ public void testJobTemplateHasClassName() {
+ IcebergRewriteManifestsJob job = new IcebergRewriteManifestsJob();
+ SparkJobTemplate template = job.jobTemplate();
+
+ assertNotNull(template.className());
+ assertEquals(IcebergRewriteManifestsJob.class.getName(),
template.className());
+ }
+
+ @Test
+ public void testJobTemplateHasArguments() {
+ IcebergRewriteManifestsJob job = new IcebergRewriteManifestsJob();
+ SparkJobTemplate template = job.jobTemplate();
+
+ assertNotNull(template.arguments());
+ assertEquals(10, template.arguments().size()); // 5 flags * 2 (flag +
value)
+
+ // Verify all expected arguments are present
+ assertTrue(template.arguments().contains("--catalog"));
+ assertTrue(template.arguments().contains("{{catalog_name}}"));
+ assertTrue(template.arguments().contains("--table"));
+ assertTrue(template.arguments().contains("{{table_identifier}}"));
+ assertTrue(template.arguments().contains("--use-caching"));
+ assertTrue(template.arguments().contains("{{use_caching}}"));
+ assertTrue(template.arguments().contains("--spec-id"));
+ assertTrue(template.arguments().contains("{{spec_id}}"));
+ assertTrue(template.arguments().contains("--spark-conf"));
+ assertTrue(template.arguments().contains("{{spark_conf}}"));
+ }
+
+ @Test
+ public void testJobTemplateHasSparkConfigs() {
+ IcebergRewriteManifestsJob job = new IcebergRewriteManifestsJob();
+ SparkJobTemplate template = job.jobTemplate();
+
+ Map<String, String> configs = template.configs();
+ assertNotNull(configs);
+ assertFalse(configs.isEmpty());
+
+ // Verify Spark runtime configs
+ assertTrue(configs.containsKey("spark.master"));
+ assertTrue(configs.containsKey("spark.executor.instances"));
+ assertTrue(configs.containsKey("spark.executor.cores"));
+ assertTrue(configs.containsKey("spark.executor.memory"));
+ assertTrue(configs.containsKey("spark.driver.memory"));
+
+ // Verify Iceberg catalog configs
+ assertTrue(configs.containsKey("spark.sql.catalog.{{catalog_name}}"));
+ assertTrue(configs.containsKey("spark.sql.catalog.{{catalog_name}}.type"));
+ assertTrue(configs.containsKey("spark.sql.catalog.{{catalog_name}}.uri"));
+
assertTrue(configs.containsKey("spark.sql.catalog.{{catalog_name}}.warehouse"));
+ }
+
+ @Test
+ public void testJobTemplateHasVersion() {
+ IcebergRewriteManifestsJob job = new IcebergRewriteManifestsJob();
+ SparkJobTemplate template = job.jobTemplate();
+
+ Map<String, String> customFields = template.customFields();
+ assertNotNull(customFields);
+
assertTrue(customFields.containsKey(JobTemplateProvider.PROPERTY_VERSION_KEY));
+
+ String version =
customFields.get(JobTemplateProvider.PROPERTY_VERSION_KEY);
+ assertEquals("v1", version);
+ assertTrue(version.matches(JobTemplateProvider.VERSION_VALUE_PATTERN));
+ }
+
+ @Test
+ public void testJobTemplateNameMatchesBuiltInPattern() {
+ IcebergRewriteManifestsJob job = new IcebergRewriteManifestsJob();
+ SparkJobTemplate template = job.jobTemplate();
+
+
assertTrue(template.name().matches(JobTemplateProvider.BUILTIN_NAME_PATTERN));
+
assertTrue(template.name().startsWith(JobTemplateProvider.BUILTIN_NAME_PREFIX));
+ }
+
+ // Test parseArguments method
+
+ @Test
+ public void testParseArgumentsWithAllRequired() {
+ String[] args = {"--catalog", "iceberg_prod", "--table", "db.sample"};
+ Map<String, String> result =
IcebergRewriteManifestsJob.parseArguments(args);
+
+ assertEquals(2, result.size());
+ assertEquals("iceberg_prod", result.get("catalog"));
+ assertEquals("db.sample", result.get("table"));
+ }
+
+ @Test
+ public void testParseArgumentsWithOptional() {
+ String[] args = {
+ "--catalog", "iceberg_prod",
+ "--table", "db.sample",
+ "--use-caching", "false"
+ };
+ Map<String, String> result =
IcebergRewriteManifestsJob.parseArguments(args);
+
+ assertEquals(3, result.size());
+ assertEquals("iceberg_prod", result.get("catalog"));
+ assertEquals("db.sample", result.get("table"));
+ assertEquals("false", result.get("use-caching"));
+ }
+
+ @Test
+ public void testParseArgumentsWithEmptyValues() {
+ String[] args = {"--catalog", "iceberg_prod", "--table", "db.sample",
"--use-caching", ""};
+ Map<String, String> result =
IcebergRewriteManifestsJob.parseArguments(args);
+
+ // Empty optional values are normalized to null.
+ assertEquals(3, result.size());
+ assertEquals("iceberg_prod", result.get("catalog"));
+ assertEquals("db.sample", result.get("table"));
+ assertNull(result.get("use-caching"));
+ }
+
+ @Test
+ public void testParseArgumentsOrderIndependent() {
+ String[] args1 = {"--catalog", "cat1", "--table", "tbl1", "--use-caching",
"true"};
+ String[] args2 = {"--use-caching", "true", "--table", "tbl1", "--catalog",
"cat1"};
+
+ Map<String, String> result1 =
IcebergRewriteManifestsJob.parseArguments(args1);
+ Map<String, String> result2 =
IcebergRewriteManifestsJob.parseArguments(args2);
+
+ assertEquals(result1, result2);
+ }
+
+ // Test unresolved placeholder handling
+
+ @Test
+ public void testUnresolvedPlaceholderIsDroppedForOptionalArguments() {
+ // The server leaves placeholders untouched when the caller omits the
parameter, so they reach
+ // the job as literal arguments.
+ String[] args = {
+ "--catalog", "iceberg_prod",
+ "--table", "db.sample",
+ "--use-caching", "{{use_caching}}",
+ "--spark-conf", "{{spark_conf}}"
+ };
+ Map<String, String> result =
IcebergRewriteManifestsJob.parseArguments(args);
+
+ assertNull(result.get("use-caching"));
+
assertNull(IcebergJobUtils.nullIfUnresolvedPlaceholder(result.get("use-caching")));
+
assertNull(IcebergJobUtils.nullIfUnresolvedPlaceholder(result.get("spark-conf")));
+ }
+
+ @Test
+ public void testRealValuesSurviveUnresolvedPlaceholderFiltering() {
+ assertEquals("false",
IcebergJobUtils.nullIfUnresolvedPlaceholder("false"));
+ assertEquals("db.sample",
IcebergJobUtils.nullIfUnresolvedPlaceholder("db.sample"));
+ // Only a value that is entirely a placeholder is dropped.
+ assertEquals("{{a}}b",
IcebergJobUtils.nullIfUnresolvedPlaceholder("{{a}}b"));
+ assertNull(IcebergJobUtils.nullIfUnresolvedPlaceholder(null));
+ }
+
+ @Test
+ public void testProcedureCallOmitsUseCachingWhenPlaceholderUnresolved() {
+ String useCaching =
IcebergJobUtils.nullIfUnresolvedPlaceholder("{{use_caching}}");
+ String sql =
+ IcebergRewriteManifestsJob.buildProcedureCall(
+ "iceberg_prod", "db.sample", useCaching, null);
+
+ // Without filtering, Boolean.parseBoolean("{{use_caching}}") would
silently emit
+ // use_caching => false instead of leaving Iceberg's default in place.
+ assertEquals("CALL `iceberg_prod`.system.rewrite_manifests(table =>
'db.sample')", sql);
+ }
+
+ // Test validateUseCaching method
+
+ @Test
+ public void testValidateUseCachingAcceptsBooleans() {
+ assertDoesNotThrow(() ->
IcebergRewriteManifestsJob.validateUseCaching("true"));
+ assertDoesNotThrow(() ->
IcebergRewriteManifestsJob.validateUseCaching("false"));
+ assertDoesNotThrow(() ->
IcebergRewriteManifestsJob.validateUseCaching("TRUE"));
+ }
+
+ @Test
+ public void testValidateUseCachingAcceptsAbsentValue() {
+ assertDoesNotThrow(() ->
IcebergRewriteManifestsJob.validateUseCaching(null));
+ assertDoesNotThrow(() ->
IcebergRewriteManifestsJob.validateUseCaching(""));
+ }
+
+ @Test
+ public void testValidateUseCachingRejectsNonBoolean() {
+ IllegalArgumentException e =
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> IcebergRewriteManifestsJob.validateUseCaching("yes"));
+
+ assertTrue(e.getMessage().contains("yes"));
+ }
+
+ // Test validateSpecId method
+
+ @Test
+ public void testValidateSpecIdAcceptsNonNegativeIntegers() {
+ assertDoesNotThrow(() -> IcebergRewriteManifestsJob.validateSpecId("0"));
+ assertDoesNotThrow(() -> IcebergRewriteManifestsJob.validateSpecId("2"));
+ }
+
+ @Test
+ public void testValidateSpecIdAcceptsAbsentValue() {
+ assertDoesNotThrow(() -> IcebergRewriteManifestsJob.validateSpecId(null));
+ assertDoesNotThrow(() -> IcebergRewriteManifestsJob.validateSpecId(""));
+ }
+
+ @Test
+ public void testValidateSpecIdRejectsNegativeValue() {
+ IllegalArgumentException e =
+ assertThrows(
+ IllegalArgumentException.class, () ->
IcebergRewriteManifestsJob.validateSpecId("-1"));
+
+ assertTrue(e.getMessage().contains("-1"));
+ }
+
+ @Test
+ public void testValidateSpecIdRejectsNonNumericValue() {
+ IllegalArgumentException e =
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> IcebergRewriteManifestsJob.validateSpecId("latest"));
+
+ assertTrue(e.getMessage().contains("latest"));
+ }
+
+ // Test buildProcedureCall method
+
+ @Test
+ public void testBuildProcedureCallMinimal() {
+ String sql =
+ IcebergRewriteManifestsJob.buildProcedureCall("iceberg_prod",
"db.sample", null, null);
+
+ assertEquals("CALL `iceberg_prod`.system.rewrite_manifests(table =>
'db.sample')", sql);
+ }
+
+ @Test
+ public void testBuildProcedureCallWithCachingTrue() {
+ String sql =
+ IcebergRewriteManifestsJob.buildProcedureCall("iceberg_prod",
"db.sample", "true", null);
+
+ assertEquals(
+ "CALL `iceberg_prod`.system.rewrite_manifests(table => 'db.sample',
use_caching => true)",
+ sql);
+ }
+
+ @Test
+ public void testBuildProcedureCallWithCachingFalse() {
+ String sql =
+ IcebergRewriteManifestsJob.buildProcedureCall("iceberg_prod",
"db.sample", "false", null);
+
+ assertEquals(
+ "CALL `iceberg_prod`.system.rewrite_manifests(table => 'db.sample',
use_caching => false)",
+ sql);
+ }
+
+ @Test
+ public void testBuildProcedureCallWithEmptyCaching() {
+ String sql =
+ IcebergRewriteManifestsJob.buildProcedureCall("iceberg_prod",
"db.sample", "", null);
+
+ assertEquals("CALL `iceberg_prod`.system.rewrite_manifests(table =>
'db.sample')", sql);
+ }
+
+ @Test
+ public void testBuildProcedureCallWithSpecId() {
+ String sql =
+ IcebergRewriteManifestsJob.buildProcedureCall("iceberg_prod",
"db.sample", null, "2");
+
+ assertEquals(
+ "CALL `iceberg_prod`.system.rewrite_manifests(table => 'db.sample',
spec_id => 2)", sql);
+ }
+
+ @Test
+ public void testBuildProcedureCallWithCachingAndSpecId() {
+ String sql =
+ IcebergRewriteManifestsJob.buildProcedureCall("iceberg_prod",
"db.sample", "false", "0");
+
+ assertEquals(
+ "CALL `iceberg_prod`.system.rewrite_manifests(table => 'db.sample', "
+ + "use_caching => false, spec_id => 0)",
+ sql);
+ }
+
+ @Test
+ public void testProcedureCallOmitsSpecIdWhenPlaceholderUnresolved() {
+ String specId = IcebergJobUtils.nullIfUnresolvedPlaceholder("{{spec_id}}");
+ String sql =
+ IcebergRewriteManifestsJob.buildProcedureCall("iceberg_prod",
"db.sample", null, specId);
+
+ // An unresolved placeholder must not reach Integer.parseInt as a SQL
literal.
+ assertEquals("CALL `iceberg_prod`.system.rewrite_manifests(table =>
'db.sample')", sql);
+ }
+
+ @Test
+ public void testBuildProcedureCallWithSqlInjectionAttempt() {
+ // Test SQL injection attempt in table name
+ String maliciousTable = "db.table' OR '1'='1";
+ String sql =
+ IcebergRewriteManifestsJob.buildProcedureCall(
+ "iceberg_catalog", maliciousTable, null, null);
+
+ // Verify single quotes are escaped (becomes '')
+ assertTrue(sql.contains("db.table'' OR ''1''=''1"));
+ assertFalse(sql.contains("' OR '1'='1"));
+
+ // Test SQL injection attempt in catalog name
+ String maliciousCatalog = "catalog`; DROP TABLE users; --";
+ sql = IcebergRewriteManifestsJob.buildProcedureCall(maliciousCatalog,
"db.table", null, null);
+
+ // Verify backticks are escaped and the whole identifier stays quoted
+ assertTrue(sql.startsWith("CALL `catalog``; DROP TABLE users;
--`.system.rewrite_manifests("));
+ }
+}
diff --git
a/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRewriteManifestsJobWithSpark.java
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRewriteManifestsJobWithSpark.java
new file mode 100644
index 0000000000..81edc2e4c1
--- /dev/null
+++
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergRewriteManifestsJobWithSpark.java
@@ -0,0 +1,258 @@
+/*
+ * 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 java.io.ByteArrayOutputStream;
+import java.io.PrintStream;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Path;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.stream.Collectors;
+import org.apache.spark.sql.Row;
+import org.apache.spark.sql.SparkSession;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+/** Exercises the job entry point against Spark 3.5 and Iceberg 1.11.0. */
+public class TestIcebergRewriteManifestsJobWithSpark {
+ private static final String TABLE = "manifest_catalog.db.events";
+
+ @TempDir Path temporaryDirectory;
+ private SparkSession spark;
+
+ /** Creates multiple manifests for both day and hour partition specs. */
+ @BeforeEach
+ public void setUp() {
+ startSpark();
+ spark.sql("CREATE NAMESPACE manifest_catalog.db");
+ spark.sql(
+ "CREATE TABLE "
+ + TABLE
+ + " (id INT, event_time TIMESTAMP) USING iceberg "
+ + "PARTITIONED BY (days(event_time)) "
+ + "TBLPROPERTIES ('commit.manifest-merge.enabled'='false')");
+ insertRows(0);
+ spark.sql(
+ "ALTER TABLE "
+ + TABLE
+ + " REPLACE PARTITION FIELD days(event_time) WITH
hours(event_time)");
+ insertRows(3);
+ }
+
+ /** Releases Spark resources after each test. */
+ @AfterEach
+ public void tearDown() {
+ if (spark != null) {
+ spark.stop();
+ SparkSession.clearActiveSession();
+ SparkSession.clearDefaultSession();
+ }
+ }
+
+ /** Omitting spec_id selects the current spec and leaves old-spec manifests
untouched. */
+ @Test
+ public void testDefaultSpecAndNoOpPreserveData() {
+ List<Row> records = records();
+ Set<String> oldManifests = manifests(0);
+ assertEquals(3, oldManifests.size());
+ assertEquals(3, manifests(1).size());
+ runJob(null, null);
+ assertEquals(oldManifests, manifests(0));
+ assertEquals(1, manifests(1).size());
+ assertEquals(records, records());
+ Set<String> allManifests = manifests(1);
+ long snapshot = currentSnapshot();
+ runJob(null, "false");
+ assertEquals(allManifests, manifests(1));
+ assertEquals(snapshot, currentSnapshot());
+ assertEquals(records, records());
+ }
+
+ /** An explicit old spec is rewritten without migrating it or changing the
current spec. */
+ @Test
+ public void testExplicitOldSpecWithCaching() {
+ List<Row> records = records();
+ Set<String> currentManifests = manifests(1);
+ assertEquals(3, manifests(0).size());
+ runJob("0", "true");
+ assertEquals(1, manifests(0).size());
+ assertEquals(currentManifests, manifests(1));
+ assertEquals(records, records());
+ // A subsequent default run still selects spec 1.
+ Set<String> oldManifests = manifests(0);
+ runJob(null, null);
+ assertEquals(oldManifests, manifests(0));
+ assertEquals(1, manifests(1).size());
+ }
+
+ /** A nonexistent spec fails the entry point, stops Spark, and leaves
metadata unchanged. */
+ @Test
+ public void testUnknownSpecFailsAndStopsSpark() {
+ long snapshot = currentSnapshot();
+ Set<String> oldManifests = manifests(0);
+ Set<String> newManifests = manifests(1);
+ SparkSession running = spark;
+ PrintStream originalErr = System.err;
+ ByteArrayOutputStream errors = new ByteArrayOutputStream();
+ try (PrintStream captured = new PrintStream(errors, true,
StandardCharsets.UTF_8)) {
+ System.setErr(captured);
+ assertEquals(1, IcebergRewriteManifestsJob.run(arguments("999", null)));
+ } finally {
+ System.setErr(originalErr);
+ }
+ String message = errors.toString(StandardCharsets.UTF_8);
+ assertTrue(message.contains("Error rewriting manifests:"));
+ assertTrue(message.contains("999"));
+ assertFalse(
+ message.contains(
+ "\tat
org.apache.gravitino.maintenance.jobs.iceberg.IcebergRewriteManifestsJob"));
+ assertTrue(running.sparkContext().isStopped());
+ startSpark();
+ assertEquals(snapshot, currentSnapshot());
+ assertEquals(oldManifests, manifests(0));
+ assertEquals(newManifests, manifests(1));
+ assertEquals(6, records().size());
+ }
+
+ private void runJob(String spec, String caching) {
+ spark.stop();
+ SparkSession.clearActiveSession();
+ SparkSession.clearDefaultSession();
+ PrintStream originalOut = System.out;
+ ByteArrayOutputStream output = new ByteArrayOutputStream();
+ try (PrintStream captured = new PrintStream(output, true,
StandardCharsets.UTF_8)) {
+ System.setOut(captured);
+ IcebergRewriteManifestsJob.main(arguments(spec, caching));
+ } finally {
+ System.setOut(originalOut);
+ }
+ String stdout = output.toString(StandardCharsets.UTF_8);
+ String expectedSql = "CALL
`manifest_catalog`.system.rewrite_manifests(table => 'db.events'";
+ if (caching != null) {
+ expectedSql += ", use_caching => " + caching;
+ }
+ if (spec != null) {
+ expectedSql += ", spec_id => " + spec;
+ }
+ expectedSql += ")";
+ int execution = stdout.indexOf("Executing Iceberg rewrite_manifests
procedure: " + expectedSql);
+ int results = stdout.indexOf("Rewrite Manifests Results: Rewritten
manifests:");
+ assertTrue(execution >= 0, stdout);
+ assertTrue(results > execution, stdout);
+ if (SparkSession.getActiveSession().isDefined()) {
+
assertTrue(SparkSession.getActiveSession().get().sparkContext().isStopped());
+ }
+ if (SparkSession.getDefaultSession().isDefined()) {
+
assertTrue(SparkSession.getDefaultSession().get().sparkContext().isStopped());
+ }
+ startSpark();
+ }
+
+ private String[] arguments(String spec, String caching) {
+ Map<String, String> jobConf = new HashMap<>();
+ jobConf.put("catalog_name", "manifest_catalog");
+ jobConf.put("table_identifier", "db.events");
+ jobConf.put(
+ "spark_conf",
+ "{\"spark.master\":\"local[2]\",\"spark.ui.enabled\":\"false\","
+ + "\"spark.sql.shuffle.partitions\":\"2\","
+ +
"\"spark.sql.extensions\":\"org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions\","
+ +
"\"spark.sql.catalog.manifest_catalog\":\"org.apache.iceberg.spark.SparkCatalog\","
+ + "\"spark.sql.catalog.manifest_catalog.type\":\"hadoop\","
+ + "\"spark.sql.catalog.manifest_catalog.warehouse\":\""
+ + temporaryDirectory.resolve("warehouse")
+ + "\"}");
+ if (spec != null) {
+ jobConf.put("spec_id", spec);
+ }
+ if (caching != null) {
+ jobConf.put("use_caching", caching);
+ }
+ // Model JobManager's template substitution, including unresolved optional
values.
+ return new IcebergRewriteManifestsJob()
+ .jobTemplate().arguments().stream()
+ .map(
+ value -> {
+ String resolved = value;
+ for (Map.Entry<String, String> entry : jobConf.entrySet()) {
+ resolved = resolved.replace("{{" + entry.getKey() + "}}",
entry.getValue());
+ }
+ return resolved;
+ })
+ .toArray(String[]::new);
+ }
+
+ private void startSpark() {
+ spark =
+ SparkSession.builder()
+ .master("local[2]")
+ .appName("TestIcebergRewriteManifestsJob")
+ .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.manifest_catalog",
"org.apache.iceberg.spark.SparkCatalog")
+ .config("spark.sql.catalog.manifest_catalog.type", "hadoop")
+ .config(
+ "spark.sql.catalog.manifest_catalog.warehouse",
+ temporaryDirectory.resolve("warehouse").toString())
+ .getOrCreate();
+ }
+
+ private void insertRows(int offset) {
+ for (int i = 0; i < 3; i++) {
+ spark.sql(
+ "INSERT INTO "
+ + TABLE
+ + " VALUES ("
+ + (offset + i)
+ + ", TIMESTAMP '2026-01-01 01:00:00')");
+ }
+ }
+
+ private Set<String> manifests(int spec) {
+ return spark
+ .sql("SELECT path FROM " + TABLE + ".manifests WHERE partition_spec_id
= " + spec)
+ .collectAsList()
+ .stream()
+ .map(row -> row.getString(0))
+ .collect(Collectors.toSet());
+ }
+
+ private List<Row> records() {
+ return spark.sql("SELECT * FROM " + TABLE + " ORDER BY
id").collectAsList();
+ }
+
+ private long currentSnapshot() {
+ return spark
+ .sql("SELECT snapshot_id FROM " + TABLE + ".refs WHERE name = 'main'")
+ .first()
+ .getLong(0);
+ }
+}