This is an automated email from the ASF dual-hosted git repository.
jerryshao pushed a commit to branch branch-1.3
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/branch-1.3 by this push:
new 17d502a9db [Cherry-pick to branch-1.3] [#11601] fix(flink-connector):
Propagate auth to Iceberg REST catalog backend (#11628) (#11665)
17d502a9db is described below
commit 17d502a9db22f209af5c8e7f39cd0118cd342201
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Mon Jun 15 21:50:12 2026 -0700
[Cherry-pick to branch-1.3] [#11601] fix(flink-connector): Propagate auth
to Iceberg REST catalog backend (#11628) (#11665)
**Cherry-pick Information:**
- Original commit: 359778b784f3cf1c1a0ef1c4839d6ecf70c97139
- Target branch: `branch-1.3`
- Status: ✅ Clean cherry-pick (no conflicts)
Co-authored-by: Yuhui <[email protected]>
---
.../connector/catalog/GravitinoCatalogManager.java | 31 +++++++---
.../iceberg/GravitinoIcebergCatalogFactory.java | 11 ++++
.../iceberg/IcebergPropertiesConverter.java | 58 ++++++++++++++++++
.../iceberg/TestIcebergPropertiesConverter.java | 70 ++++++++++++++++++++++
4 files changed, 163 insertions(+), 7 deletions(-)
diff --git
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/catalog/GravitinoCatalogManager.java
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/catalog/GravitinoCatalogManager.java
index 6dd927e961..3c2ecad540 100644
---
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/catalog/GravitinoCatalogManager.java
+++
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/catalog/GravitinoCatalogManager.java
@@ -20,6 +20,7 @@ package org.apache.gravitino.flink.connector.catalog;
import com.google.common.base.Preconditions;
import com.google.common.base.Strings;
+import com.google.common.collect.Maps;
import com.google.common.collect.Sets;
import java.security.PrivilegedAction;
import java.util.Arrays;
@@ -152,6 +153,17 @@ public class GravitinoCatalogManager {
}
}
+ /**
+ * Get the Gravitino client config, including the authentication entries.
This is used to
+ * propagate authentication to the Iceberg REST catalog backend, which
connects directly to the
+ * Iceberg REST service and therefore needs its own credentials.
+ *
+ * @return The Gravitino client config
+ */
+ public Map<String, String> getGravitinoClientConfig() {
+ return gravitinoClientConfig;
+ }
+
/**
* Get GravitinoCatalog by name.
*
@@ -261,17 +273,20 @@ public class GravitinoCatalogManager {
"Basic password is required. Please set %s",
GravitinoCatalogStoreFactoryOptions.BASIC_PASSWORD);
+ // Strip auth keys from a copy so the original config (kept in
gravitinoClientConfig) stays
+ // intact for propagating authentication to the Iceberg REST catalog
backend.
+ Map<String, String> clientConfig = Maps.newHashMap(config);
Set<String> basicConfigKeys =
Sets.newHashSet(
GravitinoCatalogStoreFactoryOptions.AUTH_TYPE,
GravitinoCatalogStoreFactoryOptions.BASIC_USERNAME,
GravitinoCatalogStoreFactoryOptions.BASIC_PASSWORD);
for (String key : basicConfigKeys) {
- config.remove(key);
+ clientConfig.remove(key);
}
return GravitinoAdminClient.builder(gravitinoUri)
.withBasicAuth(username, password)
- .withClientConfig(config)
+ .withClientConfig(clientConfig)
.build();
}
@@ -282,10 +297,12 @@ public class GravitinoCatalogManager {
String path =
config.get(GravitinoCatalogStoreFactoryOptions.OAUTH2_TOKEN_PATH);
String scope =
config.get(GravitinoCatalogStoreFactoryOptions.OAUTH2_SCOPE);
- // Remove OAuth-specific config entries from the client config map. These
keys are only
- // used to construct the OAuth2 token provider and are not valid
GravitinoAdminClient
+ // Remove OAuth-specific config entries from a copy of the client config
map. These keys are
+ // only used to construct the OAuth2 token provider and are not valid
GravitinoAdminClient
// client configuration options; passing them to withClientConfig() could
cause validation
- // errors or other unexpected behavior.
+ // errors or other unexpected behavior. The original config (kept in
gravitinoClientConfig)
+ // stays intact for propagating authentication to the Iceberg REST catalog
backend.
+ Map<String, String> clientConfig = Maps.newHashMap(config);
Set<String> oauthConfigKeys =
Sets.newHashSet(
GravitinoCatalogStoreFactoryOptions.AUTH_TYPE,
@@ -294,7 +311,7 @@ public class GravitinoCatalogManager {
GravitinoCatalogStoreFactoryOptions.OAUTH2_TOKEN_PATH,
GravitinoCatalogStoreFactoryOptions.OAUTH2_SCOPE);
for (String key : oauthConfigKeys) {
- config.remove(key);
+ clientConfig.remove(key);
}
DefaultOAuth2TokenProvider provider =
@@ -307,7 +324,7 @@ public class GravitinoCatalogManager {
return GravitinoAdminClient.builder(gravitinoUri)
.withOAuth(provider)
- .withClientConfig(config)
+ .withClientConfig(clientConfig)
.build();
}
diff --git
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/iceberg/GravitinoIcebergCatalogFactory.java
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/iceberg/GravitinoIcebergCatalogFactory.java
index 3337b469d2..b39ef2e53c 100644
---
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/iceberg/GravitinoIcebergCatalogFactory.java
+++
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/iceberg/GravitinoIcebergCatalogFactory.java
@@ -32,7 +32,9 @@ import
org.apache.gravitino.flink.connector.DefaultPartitionConverter;
import org.apache.gravitino.flink.connector.PartitionConverter;
import org.apache.gravitino.flink.connector.SchemaAndTablePropertiesConverter;
import org.apache.gravitino.flink.connector.catalog.BaseCatalogFactory;
+import org.apache.gravitino.flink.connector.catalog.GravitinoCatalogManager;
import org.apache.gravitino.flink.connector.utils.FactoryUtils;
+import org.apache.iceberg.rest.auth.AuthProperties;
public class GravitinoIcebergCatalogFactory implements BaseCatalogFactory {
@@ -127,6 +129,15 @@ public class GravitinoIcebergCatalogFactory implements
BaseCatalogFactory {
&&
!icebergCatalogOptions.containsKey(IcebergPropertiesConstants.ICEBERG_CATALOG_IMPL))
{
icebergCatalogOptions.put(IcebergPropertiesConstants.ICEBERG_CATALOG_TYPE,
catalogBackend);
}
+ // A REST backend connects directly to the Iceberg REST service, bypassing
the Gravitino
+ // server's auth proxy, so propagate the Gravitino client's authentication
to the REST client,
+ // unless the user has already configured REST auth explicitly.
+ if
(IcebergPropertiesConstants.ICEBERG_CATALOG_BACKEND_REST.equalsIgnoreCase(catalogBackend)
+ && !icebergCatalogOptions.containsKey(AuthProperties.AUTH_TYPE)) {
+ icebergCatalogOptions.putAll(
+ IcebergPropertiesConverter.INSTANCE.toRestAuthProperties(
+ GravitinoCatalogManager.get().getGravitinoClientConfig()));
+ }
// Iceberg's FlinkCatalogFactory only accepts hive/hadoop/rest as
`catalog-type`; a JDBC backend
// must be loaded through `catalog-impl` instead. The two keys are
mutually exclusive, so drop
// `catalog-type` and use `putIfAbsent` to respect an explicitly provided
`catalog-impl`.
diff --git
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/iceberg/IcebergPropertiesConverter.java
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/iceberg/IcebergPropertiesConverter.java
index 8bc75b3b2d..74d9998b9f 100644
---
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/iceberg/IcebergPropertiesConverter.java
+++
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/iceberg/IcebergPropertiesConverter.java
@@ -20,12 +20,17 @@
package org.apache.gravitino.flink.connector.iceberg;
import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.Maps;
import java.util.HashMap;
import java.util.Map;
+import org.apache.commons.lang3.StringUtils;
import org.apache.gravitino.catalog.lakehouse.iceberg.IcebergConstants;
import org.apache.gravitino.catalog.lakehouse.iceberg.IcebergPropertiesUtils;
import org.apache.gravitino.flink.connector.CatalogPropertiesConverter;
import org.apache.gravitino.flink.connector.SchemaAndTablePropertiesConverter;
+import
org.apache.gravitino.flink.connector.store.GravitinoCatalogStoreFactoryOptions;
+import org.apache.iceberg.rest.auth.AuthProperties;
+import org.apache.iceberg.rest.auth.OAuth2Properties;
public class IcebergPropertiesConverter
implements CatalogPropertiesConverter, SchemaAndTablePropertiesConverter {
@@ -37,6 +42,16 @@ public class IcebergPropertiesConverter
ImmutableMap.of(
IcebergConstants.CATALOG_BACKEND,
IcebergPropertiesConstants.ICEBERG_CATALOG_TYPE);
+ // Simple key renames from Gravitino client auth config to Iceberg REST
client auth properties.
+ // The auth type value, and the OAuth2 token endpoint (built from server uri
+ token path), need
+ // special handling and are not listed.
+ private static final Map<String, String> GRAVITINO_AUTH_TO_ICEBERG_REST =
+ ImmutableMap.of(
+ GravitinoCatalogStoreFactoryOptions.OAUTH2_CREDENTIAL,
OAuth2Properties.CREDENTIAL,
+ GravitinoCatalogStoreFactoryOptions.OAUTH2_SCOPE,
OAuth2Properties.SCOPE,
+ GravitinoCatalogStoreFactoryOptions.BASIC_USERNAME,
AuthProperties.BASIC_USERNAME,
+ GravitinoCatalogStoreFactoryOptions.BASIC_PASSWORD,
AuthProperties.BASIC_PASSWORD);
+
@Override
public String transformPropertyToGravitinoCatalog(String configKey) {
return
IcebergPropertiesUtils.ICEBERG_CATALOG_CONFIG_TO_GRAVITINO.get(configKey);
@@ -64,4 +79,47 @@ public class IcebergPropertiesConverter
public String getFlinkCatalogType() {
return GravitinoIcebergCatalogFactoryOptions.IDENTIFIER;
}
+
+ /**
+ * Maps the Gravitino client's authentication config to the Iceberg REST
client's auth properties,
+ * so the Iceberg REST catalog can authenticate against the configured
Iceberg REST service.
+ *
+ * @param gravitinoClientConfig the Gravitino client config carrying the
authentication entries
+ * @return Iceberg REST auth properties, empty when no supported auth is
configured
+ */
+ public Map<String, String> toRestAuthProperties(Map<String, String>
gravitinoClientConfig) {
+ Map<String, String> authProperties = Maps.newHashMap();
+ String authType =
gravitinoClientConfig.get(GravitinoCatalogStoreFactoryOptions.AUTH_TYPE);
+ if (GravitinoCatalogStoreFactoryOptions.OAUTH2.equalsIgnoreCase(authType))
{
+ authProperties.put(AuthProperties.AUTH_TYPE,
AuthProperties.AUTH_TYPE_OAUTH2);
+ // Iceberg expects a single token endpoint URI, while Gravitino splits
it into server uri and
+ // token path, so they are joined here with a single separating slash.
+ String serverUri =
+
gravitinoClientConfig.get(GravitinoCatalogStoreFactoryOptions.OAUTH2_SERVER_URI);
+ String tokenPath =
+
gravitinoClientConfig.get(GravitinoCatalogStoreFactoryOptions.OAUTH2_TOKEN_PATH);
+ if (StringUtils.isNotBlank(serverUri)) {
+ authProperties.put(
+ OAuth2Properties.OAUTH2_SERVER_URI,
+ StringUtils.isNotBlank(tokenPath)
+ ? StringUtils.stripEnd(serverUri, "/")
+ + "/"
+ + StringUtils.stripStart(tokenPath, "/")
+ : serverUri);
+ }
+ } else if
(GravitinoCatalogStoreFactoryOptions.BASIC.equalsIgnoreCase(authType)) {
+ authProperties.put(AuthProperties.AUTH_TYPE,
AuthProperties.AUTH_TYPE_BASIC);
+ } else {
+ // No auth or an auth type not applicable to the Iceberg REST client
(e.g. kerberos).
+ return authProperties;
+ }
+ GRAVITINO_AUTH_TO_ICEBERG_REST.forEach(
+ (gravitinoKey, icebergKey) -> {
+ String value = gravitinoClientConfig.get(gravitinoKey);
+ if (StringUtils.isNotBlank(value)) {
+ authProperties.put(icebergKey, value);
+ }
+ });
+ return authProperties;
+ }
}
diff --git
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/iceberg/TestIcebergPropertiesConverter.java
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/iceberg/TestIcebergPropertiesConverter.java
index 8287eebf25..3ffc584845 100644
---
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/iceberg/TestIcebergPropertiesConverter.java
+++
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/iceberg/TestIcebergPropertiesConverter.java
@@ -23,6 +23,9 @@ import com.google.common.collect.ImmutableMap;
import java.util.Map;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.table.catalog.CommonCatalogOptions;
+import
org.apache.gravitino.flink.connector.store.GravitinoCatalogStoreFactoryOptions;
+import org.apache.iceberg.rest.auth.AuthProperties;
+import org.apache.iceberg.rest.auth.OAuth2Properties;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
@@ -130,4 +133,71 @@ public class TestIcebergPropertiesConverter {
"value2"),
toFlinkProperties);
}
+
+ @Test
+ void testRestAuthPropertiesOAuth2() {
+ Map<String, String> authProperties =
+ CONVERTER.toRestAuthProperties(
+ ImmutableMap.of(
+ GravitinoCatalogStoreFactoryOptions.AUTH_TYPE,
+ GravitinoCatalogStoreFactoryOptions.OAUTH2,
+ GravitinoCatalogStoreFactoryOptions.OAUTH2_SERVER_URI,
+ "http://oauth-server",
+ GravitinoCatalogStoreFactoryOptions.OAUTH2_TOKEN_PATH,
+ "/token",
+ GravitinoCatalogStoreFactoryOptions.OAUTH2_CREDENTIAL,
+ "client:secret",
+ GravitinoCatalogStoreFactoryOptions.OAUTH2_SCOPE,
+ "catalog"));
+
+ Assertions.assertEquals(
+ AuthProperties.AUTH_TYPE_OAUTH2,
authProperties.get(AuthProperties.AUTH_TYPE));
+ Assertions.assertEquals(
+ "http://oauth-server/token",
authProperties.get(OAuth2Properties.OAUTH2_SERVER_URI));
+ Assertions.assertEquals("client:secret",
authProperties.get(OAuth2Properties.CREDENTIAL));
+ Assertions.assertEquals("catalog",
authProperties.get(OAuth2Properties.SCOPE));
+ }
+
+ @Test
+ void testRestAuthPropertiesOAuth2JoinsTokenEndpointSlashes() {
+ Map<String, String> authProperties =
+ CONVERTER.toRestAuthProperties(
+ ImmutableMap.of(
+ GravitinoCatalogStoreFactoryOptions.AUTH_TYPE,
+ GravitinoCatalogStoreFactoryOptions.OAUTH2,
+ GravitinoCatalogStoreFactoryOptions.OAUTH2_SERVER_URI,
+ "http://oauth-server/",
+ GravitinoCatalogStoreFactoryOptions.OAUTH2_TOKEN_PATH,
+ "/token",
+ GravitinoCatalogStoreFactoryOptions.OAUTH2_CREDENTIAL,
+ "client:secret"));
+
+ Assertions.assertEquals(
+ "http://oauth-server/token",
authProperties.get(OAuth2Properties.OAUTH2_SERVER_URI));
+ }
+
+ @Test
+ void testRestAuthPropertiesBasic() {
+ Map<String, String> authProperties =
+ CONVERTER.toRestAuthProperties(
+ ImmutableMap.of(
+ GravitinoCatalogStoreFactoryOptions.AUTH_TYPE,
+ GravitinoCatalogStoreFactoryOptions.BASIC,
+ GravitinoCatalogStoreFactoryOptions.BASIC_USERNAME,
+ "user",
+ GravitinoCatalogStoreFactoryOptions.BASIC_PASSWORD,
+ "password"));
+
+ Assertions.assertEquals(
+ AuthProperties.AUTH_TYPE_BASIC,
authProperties.get(AuthProperties.AUTH_TYPE));
+ Assertions.assertEquals("user",
authProperties.get(AuthProperties.BASIC_USERNAME));
+ Assertions.assertEquals("password",
authProperties.get(AuthProperties.BASIC_PASSWORD));
+ }
+
+ @Test
+ void testRestAuthPropertiesNoAuth() {
+ Assertions.assertTrue(
+ CONVERTER.toRestAuthProperties(ImmutableMap.of()).isEmpty(),
+ "No auth type configured should yield no REST auth properties");
+ }
}