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