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 6736406ea9 [#13208] feat(spark-connector): Allow external jars to 
register Spark catalogs through an SPI (#13457)
6736406ea9 is described below

commit 6736406ea9a3d9c3f989a38a2fb620d60ab9d6a6
Author: Yuhui <[email protected]>
AuthorDate: Wed Sep 23 18:40:04 2026 +0800

    [#13208] feat(spark-connector): Allow external jars to register Spark 
catalogs through an SPI (#13457)
    
    ### What changes were proposed in this pull request?
    
    Add a `SparkCatalogExtension` SPI, discovered through `ServiceLoader`,
    so a jar outside the connector can register a Spark catalog class for a
    provider this build has no in-tree class for.
    `SparkBindings.Builder.build()` merges discovered extensions before
    validating required bindings.
    
    ### Why are the changes needed?
    
    A catalog implementation living outside the connector (e.g. a
    vendor-specific JDBC catalog) cannot currently be plugged in without
    modifying the connector itself.
    
    Fix: #13208
    
    ### Does this PR introduce _any_ user-facing change?
    
    Adds a new extension point. No behavior change for existing builds: a
    jar with a
    
`META-INF/services/org.apache.gravitino.spark.connector.plugin.SparkCatalogExtension`
    entry can now supply a catalog for a kind the build doesn't already bind
    at compile time.
    
    ### How was this patch tested?
    
    Added `TestSparkCatalogExtensionLoader`, covering: fills an omitted
    kind, unknown-provider skip, already-bound-kind skip, throwing-extension
    skip, no-op with no extensions. Existing `TestSparkBindings` still
    passes.
    
    Co-authored-by: Jerry Shao <[email protected]>
---
 .../spark/connector/plugin/SparkBindings.java      |   4 +
 .../connector/plugin/SparkCatalogExtension.java    |  50 ++++++++
 .../plugin/SparkCatalogExtensionLoader.java        |  85 ++++++++++++++
 .../plugin/TestSparkCatalogExtensionLoader.java    | 127 +++++++++++++++++++++
 4 files changed, 266 insertions(+)

diff --git 
a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/plugin/SparkBindings.java
 
b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/plugin/SparkBindings.java
index fb158ba3b2..bebf22294f 100644
--- 
a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/plugin/SparkBindings.java
+++ 
b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/plugin/SparkBindings.java
@@ -155,11 +155,15 @@ public final class SparkBindings {
     /**
      * Builds the bindings, failing if this build left one out.
      *
+     * <p>Before validating, merges in any {@link SparkCatalogExtension} 
discovered on the classpath
+     * for a kind this build did not already bind at compile time.
+     *
      * @return the bindings
      */
     public SparkBindings build() {
       Preconditions.checkState(
           authorizationExtension != null, "No authorization session extension 
was bound");
+      SparkCatalogExtensionLoader.registerDiscoveredCatalogs(this);
       Set<SparkCatalogKind> missing = EnumSet.copyOf(REQUIRED_KINDS);
       missing.removeAll(catalogClassNames.keySet());
       Preconditions.checkState(missing.isEmpty(), "No catalog was bound for 
%s", missing);
diff --git 
a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/plugin/SparkCatalogExtension.java
 
b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/plugin/SparkCatalogExtension.java
new file mode 100644
index 0000000000..8b344bc9ef
--- /dev/null
+++ 
b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/plugin/SparkCatalogExtension.java
@@ -0,0 +1,50 @@
+/*
+ * 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.spark.connector.plugin;
+
+/**
+ * A catalog binding contributed by a jar outside this connector, discovered 
through {@link
+ * java.util.ServiceLoader}.
+ *
+ * <p>A build that carries no in-tree class for a {@link
+ * org.apache.gravitino.spark.connector.catalog.SparkCatalogKind}, such as 
Paimon on a Scala version
+ * Paimon publishes no artifact for, can still support it by shipping a 
separate jar with a {@code
+ * 
META-INF/services/org.apache.gravitino.spark.connector.plugin.SparkCatalogExtension}
 entry. See
+ * {@link SparkBindings.Builder#build()} for how discovered extensions are 
merged with the bindings
+ * a version module supplies at compile time.
+ */
+public interface SparkCatalogExtension {
+
+  /**
+   * Returns the Gravitino catalog provider this extension serves, such as 
{@code lakehouse-paimon}.
+   * Resolved to a {@link 
org.apache.gravitino.spark.connector.catalog.SparkCatalogKind} through
+   * {@link 
org.apache.gravitino.spark.connector.catalog.SparkCatalogKind#fromProvider(String)}.
+   *
+   * @return the Gravitino catalog provider
+   */
+  String provider();
+
+  /**
+   * Returns the Spark catalog class name to register for {@link #provider()}.
+   *
+   * @return the fully qualified Spark catalog class name
+   * @throws IllegalStateException if this jar has no catalog class for the 
running Spark version
+   */
+  String catalogClassName();
+}
diff --git 
a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/plugin/SparkCatalogExtensionLoader.java
 
b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/plugin/SparkCatalogExtensionLoader.java
new file mode 100644
index 0000000000..4846a3fecb
--- /dev/null
+++ 
b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/plugin/SparkCatalogExtensionLoader.java
@@ -0,0 +1,85 @@
+/*
+ * 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.spark.connector.plugin;
+
+import com.google.common.annotations.VisibleForTesting;
+import java.util.Iterator;
+import java.util.ServiceConfigurationError;
+import java.util.ServiceLoader;
+import org.apache.gravitino.spark.connector.catalog.SparkCatalogKind;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * Discovers {@link SparkCatalogExtension}s on the classpath and merges them 
into a {@link
+ * SparkBindings.Builder}: an entry that cannot be loaded, or that this build 
already binds a class
+ * for, is logged and skipped rather than failing the whole build.
+ */
+final class SparkCatalogExtensionLoader {
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(SparkCatalogExtensionLoader.class);
+
+  private SparkCatalogExtensionLoader() {}
+
+  /**
+   * Registers every discoverable {@link SparkCatalogExtension} into {@code 
builder} that this build
+   * did not already bind at compile time.
+   *
+   * @param builder the builder to register discovered extensions into
+   */
+  static void registerDiscoveredCatalogs(SparkBindings.Builder builder) {
+    registerDiscoveredCatalogs(builder, 
ServiceLoader.load(SparkCatalogExtension.class).iterator());
+  }
+
+  @VisibleForTesting
+  static void registerDiscoveredCatalogs(
+      SparkBindings.Builder builder, Iterator<SparkCatalogExtension> 
extensions) {
+    while (true) {
+      try {
+        if (!extensions.hasNext()) {
+          return;
+        }
+        registerOne(builder, extensions.next());
+      } catch (ServiceConfigurationError | LinkageError | RuntimeException e) {
+        // A provider that cannot be instantiated cannot report which provider 
it was for, so the
+        // most useful thing left to log is that one entry was skipped.
+        LOG.warn(
+            "Skip a {} entry that could not be loaded.", 
SparkCatalogExtension.class.getName(), e);
+      }
+    }
+  }
+
+  private static void registerOne(SparkBindings.Builder builder, 
SparkCatalogExtension extension) {
+    String provider = extension.provider();
+    SparkCatalogKind kind = SparkCatalogKind.fromProvider(provider);
+    if (kind == null) {
+      LOG.warn(
+          "Skip {} because provider {} is not supported by this connector.",
+          extension.getClass().getName(),
+          provider);
+      return;
+    }
+    try {
+      builder.catalog(kind, extension.catalogClassName());
+    } catch (RuntimeException e) {
+      LOG.warn(
+          "Skip {} for provider {}: {}", extension.getClass().getName(), 
provider, e.getMessage());
+    }
+  }
+}
diff --git 
a/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/plugin/TestSparkCatalogExtensionLoader.java
 
b/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/plugin/TestSparkCatalogExtensionLoader.java
new file mode 100644
index 0000000000..c328cbc63e
--- /dev/null
+++ 
b/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/plugin/TestSparkCatalogExtensionLoader.java
@@ -0,0 +1,127 @@
+/*
+ * 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.spark.connector.plugin;
+
+import java.util.Collections;
+import java.util.List;
+import org.apache.gravitino.spark.connector.catalog.SparkCatalogKind;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Tests that discovered {@link SparkCatalogExtension}s are merged into a 
{@link
+ * SparkBindings.Builder}, or skipped, correctly. Registration runs through 
{@link
+ * SparkBindings.Builder#build()} in these tests, the same path a real build 
exercises, since {@link
+ * SparkBindings.Builder}'s catalog map is private.
+ */
+public class TestSparkCatalogExtensionLoader {
+
+  @Test
+  void testADiscoveredExtensionFillsAnOmittedKind() {
+    SparkBindings bindings =
+        buildWithExtensions(
+            List.of(fixedExtension("lakehouse-paimon", 
"org.example.PaimonCatalog")));
+
+    Assertions.assertEquals(
+        "org.example.PaimonCatalog",
+        bindings.catalogClassNames().get(SparkCatalogKind.LAKEHOUSE_PAIMON));
+  }
+
+  @Test
+  void testAnExtensionForAnUnknownProviderIsSkipped() {
+    SparkBindings bindings =
+        buildWithExtensions(List.of(fixedExtension("no-such-provider", 
"org.example.Catalog")));
+
+    Assertions.assertFalse(
+        
bindings.catalogClassNames().containsKey(SparkCatalogKind.LAKEHOUSE_PAIMON));
+    Assertions.assertEquals(5, bindings.catalogClassNames().size());
+  }
+
+  @Test
+  void testAnExtensionForAnAlreadyBoundKindIsSkipped() {
+    SparkBindings bindings =
+        buildWithExtensions(List.of(fixedExtension("hive", 
"org.example.OtherHiveCatalog")));
+
+    Assertions.assertEquals(
+        "org.example.HiveCatalog", 
bindings.catalogClassNames().get(SparkCatalogKind.HIVE));
+  }
+
+  @Test
+  void testAnExtensionThatThrowsIsSkippedRatherThanFailingTheWholeLoad() {
+    SparkBindings bindings =
+        buildWithExtensions(
+            List.of(
+                throwingExtension(),
+                fixedExtension("lakehouse-paimon", 
"org.example.PaimonCatalog")));
+
+    Assertions.assertEquals(
+        "org.example.PaimonCatalog",
+        bindings.catalogClassNames().get(SparkCatalogKind.LAKEHOUSE_PAIMON));
+  }
+
+  @Test
+  void testNoExtensionsIsANoOp() {
+    SparkBindings bindings = buildWithExtensions(Collections.emptyList());
+
+    Assertions.assertFalse(
+        
bindings.catalogClassNames().containsKey(SparkCatalogKind.LAKEHOUSE_PAIMON));
+    Assertions.assertEquals(5, bindings.catalogClassNames().size());
+  }
+
+  private static SparkBindings buildWithExtensions(List<SparkCatalogExtension> 
extensions) {
+    SparkBindings.Builder builder =
+        SparkBindings.builder()
+            .authorizationExtension("org.example.AuthorizationExtensions")
+            .catalog(SparkCatalogKind.HIVE, "org.example.HiveCatalog")
+            .catalog(SparkCatalogKind.LAKEHOUSE_ICEBERG, 
"org.example.IcebergCatalog")
+            .catalog(SparkCatalogKind.GLUE, "org.example.GlueCatalog")
+            .catalog(SparkCatalogKind.JDBC, "org.example.JdbcCatalog")
+            .catalog(SparkCatalogKind.JDBC_POSTGRESQL, 
"org.example.PostgreSqlCatalog");
+    SparkCatalogExtensionLoader.registerDiscoveredCatalogs(builder, 
extensions.iterator());
+    return builder.build();
+  }
+
+  private static SparkCatalogExtension fixedExtension(String provider, String 
catalogClassName) {
+    return new SparkCatalogExtension() {
+      @Override
+      public String provider() {
+        return provider;
+      }
+
+      @Override
+      public String catalogClassName() {
+        return catalogClassName;
+      }
+    };
+  }
+
+  private static SparkCatalogExtension throwingExtension() {
+    return new SparkCatalogExtension() {
+      @Override
+      public String provider() {
+        throw new IllegalStateException("boom");
+      }
+
+      @Override
+      public String catalogClassName() {
+        throw new IllegalStateException("boom");
+      }
+    };
+  }
+}

Reply via email to