This is an automated email from the ASF dual-hosted git repository.
FANNG1 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 4d19799ef7 [#10394] fix(optimizer): Normalize CLI table identifiers
(#12935)
4d19799ef7 is described below
commit 4d19799ef79aeaac1daae43e430b73d652da2796
Author: Sasali <[email protected]>
AuthorDate: Fri Sep 18 15:20:27 2026 +0800
[#10394] fix(optimizer): Normalize CLI table identifiers (#12935)
### What changes were proposed in this pull request?
Normalize table identifiers used by `submit-strategy-jobs`,
`monitor-metrics`, and `list-table-metrics`. Fully qualified
`catalog.schema.table` identifiers pass through unchanged, while
`schema.table` identifiers use
`gravitino.optimizer.gravitinoDefaultCatalog`.
Keep job identifier parsing unchanged and add CLI regression tests for
all affected commands.
### Why are the changes needed?
The table-oriented commands currently parse identifiers directly, so
they pass `schema.table` to downstream components even when a default
catalog is configured. This differs from `submit-update-stats-job` and
from the documented default-catalog behavior. Without a default catalog,
the same input is accepted and fails later or is interpreted
inconsistently.
Fix: #10394
### Does this PR introduce _any_ user-facing change?
Table-oriented optimizer commands now accept `schema.table` when
`gravitino.optimizer.gravitinoDefaultCatalog` is configured and report a
clear validation error when it is not. Fully qualified table identifiers
and job identifiers keep their existing behavior. There are no API or
configuration changes.
### How was this patch tested?
- Added five CLI regression tests covering default-catalog normalization
in the three affected table commands, missing-default validation, and
unchanged job identifier parsing.
- Verified that four table-command regression tests fail on the original
implementation while the job-identifier regression test passes.
- Ran `:maintenance:optimizer:check -PskipTrinoConnector=true -PskipITs
-PskipDockerTests=true`: 225 tests passed with no failures or errors.
Five pre-existing environment integration tests were skipped because
`GRAVITINO_ENV_IT` was not enabled.
---
.../optimizer-cli-reference.md | 7 +-
.../optimizer/command/ListTableMetricsCommand.java | 3 +-
.../optimizer/command/MonitorMetricsCommand.java | 2 +-
.../optimizer/command/OptimizerCommandContext.java | 11 ++
.../optimizer/command/OptimizerCommandUtils.java | 20 ++++
.../command/SubmitStrategyJobsCommand.java | 4 +-
.../optimizer/command/UpdateStatisticsCommand.java | 2 +-
.../maintenance/optimizer/TestOptimizerCmd.java | 119 +++++++++++++++++++++
.../updater/StatisticsUpdaterForTest.java | 7 ++
9 files changed, 168 insertions(+), 7 deletions(-)
diff --git a/docs/table-maintenance-service/optimizer-cli-reference.md
b/docs/table-maintenance-service/optimizer-cli-reference.md
index 2389cbb956..81c35f0a1a 100644
--- a/docs/table-maintenance-service/optimizer-cli-reference.md
+++ b/docs/table-maintenance-service/optimizer-cli-reference.md
@@ -31,7 +31,7 @@ directory. Use `--conf-path` only when you need a custom
config file.
| Option | Meaning | Used by |
| --- | --- | --- |
-| `--identifiers` | Comma-separated identifiers. Table format supports
`catalog.schema.table` (or `schema.table` when default catalog is configured).
| Most commands |
+| `--identifiers` | Comma-separated identifiers. See [Identifier
Rules](#identifier-rules) for command-specific formats. | All commands |
| `--strategy-name` | Policy name to evaluate, for example
`iceberg_compaction_default`. | `submit-strategy-jobs` |
| `--dry-run` | Preview mode. Prints recommendations or job configs without
submitting jobs. | `submit-strategy-jobs`, `submit-update-stats-job` |
| `--limit` | Maximum number of strategy jobs to process. Must be `> 0`. |
`submit-strategy-jobs` |
@@ -79,7 +79,10 @@ job scopes with multiple metric/statistic fields:
## Identifier Rules
- Table and partition records: `catalog.schema.table`
-- If `gravitino.optimizer.gravitinoDefaultCatalog` is set, `schema.table` is
also accepted
+- If `gravitino.optimizer.gravitinoDefaultCatalog` is set, `schema.table` is
also accepted by
+ `submit-strategy-jobs`, `update-statistics`, `monitor-metrics`,
`list-table-metrics`, and
+ `submit-update-stats-job`
+- `append-metrics` accepts both table and job identifiers, so it does not
apply the default catalog
- Job records: parsed as a regular Gravitino `NameIdentifier`
## CLI Workflow Examples
diff --git
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/ListTableMetricsCommand.java
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/ListTableMetricsCommand.java
index 5e577a90f4..302e32e19c 100644
---
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/ListTableMetricsCommand.java
+++
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/ListTableMetricsCommand.java
@@ -38,9 +38,10 @@ public class ListTableMetricsCommand implements
OptimizerCommandExecutor {
Preconditions.checkArgument(
partitionPath.isEmpty() || context.identifiers().length == 1,
"--partition-path requires exactly one identifier");
+ List<NameIdentifier> identifiers = context.parsedTableIdentifiers();
try (MetricsProvider metricsProvider =
OptimizerCommandUtils.createMetricsProvider(context.optimizerEnv())) {
- for (NameIdentifier identifier : context.parsedIdentifiers()) {
+ for (NameIdentifier identifier : identifiers) {
List<MetricPoint> metrics =
partitionPath
.map(
diff --git
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/MonitorMetricsCommand.java
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/MonitorMetricsCommand.java
index 809da835b6..72ab3b5c76 100644
---
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/MonitorMetricsCommand.java
+++
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/MonitorMetricsCommand.java
@@ -42,7 +42,7 @@ public class MonitorMetricsCommand implements
OptimizerCommandExecutor {
partitionPath.isEmpty() || context.identifiers().length == 1,
"--partition-path requires exactly one identifier");
try (Monitor monitor = new Monitor(context.optimizerEnv())) {
- for (NameIdentifier identifier : context.parsedIdentifiers()) {
+ for (NameIdentifier identifier : context.parsedTableIdentifiers()) {
List<EvaluationResult> results =
monitor.evaluateMetrics(identifier, actionTimeSeconds,
rangeSeconds, partitionPath);
results.forEach(
diff --git
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/OptimizerCommandContext.java
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/OptimizerCommandContext.java
index e1e9d62e03..a7eb635f3b 100644
---
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/OptimizerCommandContext.java
+++
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/OptimizerCommandContext.java
@@ -26,6 +26,7 @@ import java.util.Optional;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.maintenance.optimizer.common.OptimizerEnv;
import
org.apache.gravitino.maintenance.optimizer.common.StatisticsInputContent;
+import org.apache.gravitino.maintenance.optimizer.common.conf.OptimizerConfig;
/** Context shared by optimizer command executors. */
public final class OptimizerCommandContext {
@@ -82,6 +83,16 @@ public final class OptimizerCommandContext {
return OptimizerCommandUtils.parseIdentifiers(identifiers);
}
+ /**
+ * Parses table identifiers, applying the configured default catalog to
schema-qualified names.
+ *
+ * @return normalized table identifiers
+ */
+ public List<NameIdentifier> parsedTableIdentifiers() {
+ return OptimizerCommandUtils.parseTableIdentifiers(
+ identifiers,
optimizerEnv.config().get(OptimizerConfig.GRAVITINO_DEFAULT_CATALOG_CONFIG));
+ }
+
public boolean hasIdentifiers() {
return identifiers != null && identifiers.length > 0;
}
diff --git
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/OptimizerCommandUtils.java
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/OptimizerCommandUtils.java
index 874ac537ce..87ef159bde 100644
---
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/OptimizerCommandUtils.java
+++
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/OptimizerCommandUtils.java
@@ -30,6 +30,7 @@ import
org.apache.gravitino.maintenance.optimizer.api.monitor.MetricsProvider;
import org.apache.gravitino.maintenance.optimizer.common.OptimizerEnv;
import
org.apache.gravitino.maintenance.optimizer.common.StatisticsInputContent;
import org.apache.gravitino.maintenance.optimizer.common.conf.OptimizerConfig;
+import org.apache.gravitino.maintenance.optimizer.common.util.IdentifierUtils;
import org.apache.gravitino.maintenance.optimizer.common.util.ProviderUtils;
import
org.apache.gravitino.maintenance.optimizer.recommender.util.PartitionUtils;
@@ -47,6 +48,25 @@ public final class OptimizerCommandUtils {
return Arrays.stream(identifiers).map(NameIdentifier::parse).toList();
}
+ static List<NameIdentifier> parseTableIdentifiers(
+ String[] identifiers, String defaultCatalogName) {
+ if (identifiers == null) {
+ return List.of();
+ }
+ return Arrays.stream(identifiers)
+ .map(
+ identifier ->
+ IdentifierUtils.parseTableIdentifier(identifier,
defaultCatalogName)
+ .orElseThrow(
+ () ->
+ new IllegalArgumentException(
+ String.format(
+ "Identifier '%s' is invalid. Use
catalog.schema.table, or "
+ + "configure %s when using
schema.table",
+ identifier,
OptimizerConfig.GRAVITINO_DEFAULT_CATALOG))))
+ .toList();
+ }
+
static Optional<PartitionPath> parsePartitionPath(String partitionPathStr) {
if (StringUtils.isBlank(partitionPathStr)) {
return Optional.empty();
diff --git
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/SubmitStrategyJobsCommand.java
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/SubmitStrategyJobsCommand.java
index 3e14328439..f4ae71120c 100644
---
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/SubmitStrategyJobsCommand.java
+++
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/SubmitStrategyJobsCommand.java
@@ -36,13 +36,13 @@ public class SubmitStrategyJobsCommand implements
OptimizerCommandExecutor {
if (context.dryRun()) {
results =
recommender.recommendForStrategyName(
- context.parsedIdentifiers(), context.strategyName(), limit);
+ context.parsedTableIdentifiers(), context.strategyName(),
limit);
results.forEach(
result ->
OptimizerOutputPrinter.printDryRunResult(context.output(), result));
} else {
results =
recommender.submitForStrategyName(
- context.parsedIdentifiers(), context.strategyName(), limit);
+ context.parsedTableIdentifiers(), context.strategyName(),
limit);
results.forEach(
result ->
OptimizerOutputPrinter.printSubmitResult(context.output(), result));
}
diff --git
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/UpdateStatisticsCommand.java
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/UpdateStatisticsCommand.java
index c73bc054c1..468c2707f7 100644
---
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/UpdateStatisticsCommand.java
+++
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/UpdateStatisticsCommand.java
@@ -38,7 +38,7 @@ public class UpdateStatisticsCommand implements
OptimizerCommandExecutor {
} else {
summary =
updater.update(
- context.calculatorName(), context.parsedIdentifiers(),
UpdateType.STATISTICS);
+ context.calculatorName(), context.parsedTableIdentifiers(),
UpdateType.STATISTICS);
}
OptimizerOutputPrinter.printUpdateSummary(context.output(), summary);
}
diff --git
a/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/TestOptimizerCmd.java
b/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/TestOptimizerCmd.java
index ce3d33b433..8ea3772db5 100644
---
a/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/TestOptimizerCmd.java
+++
b/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/TestOptimizerCmd.java
@@ -24,6 +24,9 @@ import java.io.PrintStream;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
+import java.nio.file.StandardOpenOption;
+import java.util.List;
+import org.apache.gravitino.NameIdentifier;
import
org.apache.gravitino.maintenance.optimizer.monitor.evaluator.MetricsEvaluatorForTest;
import
org.apache.gravitino.maintenance.optimizer.monitor.job.TableJobRelationProviderForTest;
import
org.apache.gravitino.maintenance.optimizer.monitor.metrics.MetricsProviderForTest;
@@ -156,6 +159,23 @@ class TestOptimizerCmd {
Assertions.assertTrue(output[0].contains("jobId="));
}
+ @Test
+ void testSubmitStrategyJobsNormalizesIdentifierWithDefaultCatalog() throws
Exception {
+ Path confPath = addDefaultCatalog(createOptimizerConfForSubmitStrategy(),
"test");
+ String[] output =
+ runCommand(
+ "--type",
+ "submit-strategy-jobs",
+ "--identifiers",
+ "db.table",
+ "--strategy-name",
+ StrategyProviderForCmdTest.STRATEGY_NAME,
+ "--conf-path",
+ confPath.toString());
+ Assertions.assertTrue(output[1].isEmpty(), "stderr=" + output[1] + ",
stdout=" + output[0]);
+ Assertions.assertTrue(output[0].contains("identifier=test.db.table"));
+ }
+
@Test
void testSubmitStrategyJobsDryRunDoesNotSubmit() throws Exception {
Path confPath = createOptimizerConfForSubmitStrategy();
@@ -248,6 +268,38 @@ class TestOptimizerCmd {
Assertions.assertTrue(output[0].contains("row_count=["));
}
+ @Test
+ void testListTableMetricsNormalizesIdentifierWithDefaultCatalog() throws
Exception {
+ Path confPath = addDefaultCatalog(createOptimizerConfForMetricsProvider(),
"test");
+ String[] output =
+ runCommand(
+ "--type",
+ "list-table-metrics",
+ "--identifiers",
+ "db.table",
+ "--conf-path",
+ confPath.toString());
+ Assertions.assertTrue(output[1].isEmpty(), "stderr=" + output[1] + ",
stdout=" + output[0]);
+ Assertions.assertTrue(output[0].contains("identifier=test.db.table"));
+ }
+
+ @Test
+ void testListTableMetricsRejectsTwoLevelIdentifierWithoutDefaultCatalog()
throws Exception {
+ Path confPath = createOptimizerConfForMetricsProvider();
+ String[] output =
+ runCommand(
+ "--type",
+ "list-table-metrics",
+ "--identifiers",
+ "db.table",
+ "--conf-path",
+ confPath.toString());
+ Assertions.assertTrue(
+ output[1].contains(
+ "configure gravitino.optimizer.gravitinoDefaultCatalog when using
schema.table"),
+ "stderr=" + output[1] + ", stdout=" + output[0]);
+ }
+
@Test
void testListTableMetricsWithPartitionPathImplemented() throws Exception {
Path confPath = createOptimizerConfForMetricsProvider();
@@ -288,6 +340,22 @@ class TestOptimizerCmd {
Assertions.assertTrue(output[0].contains("duration=["));
}
+ @Test
+ void testListJobMetricsDoesNotApplyDefaultCatalog() throws Exception {
+ Path confPath = addDefaultCatalog(createOptimizerConfForMetricsProvider(),
"catalog");
+ String[] output =
+ runCommand(
+ "--type",
+ "list-job-metrics",
+ "--identifiers",
+ "db.job1",
+ "--conf-path",
+ confPath.toString());
+ Assertions.assertTrue(output[1].isEmpty(), "stderr=" + output[1] + ",
stdout=" + output[0]);
+ Assertions.assertTrue(output[0].contains("identifier=db.job1"));
+ Assertions.assertFalse(output[0].contains("identifier=catalog.db.job1"));
+ }
+
@Test
void testMonitorMetricsImplemented() throws Exception {
Path confPath = createOptimizerConfForMonitor();
@@ -310,6 +378,25 @@ class TestOptimizerCmd {
Assertions.assertTrue(output[0].contains("identifier=test.db.job2"));
}
+ @Test
+ void testMonitorMetricsNormalizesIdentifierWithDefaultCatalog() throws
Exception {
+ Path confPath = addDefaultCatalog(createOptimizerConfForMonitor(), "test");
+ String[] output =
+ runCommand(
+ "--type",
+ "monitor-metrics",
+ "--identifiers",
+ "db.table",
+ "--action-time",
+ "100",
+ "--range-seconds",
+ "10",
+ "--conf-path",
+ confPath.toString());
+ Assertions.assertTrue(output[1].isEmpty(), "stderr=" + output[1] + ",
stdout=" + output[0]);
+ Assertions.assertTrue(output[0].contains("identifier=test.db.table"));
+ }
+
@Test
void testUpdateStatisticsWithoutIdentifiersUsesUpdateAllPath() throws
Exception {
StatisticsUpdaterForTest.reset();
@@ -340,6 +427,29 @@ class TestOptimizerCmd {
"Expected partition updates from updateAll, but got " +
totalPartitionUpdates);
}
+ @Test
+ void testUpdateStatisticsNormalizesIdentifierWithDefaultCatalog() throws
Exception {
+ StatisticsUpdaterForTest.reset();
+ Path confPath = addDefaultCatalog(createOptimizerConfForUpdater(),
"catalog");
+ String[] output =
+ runCommand(
+ "--type",
+ "update-statistics",
+ "--identifiers",
+ "schema.table",
+ "--calculator-name",
+ StatisticsCalculatorForTest.NAME,
+ "--conf-path",
+ confPath.toString());
+ Assertions.assertTrue(output[1].isEmpty(), "stderr=" + output[1] + ",
stdout=" + output[0]);
+ List<NameIdentifier> updatedIdentifiers =
+ StatisticsUpdaterForTest.instances().stream()
+ .flatMap(updater -> updater.tableIdentifiers().stream())
+ .toList();
+ Assertions.assertEquals(
+ List.of(NameIdentifier.of("catalog", "schema", "table")),
updatedIdentifiers);
+ }
+
@Test
void testAppendMetricsWithoutIdentifiersUsesUpdateAllPath() throws Exception
{
StatisticsUpdaterForTest.reset();
@@ -643,6 +753,15 @@ class TestOptimizerCmd {
return confPath;
}
+ private Path addDefaultCatalog(Path confPath, String catalogName) throws
Exception {
+ Files.writeString(
+ confPath,
+ "gravitino.optimizer.gravitinoDefaultCatalog = " + catalogName +
System.lineSeparator(),
+ StandardCharsets.UTF_8,
+ StandardOpenOption.APPEND);
+ return confPath;
+ }
+
private String[] runCommand(String... args) {
ByteArrayOutputStream outBuffer = new ByteArrayOutputStream();
ByteArrayOutputStream errBuffer = new ByteArrayOutputStream();
diff --git
a/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/updater/StatisticsUpdaterForTest.java
b/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/updater/StatisticsUpdaterForTest.java
index a2e19a3d15..8d9aa1c44e 100644
---
a/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/updater/StatisticsUpdaterForTest.java
+++
b/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/updater/StatisticsUpdaterForTest.java
@@ -35,6 +35,8 @@ public class StatisticsUpdaterForTest implements
StatisticsUpdater {
public static final String NAME = "test-statistics-updater";
private static final List<StatisticsUpdaterForTest> INSTANCES =
Collections.synchronizedList(new ArrayList<>());
+ private final List<NameIdentifier> tableIdentifiers =
+ Collections.synchronizedList(new ArrayList<>());
private final AtomicInteger tableUpdates = new AtomicInteger();
private final AtomicInteger partitionUpdates = new AtomicInteger();
private final AtomicInteger closeCalls = new AtomicInteger();
@@ -55,6 +57,10 @@ public class StatisticsUpdaterForTest implements
StatisticsUpdater {
return tableUpdates.get();
}
+ public List<NameIdentifier> tableIdentifiers() {
+ return new ArrayList<>(tableIdentifiers);
+ }
+
public int partitionUpdates() {
return partitionUpdates.get();
}
@@ -74,6 +80,7 @@ public class StatisticsUpdaterForTest implements
StatisticsUpdater {
@Override
public void updateTableStatistics(
NameIdentifier tableIdentifier, List<StatisticEntry<?>> tableStatistics)
{
+ tableIdentifiers.add(tableIdentifier);
tableUpdates.incrementAndGet();
}