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 f3a41e6128 [#13137] fix(maintenance): authenticate
builtin-iceberg-update-stats against secured servers (#13138)
f3a41e6128 is described below
commit f3a41e6128c2667f70fcb0a169ceb18d4434a8b6
Author: MaSai <[email protected]>
AuthorDate: Thu Sep 17 16:03:37 2026 +0800
[#13137] fix(maintenance): authenticate builtin-iceberg-update-stats
against secured servers (#13138)
### What changes were proposed in this pull request?
Allow `builtin-iceberg-update-stats` to authenticate when Gravitino /
Iceberg REST requires credentials:
- Add `GravitinoAuthSettings` to resolve auth from `--updater-options` /
`updater_options` (`auth_type`, `username`, `password`, or OAuth
client-credentials fields).
- Apply credentials to the Gravitino client used by the statistics
updater and to Spark Iceberg REST catalog configs (`rest.auth.*`). OAuth
joins `oauth_server_uri` + `oauth_path` into Iceberg
`oauth2-server-uri`.
- `submit-update-stats-job` reuses the same `updater_options` auth when
calling Gravitino `runJob` (no separate CLI auth config).
- Expire-snapshots / rewrite-data-files are unchanged; callers that need
Iceberg REST auth should set `rest.auth.*` in `spark-conf`.
- Redact passwords / credentials in `updater_options` and `spark_conf`
in DRY-RUN / SUBMIT log output.
Fix: #13137
### Why are the changes needed?
With authenticators enabled, update-stats fails on unauthenticated
Gravitino API callbacks and against secured Iceberg REST catalogs.
### Does this PR introduce _any_ user-facing change?
Yes. Optional auth fields in `updater_options` for
`builtin-iceberg-update-stats` (`none` / `simple` / `basic` / `oauth`
client-credentials).
### How was this patch tested?
```bash
./gradlew spotlessApply \
:maintenance:optimizer-api:test --tests
'org.apache.gravitino.maintenance.optimizer.common.util.TestGravitinoAuthSettings'
\
:maintenance:jobs:test --tests
'org.apache.gravitino.maintenance.jobs.iceberg.TestIcebergUpdateStatsJob' \
--tests
'org.apache.gravitino.maintenance.jobs.iceberg.TestIcebergJobUtils' \
:maintenance:optimizer:test --tests
'org.apache.gravitino.maintenance.optimizer.command.TestSubmitUpdateStatsJobCommand'
\
-PskipITs -PskipDockerTests=true
```
---------
Co-authored-by: Cursor <[email protected]>
---
.../optimizer-configuration.md | 21 ++-
.../optimizer-troubleshooting.md | 2 +
.../iceberg/IcebergUpdateStatsAndMetricsJob.java | 5 +-
.../jobs/iceberg/TestIcebergUpdateStatsJob.java | 16 ++
maintenance/optimizer-api/build.gradle.kts | 1 +
.../optimizer/common/conf/OptimizerConfig.java | 12 ++
.../common/util/GravitinoAuthSettings.java | 203 +++++++++++++++++++++
.../common/util/GravitinoClientUtils.java | 33 +++-
.../common/util/TestGravitinoAuthSettings.java | 123 +++++++++++++
.../command/SubmitUpdateStatsJobCommand.java | 67 ++++++-
.../command/TestSubmitUpdateStatsJobCommand.java | 57 ++++++
11 files changed, 534 insertions(+), 6 deletions(-)
diff --git a/docs/table-maintenance-service/optimizer-configuration.md
b/docs/table-maintenance-service/optimizer-configuration.md
index cc379da68b..ae6350e9b2 100644
--- a/docs/table-maintenance-service/optimizer-configuration.md
+++ b/docs/table-maintenance-service/optimizer-configuration.md
@@ -53,6 +53,25 @@ gravitino.optimizer.jobSubmitterConfig.warehouse_location =
gravitino.optimizer.jobSubmitterConfig.spark_conf =
{"spark.master":"local[2]","spark.hadoop.fs.defaultFS":"file:///"}
```
+When the Gravitino server has authentication enabled,
`builtin-iceberg-update-stats` needs
+credentials in `--updater-options` / `updater_options` for the Gravitino
client (statistics
+updater and `submit-update-stats-job` `runJob` calls):
+
+| `auth_type` | Fields
|
+|------------------|------------------------------------------------------------------------------------------------|
+| `none` (default) | (none)
|
+| `simple` | `username` (optional)
|
+| `basic` | `username`, `password`
|
+| `oauth` | client-credentials only: `oauth_server_uri`,
`oauth_path`, `oauth_credential`, `oauth_scope` |
+
+Iceberg REST catalog authentication is separate: set `rest.auth.*` in
`spark-conf` for any built-in
+job that talks to a secured IRC (including update-stats). Expire-snapshots and
rewrite-data-files
+do not read `updater_options` auth fields.
+
+Passwords and OAuth credentials in `updater_options` (and secrets in
`spark_conf`) travel with
+the job command line and `jobConf`; avoid logging raw `jobConf` (the
submit-update-stats CLI
+redacts them in DRY-RUN / SUBMIT output).
+
Everything under `gravitino.optimizer.jobSubmitterConfig.` becomes the
`jobConf` of jobs this CLI submits, so the two layers carry the same keys under
different names.
## Job Submission Configuration
@@ -64,7 +83,7 @@ A direct job submission carries its own `jobConf`. This is
`builtin-iceberg-upda
"catalog_name": "rest_catalog",
"table_identifier": "db.t1",
"update_mode": "all",
- "updater_options":
"{\"gravitino_uri\":\"http://localhost:8090\",\"metalake\":\"test\",\"statistics_updater\":\"gravitino-statistics-updater\",\"metrics_updater\":\"gravitino-metrics-updater\"}",
+ "updater_options":
"{\"gravitino_uri\":\"http://localhost:8090\",\"metalake\":\"test\",\"statistics_updater\":\"gravitino-statistics-updater\",\"metrics_updater\":\"gravitino-metrics-updater\",\"auth_type\":\"basic\",\"username\":\"admin\",\"password\":\"YourSecureGravitinoPassword\"}",
"spark_conf":
"{\"spark.master\":\"local[2]\",\"spark.hadoop.fs.defaultFS\":\"file:///\"}",
"spark_master": "local[2]",
"spark_executor_instances": "1",
diff --git a/docs/table-maintenance-service/optimizer-troubleshooting.md
b/docs/table-maintenance-service/optimizer-troubleshooting.md
index e9d382e7e0..099b9b4814 100644
--- a/docs/table-maintenance-service/optimizer-troubleshooting.md
+++ b/docs/table-maintenance-service/optimizer-troubleshooting.md
@@ -72,6 +72,8 @@ spark.hadoop.fs.defaultFS=file:///
}
```
+**`The provided credentials did not support` or Iceberg `Not authorized`** — a
secured Gravitino endpoint needs `auth_type` / `username` / `password` (or
OAuth client-credentials fields) in `updater_options` for
`builtin-iceberg-update-stats`. A secured Iceberg REST catalog needs
`rest.auth.*` in `spark-conf` for any built-in job that reads the table. See
[Configuration](./optimizer-configuration.md).
+
**Built-in Iceberg jobs fail with `Missing Iceberg Spark session extensions`**
—
Spark only warns when `IcebergSparkSessionExtensions` is missing, so built-in
jobs check the
classpath after `SparkSession` starts and exit with a non-zero status when the
Iceberg Spark
diff --git
a/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergUpdateStatsAndMetricsJob.java
b/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergUpdateStatsAndMetricsJob.java
index 7253cb07ec..21db50d6f6 100644
---
a/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergUpdateStatsAndMetricsJob.java
+++
b/maintenance/jobs/src/main/java/org/apache/gravitino/maintenance/jobs/iceberg/IcebergUpdateStatsAndMetricsJob.java
@@ -43,6 +43,7 @@ import
org.apache.gravitino.maintenance.optimizer.common.OptimizerEnv;
import org.apache.gravitino.maintenance.optimizer.common.PartitionEntryImpl;
import org.apache.gravitino.maintenance.optimizer.common.StatisticEntryImpl;
import org.apache.gravitino.maintenance.optimizer.common.conf.OptimizerConfig;
+import
org.apache.gravitino.maintenance.optimizer.common.util.GravitinoAuthSettings;
import
org.apache.gravitino.maintenance.optimizer.common.util.IcebergSparkConfigUtils;
import org.apache.gravitino.maintenance.optimizer.common.util.ProviderUtils;
import org.apache.gravitino.stats.StatisticValues;
@@ -434,6 +435,7 @@ public class IcebergUpdateStatsAndMetricsJob implements
BuiltInJob {
gravitinoUri.ifPresent(uri ->
optimizerProperties.put(OptimizerConfig.GRAVITINO_URI, uri));
metalake.ifPresent(value ->
optimizerProperties.put(OptimizerConfig.GRAVITINO_METALAKE, value));
+ GravitinoAuthSettings.copyAliases(optimizerProperties);
return optimizerProperties;
}
@@ -582,7 +584,8 @@ public class IcebergUpdateStatsAndMetricsJob implements
BuiltInJob {
+ " --updater-options <json> JSON map for updater and
repository settings\\n"
+ " Example:
'{\"gravitino_uri\":\"http://localhost:8090\",\\n"
+ "
\"metalake\":\"test\",\"statistics_updater\":\"gravitino-statistics-updater\",\\n"
- + "
\"metrics_updater\":\"gravitino-metrics-updater\"}'\\n"
+ + "
\"metrics_updater\":\"gravitino-metrics-updater\",\\n"
+ + "
\"auth_type\":\"basic\",\"username\":\"admin\",\"password\":\"YourSecureGravitinoPassword\"}'\\n"
+ " --spark-conf <json> JSON map of custom Spark
configs\\n"
+ " Must include Iceberg
catalog configs for --catalog\\n"
+ " Example:
'{\"spark.master\":\"local[2]\","
diff --git
a/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergUpdateStatsJob.java
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergUpdateStatsJob.java
index 9fd9307e55..2baefa9043 100644
---
a/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergUpdateStatsJob.java
+++
b/maintenance/jobs/src/test/java/org/apache/gravitino/maintenance/jobs/iceberg/TestIcebergUpdateStatsJob.java
@@ -203,6 +203,22 @@ public class TestIcebergUpdateStatsJob {
optimizerProperties.get("gravitino.optimizer.jdbcMetrics.jdbcUrl"));
}
+ @Test
+ public void testBuildOptimizerPropertiesCopiesAuthAliases() {
+ Map<String, String> options = new HashMap<>();
+ options.put("gravitino_uri", "http://localhost:8090");
+ options.put("metalake", "ml");
+ options.put("auth_type", "basic");
+ options.put("username", "admin");
+ options.put("password", "YourSecureGravitinoPassword");
+ Map<String, String> optimizerProperties =
+ IcebergUpdateStatsAndMetricsJob.buildOptimizerProperties(options);
+ assertEquals("basic", optimizerProperties.get(OptimizerConfig.AUTH_TYPE));
+ assertEquals("admin",
optimizerProperties.get(OptimizerConfig.AUTH_USERNAME));
+ assertEquals(
+ "YourSecureGravitinoPassword",
optimizerProperties.get(OptimizerConfig.AUTH_PASSWORD));
+ }
+
@Test
public void testRequireGravitinoConfig() {
Map<String, String> optimizerProperties = new HashMap<>();
diff --git a/maintenance/optimizer-api/build.gradle.kts
b/maintenance/optimizer-api/build.gradle.kts
index 436fa18ab3..07d507dfcd 100644
--- a/maintenance/optimizer-api/build.gradle.kts
+++ b/maintenance/optimizer-api/build.gradle.kts
@@ -38,6 +38,7 @@ dependencies {
testCompileOnly(libs.lombok)
testImplementation(libs.junit.jupiter.api)
+ testImplementation(libs.mockito.core)
testRuntimeOnly(libs.junit.jupiter.engine)
}
diff --git
a/maintenance/optimizer-api/src/main/java/org/apache/gravitino/maintenance/optimizer/common/conf/OptimizerConfig.java
b/maintenance/optimizer-api/src/main/java/org/apache/gravitino/maintenance/optimizer/common/conf/OptimizerConfig.java
index 0cb063f304..bd1fb7b397 100644
---
a/maintenance/optimizer-api/src/main/java/org/apache/gravitino/maintenance/optimizer/common/conf/OptimizerConfig.java
+++
b/maintenance/optimizer-api/src/main/java/org/apache/gravitino/maintenance/optimizer/common/conf/OptimizerConfig.java
@@ -41,6 +41,18 @@ public class OptimizerConfig extends Config {
public static final String GRAVITINO_METALAKE = OPTIMIZER_PREFIX +
"gravitinoMetalake";
public static final String GRAVITINO_DEFAULT_CATALOG =
OPTIMIZER_PREFIX + "gravitinoDefaultCatalog";
+ /**
+ * Canonical auth property keys used when {@code
builtin-iceberg-update-stats} loads {@code
+ * --updater-options} into {@link OptimizerConfig}. Not CLI configuration
entries.
+ */
+ public static final String AUTH_TYPE = OPTIMIZER_PREFIX + "auth.type";
+
+ public static final String AUTH_USERNAME = OPTIMIZER_PREFIX +
"auth.username";
+ public static final String AUTH_PASSWORD = OPTIMIZER_PREFIX +
"auth.password";
+ public static final String AUTH_OAUTH_SERVER_URI = OPTIMIZER_PREFIX +
"auth.oauth.serverUri";
+ public static final String AUTH_OAUTH_PATH = OPTIMIZER_PREFIX +
"auth.oauth.path";
+ public static final String AUTH_OAUTH_CREDENTIAL = OPTIMIZER_PREFIX +
"auth.oauth.credential";
+ public static final String AUTH_OAUTH_SCOPE = OPTIMIZER_PREFIX +
"auth.oauth.scope";
public static final String JOB_ADAPTER_PREFIX = OPTIMIZER_PREFIX +
"jobAdapter.";
public static final String JOB_SUBMITTER_CONFIG_PREFIX = OPTIMIZER_PREFIX +
"jobSubmitterConfig.";
diff --git
a/maintenance/optimizer-api/src/main/java/org/apache/gravitino/maintenance/optimizer/common/util/GravitinoAuthSettings.java
b/maintenance/optimizer-api/src/main/java/org/apache/gravitino/maintenance/optimizer/common/util/GravitinoAuthSettings.java
new file mode 100644
index 0000000000..fb9640f9cc
--- /dev/null
+++
b/maintenance/optimizer-api/src/main/java/org/apache/gravitino/maintenance/optimizer/common/util/GravitinoAuthSettings.java
@@ -0,0 +1,203 @@
+/*
+ * 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.optimizer.common.util;
+
+import java.util.Locale;
+import java.util.Map;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.gravitino.client.DefaultOAuth2TokenProvider;
+import org.apache.gravitino.client.GravitinoClientBase;
+import org.apache.gravitino.maintenance.optimizer.common.conf.OptimizerConfig;
+
+/**
+ * Authentication settings for {@code builtin-iceberg-update-stats}.
+ *
+ * <p>Put credentials in {@code --updater-options} (short names such as {@code
auth_type}). {@link
+ * #copyAliases(Map)} promotes them onto canonical {@code
gravitino.optimizer.auth.*} keys so {@link
+ * #from(OptimizerConfig)} can apply them to the Gravitino client used by the
statistics updater.
+ * Iceberg REST catalog auth is configured separately via {@code spark-conf}
({@code rest.auth.*}).
+ */
+public final class GravitinoAuthSettings {
+
+ public static final String TYPE_NONE = "none";
+ public static final String TYPE_SIMPLE = "simple";
+ public static final String TYPE_BASIC = "basic";
+ public static final String TYPE_OAUTH = "oauth";
+
+ private final String authType;
+ private final String username;
+ private final String password;
+ private final String oauthServerUri;
+ private final String oauthPath;
+ private final String oauthCredential;
+ private final String oauthScope;
+
+ private GravitinoAuthSettings(
+ String authType,
+ String username,
+ String password,
+ String oauthServerUri,
+ String oauthPath,
+ String oauthCredential,
+ String oauthScope) {
+ this.authType = authType;
+ this.username = username;
+ this.password = password;
+ this.oauthServerUri = oauthServerUri;
+ this.oauthPath = oauthPath;
+ this.oauthCredential = oauthCredential;
+ this.oauthScope = oauthScope;
+ }
+
+ /**
+ * Resolves settings from optimizer configuration.
+ *
+ * @param config optimizer configuration, may be {@code null}
+ * @return resolved authentication settings
+ */
+ public static GravitinoAuthSettings from(OptimizerConfig config) {
+ return new GravitinoAuthSettings(
+ configValue(config, OptimizerConfig.AUTH_TYPE),
+ configValue(config, OptimizerConfig.AUTH_USERNAME),
+ configValue(config, OptimizerConfig.AUTH_PASSWORD),
+ configValue(config, OptimizerConfig.AUTH_OAUTH_SERVER_URI),
+ configValue(config, OptimizerConfig.AUTH_OAUTH_PATH),
+ configValue(config, OptimizerConfig.AUTH_OAUTH_CREDENTIAL),
+ configValue(config, OptimizerConfig.AUTH_OAUTH_SCOPE));
+ }
+
+ /**
+ * Copies shorthand updater-option keys onto canonical {@code
gravitino.optimizer.auth.*} keys.
+ *
+ * @param properties updater options / optimizer properties
+ */
+ public static void copyAliases(Map<String, String> properties) {
+ if (properties == null) {
+ return;
+ }
+ copyAlias(properties, OptimizerConfig.AUTH_TYPE, "auth_type");
+ copyAlias(properties, OptimizerConfig.AUTH_USERNAME, "username");
+ copyAlias(properties, OptimizerConfig.AUTH_PASSWORD, "password");
+ copyAlias(properties, OptimizerConfig.AUTH_OAUTH_SERVER_URI,
"oauth_server_uri");
+ copyAlias(properties, OptimizerConfig.AUTH_OAUTH_PATH, "oauth_path");
+ copyAlias(properties, OptimizerConfig.AUTH_OAUTH_CREDENTIAL,
"oauth_credential");
+ copyAlias(properties, OptimizerConfig.AUTH_OAUTH_SCOPE, "oauth_scope");
+ }
+
+ /**
+ * Applies the matching authenticator to a Gravitino client builder.
+ *
+ * @param builder Gravitino client builder
+ * @param <T> client type
+ * @return the same builder
+ */
+ public <T extends GravitinoClientBase> GravitinoClientBase.Builder<T>
applyTo(
+ GravitinoClientBase.Builder<T> builder) {
+ switch (normalizedType()) {
+ case TYPE_SIMPLE:
+ if (StringUtils.isNotBlank(username)) {
+ builder.withSimpleAuth(username);
+ } else {
+ builder.withSimpleAuth();
+ }
+ break;
+ case TYPE_BASIC:
+ requireValue(username, "username is required for basic
authentication");
+ requireValue(password, "password is required for basic
authentication");
+ builder.withBasicAuth(username, password);
+ break;
+ case TYPE_OAUTH:
+ requireValue(oauthServerUri, "oauth server URI is required for oauth
authentication");
+ requireValue(oauthPath, "oauth path is required for oauth
authentication");
+ requireValue(oauthCredential, "oauth credential is required for oauth
authentication");
+ requireValue(oauthScope, "oauth scope is required for oauth
authentication");
+ builder.withOAuth(
+ DefaultOAuth2TokenProvider.builder()
+ .withUri(oauthServerUri)
+ .withPath(oauthPath)
+ .withCredential(oauthCredential)
+ .withScope(oauthScope)
+ .build());
+ break;
+ case TYPE_NONE:
+ break;
+ default:
+ throw new IllegalArgumentException(
+ "Unsupported auth_type: "
+ + authType
+ + ". Supported values are: none, simple, basic, oauth");
+ }
+ return builder;
+ }
+
+ /**
+ * Returns whether an authenticator should be configured.
+ *
+ * @return true when auth type is simple, basic, or oauth
+ */
+ public boolean hasAuth() {
+ String type = normalizedType();
+ return TYPE_SIMPLE.equals(type) || TYPE_BASIC.equals(type) ||
TYPE_OAUTH.equals(type);
+ }
+
+ /**
+ * Returns the normalized auth type, or {@link #TYPE_NONE} when unset.
+ *
+ * @return auth type
+ */
+ public String authType() {
+ return normalizedType();
+ }
+
+ private String normalizedType() {
+ if (StringUtils.isBlank(authType)) {
+ return TYPE_NONE;
+ }
+ return authType.trim().toLowerCase(Locale.ROOT);
+ }
+
+ private static String configValue(OptimizerConfig config, String key) {
+ if (config == null) {
+ return null;
+ }
+ String value = config.getRawString(key);
+ return StringUtils.isBlank(value) ? null : value.trim();
+ }
+
+ private static void copyAlias(
+ Map<String, String> properties, String canonical, String... aliases) {
+ if (StringUtils.isNotBlank(properties.get(canonical))) {
+ return;
+ }
+ for (String alias : aliases) {
+ String value = properties.get(alias);
+ if (StringUtils.isNotBlank(value)) {
+ properties.put(canonical, value.trim());
+ return;
+ }
+ }
+ }
+
+ private static void requireValue(String value, String message) {
+ if (StringUtils.isBlank(value)) {
+ throw new IllegalArgumentException(message);
+ }
+ }
+}
diff --git
a/maintenance/optimizer-api/src/main/java/org/apache/gravitino/maintenance/optimizer/common/util/GravitinoClientUtils.java
b/maintenance/optimizer-api/src/main/java/org/apache/gravitino/maintenance/optimizer/common/util/GravitinoClientUtils.java
index 052d7b2d72..9154f8fc80 100644
---
a/maintenance/optimizer-api/src/main/java/org/apache/gravitino/maintenance/optimizer/common/util/GravitinoClientUtils.java
+++
b/maintenance/optimizer-api/src/main/java/org/apache/gravitino/maintenance/optimizer/common/util/GravitinoClientUtils.java
@@ -20,6 +20,9 @@
package org.apache.gravitino.maintenance.optimizer.common.util;
import com.google.common.base.Preconditions;
+import java.util.HashMap;
+import java.util.Map;
+import javax.annotation.Nullable;
import org.apache.gravitino.client.GravitinoClient;
import org.apache.gravitino.maintenance.optimizer.common.OptimizerEnv;
import org.apache.gravitino.maintenance.optimizer.common.conf.OptimizerConfig;
@@ -32,14 +35,42 @@ public final class GravitinoClientUtils {
/**
* Creates a {@link GravitinoClient} using optimizer configuration.
*
+ * <p>When the config carries auth keys (for example from update-stats
{@code updater-options}
+ * loaded into {@link OptimizerConfig}), those credentials are applied to
the client builder.
+ *
* @param optimizerEnv optimizer environment
* @return configured Gravitino client
*/
public static GravitinoClient createClient(OptimizerEnv optimizerEnv) {
+ return createClient(optimizerEnv, null);
+ }
+
+ /**
+ * Creates a {@link GravitinoClient} using URI/metalake from {@code
optimizerEnv} and optional
+ * auth from {@code authFromUpdaterOptions} (short names such as {@code
auth_type}).
+ *
+ * <p>Used by {@code submit-update-stats-job} so the CLI can call {@code
runJob} against a secured
+ * server with the same credentials that the job will use.
+ *
+ * @param optimizerEnv optimizer environment (URI and metalake)
+ * @param authFromUpdaterOptions updater-options map, may be {@code null}
+ * @return configured Gravitino client
+ */
+ public static GravitinoClient createClient(
+ OptimizerEnv optimizerEnv, @Nullable Map<String, String>
authFromUpdaterOptions) {
Preconditions.checkArgument(optimizerEnv != null, "optimizerEnv must not
be null");
OptimizerConfig config = optimizerEnv.config();
String uri = config.get(OptimizerConfig.GRAVITINO_URI_CONFIG);
String metalake = config.get(OptimizerConfig.GRAVITINO_METALAKE_CONFIG);
- return GravitinoClient.builder(uri).withMetalake(metalake).build();
+ GravitinoClient.ClientBuilder builder =
GravitinoClient.builder(uri).withMetalake(metalake);
+
+ OptimizerConfig authConfig = config;
+ if (authFromUpdaterOptions != null && !authFromUpdaterOptions.isEmpty()) {
+ Map<String, String> authProperties = new
HashMap<>(authFromUpdaterOptions);
+ GravitinoAuthSettings.copyAliases(authProperties);
+ authConfig = new OptimizerConfig(authProperties);
+ }
+ GravitinoAuthSettings.from(authConfig).applyTo(builder);
+ return builder.build();
}
}
diff --git
a/maintenance/optimizer-api/src/test/java/org/apache/gravitino/maintenance/optimizer/common/util/TestGravitinoAuthSettings.java
b/maintenance/optimizer-api/src/test/java/org/apache/gravitino/maintenance/optimizer/common/util/TestGravitinoAuthSettings.java
new file mode 100644
index 0000000000..b6cc121923
--- /dev/null
+++
b/maintenance/optimizer-api/src/test/java/org/apache/gravitino/maintenance/optimizer/common/util/TestGravitinoAuthSettings.java
@@ -0,0 +1,123 @@
+/*
+ * 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.optimizer.common.util;
+
+import java.util.HashMap;
+import java.util.Map;
+import org.apache.gravitino.client.GravitinoClient;
+import org.apache.gravitino.maintenance.optimizer.common.conf.OptimizerConfig;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+class TestGravitinoAuthSettings {
+
+ @Test
+ void fromReturnsNoneWhenUnset() {
+ GravitinoAuthSettings settings = GravitinoAuthSettings.from(new
OptimizerConfig());
+ Assertions.assertEquals(GravitinoAuthSettings.TYPE_NONE,
settings.authType());
+ Assertions.assertFalse(settings.hasAuth());
+ }
+
+ @Test
+ void applyToUsesBasicAuth() {
+ Map<String, String> properties = new HashMap<>();
+ properties.put(OptimizerConfig.AUTH_TYPE, "basic");
+ properties.put(OptimizerConfig.AUTH_USERNAME, "admin");
+ properties.put(OptimizerConfig.AUTH_PASSWORD,
"YourSecureGravitinoPassword");
+ GravitinoAuthSettings settings = GravitinoAuthSettings.from(new
OptimizerConfig(properties));
+
+ @SuppressWarnings("unchecked")
+ GravitinoClient.ClientBuilder builder =
Mockito.mock(GravitinoClient.ClientBuilder.class);
+ Mockito.when(builder.withBasicAuth(Mockito.anyString(),
Mockito.anyString()))
+ .thenReturn(builder);
+
+ settings.applyTo(builder);
+ Mockito.verify(builder).withBasicAuth("admin",
"YourSecureGravitinoPassword");
+ }
+
+ @Test
+ void applyToUsesSimpleAuth() {
+ Map<String, String> properties = new HashMap<>();
+ properties.put(OptimizerConfig.AUTH_TYPE, "simple");
+ properties.put(OptimizerConfig.AUTH_USERNAME, "alice");
+ GravitinoAuthSettings settings = GravitinoAuthSettings.from(new
OptimizerConfig(properties));
+
+ @SuppressWarnings("unchecked")
+ GravitinoClient.ClientBuilder builder =
Mockito.mock(GravitinoClient.ClientBuilder.class);
+
Mockito.when(builder.withSimpleAuth(Mockito.anyString())).thenReturn(builder);
+
+ settings.applyTo(builder);
+ Mockito.verify(builder).withSimpleAuth("alice");
+ }
+
+ @Test
+ void oauthRequiresPath() {
+ Map<String, String> properties = new HashMap<>();
+ properties.put(OptimizerConfig.AUTH_TYPE, "oauth");
+ properties.put(OptimizerConfig.AUTH_OAUTH_SERVER_URI, "http://idp");
+ properties.put(OptimizerConfig.AUTH_OAUTH_CREDENTIAL, "id:secret");
+ properties.put(OptimizerConfig.AUTH_OAUTH_SCOPE, "catalog");
+ GravitinoAuthSettings settings = GravitinoAuthSettings.from(new
OptimizerConfig(properties));
+
+ @SuppressWarnings("unchecked")
+ GravitinoClient.ClientBuilder builder =
Mockito.mock(GravitinoClient.ClientBuilder.class);
+ Assertions.assertThrows(IllegalArgumentException.class, () ->
settings.applyTo(builder));
+ }
+
+ @Test
+ void oauthRequiresCredential() {
+ Map<String, String> properties = new HashMap<>();
+ properties.put(OptimizerConfig.AUTH_TYPE, "oauth");
+ properties.put(OptimizerConfig.AUTH_OAUTH_SERVER_URI, "http://idp");
+ properties.put(OptimizerConfig.AUTH_OAUTH_PATH, "oauth2/token");
+ properties.put(OptimizerConfig.AUTH_OAUTH_SCOPE, "catalog");
+ GravitinoAuthSettings settings = GravitinoAuthSettings.from(new
OptimizerConfig(properties));
+
+ @SuppressWarnings("unchecked")
+ GravitinoClient.ClientBuilder builder =
Mockito.mock(GravitinoClient.ClientBuilder.class);
+ Assertions.assertThrows(IllegalArgumentException.class, () ->
settings.applyTo(builder));
+ }
+
+ @Test
+ void basicAuthRequiresPassword() {
+ Map<String, String> properties = new HashMap<>();
+ properties.put(OptimizerConfig.AUTH_TYPE, "basic");
+ properties.put(OptimizerConfig.AUTH_USERNAME, "admin");
+ GravitinoAuthSettings settings = GravitinoAuthSettings.from(new
OptimizerConfig(properties));
+
+ @SuppressWarnings("unchecked")
+ GravitinoClient.ClientBuilder builder =
Mockito.mock(GravitinoClient.ClientBuilder.class);
+ Assertions.assertThrows(IllegalArgumentException.class, () ->
settings.applyTo(builder));
+ }
+
+ @Test
+ void copyAliasesPromotesShortNames() {
+ Map<String, String> properties = new HashMap<>();
+ properties.put("auth_type", "basic");
+ properties.put("username", "admin");
+ properties.put("password", "YourSecureGravitinoPassword");
+ GravitinoAuthSettings.copyAliases(properties);
+ Assertions.assertEquals("basic",
properties.get(OptimizerConfig.AUTH_TYPE));
+ Assertions.assertEquals("admin",
properties.get(OptimizerConfig.AUTH_USERNAME));
+ Assertions.assertEquals(
+ "YourSecureGravitinoPassword",
properties.get(OptimizerConfig.AUTH_PASSWORD));
+ }
+}
diff --git
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/SubmitUpdateStatsJobCommand.java
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/SubmitUpdateStatsJobCommand.java
index 581ef9666a..1f08d8dfd6 100644
---
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/SubmitUpdateStatsJobCommand.java
+++
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/command/SubmitUpdateStatsJobCommand.java
@@ -19,13 +19,16 @@
package org.apache.gravitino.maintenance.optimizer.command;
+import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.common.base.Preconditions;
import java.util.ArrayList;
+import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
+import java.util.Set;
import java.util.TreeMap;
import org.apache.commons.lang3.StringUtils;
import org.apache.gravitino.client.GravitinoClient;
@@ -56,6 +59,15 @@ public class SubmitUpdateStatsJobCommand implements
OptimizerCommandExecutor {
private static final String SPARK_DRIVER_MEMORY_KEY = "spark.driver.memory";
private static final String SPARK_SQL_CATALOG_PREFIX = "spark.sql.catalog.";
private static final String ICEBERG_SPARK_CATALOG_IMPL =
"org.apache.iceberg.spark.SparkCatalog";
+ private static final String REDACTED = "******";
+ private static final Set<String> SENSITIVE_UPDATER_OPTION_KEYS =
+ new HashSet<>(
+ List.of(
+ "password",
+ "oauth_credential",
+ "oauth_token",
+ OptimizerConfig.AUTH_PASSWORD,
+ OptimizerConfig.AUTH_OAUTH_CREDENTIAL));
private static final ObjectMapper MAPPER = new ObjectMapper();
@@ -95,7 +107,7 @@ public class SubmitUpdateStatsJobCommand implements
OptimizerCommandExecutor {
.output()
.printf(
"DRY-RUN: identifier=%s jobTemplate=%s jobConfig=%s%n",
- tableTarget.fullIdentifier, JOB_TEMPLATE_NAME, jobConfig);
+ tableTarget.fullIdentifier, JOB_TEMPLATE_NAME,
redactJobConfigForLog(jobConfig));
}
context
.output()
@@ -103,7 +115,8 @@ public class SubmitUpdateStatsJobCommand implements
OptimizerCommandExecutor {
return;
}
- try (GravitinoClient client =
GravitinoClientUtils.createClient(context.optimizerEnv())) {
+ try (GravitinoClient client =
+ GravitinoClientUtils.createClient(context.optimizerEnv(),
updaterOptions)) {
int submitted = 0;
for (TableTarget tableTarget : tableTargets) {
Map<String, String> jobConfig =
@@ -114,7 +127,10 @@ public class SubmitUpdateStatsJobCommand implements
OptimizerCommandExecutor {
.output()
.printf(
"SUBMIT: identifier=%s jobTemplate=%s jobId=%s jobConfig=%s%n",
- tableTarget.fullIdentifier, JOB_TEMPLATE_NAME,
jobHandle.jobId(), jobConfig);
+ tableTarget.fullIdentifier,
+ JOB_TEMPLATE_NAME,
+ jobHandle.jobId(),
+ redactJobConfigForLog(jobConfig));
}
context
.output()
@@ -141,6 +157,51 @@ public class SubmitUpdateStatsJobCommand implements
OptimizerCommandExecutor {
return jobConfig;
}
+ /**
+ * Returns a copy of {@code jobConfig} safe for logging. Sensitive fields
inside {@code
+ * updater_options} and {@code spark_conf} JSON maps are replaced with
{@code ***}.
+ */
+ static Map<String, String> redactJobConfigForLog(Map<String, String>
jobConfig) {
+ if (jobConfig == null || jobConfig.isEmpty()) {
+ return jobConfig;
+ }
+ Map<String, String> redacted = new LinkedHashMap<>(jobConfig);
+ redactJsonMapField(redacted, "updater_options");
+ redactJsonMapField(redacted, "spark_conf");
+ return redacted;
+ }
+
+ private static void redactJsonMapField(Map<String, String> jobConfig, String
fieldName) {
+ String json = jobConfig.get(fieldName);
+ if (StringUtils.isBlank(json)) {
+ return;
+ }
+ try {
+ Map<String, String> parsed =
+ MAPPER.readValue(json, new TypeReference<Map<String, String>>() {});
+ Map<String, String> safe = new LinkedHashMap<>();
+ for (Map.Entry<String, String> entry : parsed.entrySet()) {
+ if (isSensitiveConfigKey(entry.getKey())) {
+ safe.put(entry.getKey(), REDACTED);
+ } else {
+ safe.put(entry.getKey(), entry.getValue());
+ }
+ }
+ jobConfig.put(fieldName, toCanonicalJson(safe));
+ } catch (Exception ignored) {
+ // Keep the original value if JSON cannot be parsed for display.
+ }
+ }
+
+ private static boolean isSensitiveConfigKey(String key) {
+ if (StringUtils.isBlank(key)) {
+ return false;
+ }
+ String normalized = key.trim().toLowerCase(Locale.ROOT);
+ return SENSITIVE_UPDATER_OPTION_KEYS.contains(key)
+ || SENSITIVE_UPDATER_OPTION_KEYS.contains(normalized);
+ }
+
private static String resolveScalarOption(String cliValue, String confValue)
{
if (StringUtils.isNotBlank(cliValue)) {
return cliValue.trim();
diff --git
a/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/command/TestSubmitUpdateStatsJobCommand.java
b/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/command/TestSubmitUpdateStatsJobCommand.java
new file mode 100644
index 0000000000..879bb309fa
--- /dev/null
+++
b/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/command/TestSubmitUpdateStatsJobCommand.java
@@ -0,0 +1,57 @@
+/*
+ * 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.optimizer.command;
+
+import java.util.LinkedHashMap;
+import java.util.Map;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+class TestSubmitUpdateStatsJobCommand {
+
+ @Test
+ void redactJobConfigForLogMasksSensitiveUpdaterOptions() {
+ Map<String, String> jobConfig = new LinkedHashMap<>();
+ jobConfig.put("catalog_name", "rest");
+ jobConfig.put(
+ "updater_options",
+ "{\"gravitino_uri\":\"http://localhost:8090\",\"metalake\":\"test\","
+ + "\"auth_type\":\"basic\",\"username\":\"admin\","
+ + "\"password\":\"YourSecureGravitinoPassword\","
+ + "\"oauth_credential\":\"id:secret\"}");
+ jobConfig.put(
+ "spark_conf",
+ "{\"spark.sql.catalog.rest.rest.auth.type\":\"basic\","
+ +
"\"spark.sql.catalog.rest.rest.auth.basic.password\":\"spark-secret\"}");
+
+ Map<String, String> redacted =
SubmitUpdateStatsJobCommand.redactJobConfigForLog(jobConfig);
+ String updaterOptions = redacted.get("updater_options");
+ Assertions.assertTrue(updaterOptions.contains("\"password\":\"******\""));
+
Assertions.assertTrue(updaterOptions.contains("\"oauth_credential\":\"******\""));
+ Assertions.assertTrue(updaterOptions.contains("\"username\":\"admin\""));
+ Assertions.assertTrue(
+
jobConfig.get("updater_options").contains("\"password\":\"YourSecureGravitinoPassword\""));
+ // spark_conf keys are catalog-scoped (e.g. rest.auth.basic.password);
only exact sensitive
+ // updater-option keys are redacted.
+ Assertions.assertTrue(
+ redacted
+ .get("spark_conf")
+
.contains("\"spark.sql.catalog.rest.rest.auth.basic.password\":\"spark-secret\""));
+ }
+}