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");
+  }
 }

Reply via email to