This is an automated email from the ASF dual-hosted git repository.
diqiu50 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 2e915108f5 [#13583] fix(spark-connector): Skip backend JDBC driver
preload when routing through Iceberg REST server (#13584)
2e915108f5 is described below
commit 2e915108f5ab0c95468bb7432bec2dcb26b3a1d2
Author: Yuhui <[email protected]>
AuthorDate: Tue Sep 29 16:49:22 2026 +0800
[#13583] fix(spark-connector): Skip backend JDBC driver preload when
routing through Iceberg REST server (#13584)
### What changes were proposed in this pull request?
Resolve Iceberg REST routing before preloading the catalog's
`jdbc-driver` in `GravitinoIcebergCatalog`, and only preload the driver
on the legacy Hive/JDBC backend translation path.
### Why are the changes needed?
A jdbc-backed Iceberg catalog routed through the Iceberg REST server
failed with `ClassNotFoundException` when the backend JDBC driver was
not on the Spark classpath, although the routed client never uses it.
Fix: #13583
### Does this PR introduce _any_ user-facing change?
No.
### How was this patch tested?
Added unit tests in `TestGravitinoIcebergCatalogRestRouting` for both
the routed and legacy paths.
---
.../connector/iceberg/GravitinoIcebergCatalog.java | 53 ++++++++++-------
.../TestGravitinoIcebergCatalogRestRouting.java | 68 ++++++++++++++++++++++
2 files changed, 101 insertions(+), 20 deletions(-)
diff --git
a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/iceberg/GravitinoIcebergCatalog.java
b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/iceberg/GravitinoIcebergCatalog.java
index 80169829f5..2b7ac3687f 100644
---
a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/iceberg/GravitinoIcebergCatalog.java
+++
b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/iceberg/GravitinoIcebergCatalog.java
@@ -74,6 +74,34 @@ public class GravitinoIcebergCatalog extends BaseCatalog
@Override
protected TableCatalog createAndInitSparkCatalog(
String name, CaseInsensitiveStringMap options, Map<String, String>
properties) {
+ String catalogBackendName =
IcebergPropertiesUtils.getCatalogBackendName(properties);
+ Map<String, String> all =
+ buildSparkCatalogProperties(
+ name,
+ options,
+ properties,
+ SparkSession.active().sparkContext().conf(),
+ () -> GravitinoCatalogManager.get().getIcebergRestUri());
+ TableCatalog icebergCatalog = new SparkCatalog();
+ icebergCatalog.initialize(catalogBackendName, new
CaseInsensitiveStringMap(all));
+ return icebergCatalog;
+ }
+
+ Map<String, String> buildSparkCatalogProperties(
+ String name,
+ CaseInsensitiveStringMap options,
+ Map<String, String> properties,
+ SparkConf sparkConf,
+ Supplier<Optional<String>> endpointDiscovery) {
+ Optional<String> icebergRestUri =
+ resolveIcebergRestUri(properties, key -> sparkConf.get(key, null),
endpointDiscovery);
+ if (icebergRestUri.isPresent()) {
+ // The routed client only talks to the Iceberg REST server, so the
backend JDBC driver is
+ // not needed on the Spark classpath.
+ return buildAutoRoutedIcebergRestProperties(
+ name, options, properties, icebergRestUri.get(), sparkConf);
+ }
+
String jdbcDriver = properties.get(IcebergConstants.GRAVITINO_JDBC_DRIVER);
if (StringUtils.isNotBlank(jdbcDriver)) {
// If `spark.sql.hive.metastore.jars` is set, Spark will use an isolated
client class loader
@@ -84,26 +112,11 @@ public class GravitinoIcebergCatalog extends BaseCatalog
throw new RuntimeException(e);
}
}
- String catalogBackendName =
IcebergPropertiesUtils.getCatalogBackendName(properties);
- SparkConf sparkConf = SparkSession.active().sparkContext().conf();
- Optional<String> icebergRestUri =
- resolveIcebergRestUri(
- properties,
- key -> sparkConf.get(key, null),
- () -> GravitinoCatalogManager.get().getIcebergRestUri());
- Map<String, String> all;
- if (icebergRestUri.isPresent()) {
- all =
- buildAutoRoutedIcebergRestProperties(
- name, options, properties, icebergRestUri.get(), sparkConf);
- } else {
- all = getPropertiesConverter().toSparkCatalogProperties(options,
properties);
- CredentialPropertyUtils.applyIcebergCredentials(
- CredentialPropertyUtils.getCredentials(gravitinoCatalogClient), all);
- }
- TableCatalog icebergCatalog = new SparkCatalog();
- icebergCatalog.initialize(catalogBackendName, new
CaseInsensitiveStringMap(all));
- return icebergCatalog;
+ Map<String, String> all =
+ getPropertiesConverter().toSparkCatalogProperties(options, properties);
+ CredentialPropertyUtils.applyIcebergCredentials(
+ CredentialPropertyUtils.getCredentials(gravitinoCatalogClient), all);
+ return all;
}
/**
diff --git
a/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/iceberg/TestGravitinoIcebergCatalogRestRouting.java
b/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/iceberg/TestGravitinoIcebergCatalogRestRouting.java
index 5261ab73b9..415b84029d 100644
---
a/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/iceberg/TestGravitinoIcebergCatalogRestRouting.java
+++
b/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/iceberg/TestGravitinoIcebergCatalogRestRouting.java
@@ -19,20 +19,38 @@
package org.apache.gravitino.spark.connector.iceberg;
+import static org.mockito.Mockito.mock;
+
import com.google.common.collect.ImmutableMap;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.atomic.AtomicBoolean;
import org.apache.gravitino.catalog.lakehouse.iceberg.IcebergConstants;
+import org.apache.gravitino.client.GravitinoClient;
import org.apache.gravitino.credential.CredentialConstants;
import org.apache.gravitino.spark.connector.GravitinoSparkConfig;
+import org.apache.gravitino.spark.connector.catalog.GravitinoCatalogManager;
import org.apache.spark.SparkConf;
+import org.apache.spark.sql.util.CaseInsensitiveStringMap;
+import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
/** Tests Iceberg REST endpoint selection for Spark catalogs. */
public class TestGravitinoIcebergCatalogRestRouting {
+ @BeforeAll
+ static void initCatalogManager() {
+ GravitinoClient gravitinoClient = mock(GravitinoClient.class);
+ GravitinoCatalogManager.create(new SparkConf(false), "user", identity ->
gravitinoClient);
+ }
+
+ @AfterAll
+ static void cleanupCatalogManager() {
+ GravitinoCatalogManager.get().close();
+ }
+
@Test
void testRoutingDisabledSkipsDiscovery() {
Map<String, String> sessionConfig =
@@ -226,6 +244,42 @@ public class TestGravitinoIcebergCatalogRestRouting {
Assertions.assertEquals("admin", result.get("rest.auth.basic.username"));
}
+ @Test
+ void testRoutedCatalogDoesNotRequireBackendJdbcDriver() {
+ SparkConf sparkConf = new SparkConf(false);
+ sparkConf.set(GravitinoSparkConfig.GRAVITINO_ICEBERG_REST_URI,
"http://manual/iceberg");
+
+ Map<String, String> result =
+ new GravitinoIcebergCatalog()
+ .buildSparkCatalogProperties(
+ "iceberg_jdbc",
+ CaseInsensitiveStringMap.empty(),
+ jdbcPropertiesWithMissingDriver(),
+ sparkConf,
+ Optional::empty);
+
+ Assertions.assertEquals("http://manual/iceberg",
result.get(IcebergConstants.URI));
+ }
+
+ @Test
+ void testLegacyCatalogPreloadsBackendJdbcDriver() {
+ SparkConf sparkConf = new SparkConf(false);
+ sparkConf.set(GravitinoSparkConfig.GRAVITINO_ICEBERG_REST_ROUTING_ENABLED,
"false");
+
+ RuntimeException exception =
+ Assertions.assertThrows(
+ RuntimeException.class,
+ () ->
+ new GravitinoIcebergCatalog()
+ .buildSparkCatalogProperties(
+ "iceberg_jdbc",
+ CaseInsensitiveStringMap.empty(),
+ jdbcPropertiesWithMissingDriver(),
+ sparkConf,
+ Optional::empty));
+ Assertions.assertInstanceOf(ClassNotFoundException.class,
exception.getCause());
+ }
+
private static ImmutableMap<String, String> hiveProperties() {
return ImmutableMap.of(IcebergConstants.CATALOG_BACKEND, "hive");
}
@@ -243,4 +297,18 @@ public class TestGravitinoIcebergCatalogRestRouting {
IcebergConstants.CATALOG_BACKEND, "hive",
IcebergConstants.WAREHOUSE, "s3://bucket/path");
}
+
+ private static ImmutableMap<String, String>
jdbcPropertiesWithMissingDriver() {
+ return ImmutableMap.of(
+ IcebergConstants.CATALOG_BACKEND,
+ "jdbc",
+ IcebergConstants.URI,
+ "jdbc:postgresql://localhost:5432/iceberg",
+ IcebergConstants.GRAVITINO_JDBC_DRIVER,
+ "org.example.MissingDriver",
+ IcebergConstants.WAREHOUSE,
+ "s3://bucket/path",
+ CredentialConstants.CREDENTIAL_PROVIDERS,
+ "s3-token");
+ }
}