This is an automated email from the ASF dual-hosted git repository.
FANNG1 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 374407088a [#4723] feat(spark): support Paimon Spark Procedure (#11845)
374407088a is described below
commit 374407088a99550776aaa3dbc2e3e01beb91b4aa
Author: Yang Zhang <[email protected]>
AuthorDate: Thu Jul 30 17:46:45 2026 +0800
[#4723] feat(spark): support Paimon Spark Procedure (#11845)
### What changes were proposed in this pull request?
Make `GravitinoPaimonCatalog` implement Paimon's `ProcedureCatalog`
interface, enabling Spark `CALL` statements for Paimon procedures.
- Add `ProcedureCatalog` interface to `GravitinoPaimonCatalog`
- Implement `loadProcedure()` by delegating to the underlying Paimon
`SparkBaseCatalog`
### Why are the changes needed?
Paimon provides a set of stored procedures (e.g., `compact`,
`create_tag`, `delete_orphan_files`) that can be invoked via Spark
`CALL` statements. Without implementing `ProcedureCatalog`, these
procedures are not accessible when using Gravitino as the catalog
wrapper.
Fix: #4723
### Does this PR introduce _any_ user-facing change?
Yes. Users can now call Paimon procedures via `CALL
<catalog>.<namespace>.<procedure>(args)` when using the Gravitino Spark
connector with a Paimon catalog.
### How was this patch tested?
Compiled successfully with Spark 3.5 and passed spotlessCheck.
🤖 Generated with Qwen Code (15-shotted by Qwen-Coder)
---
.../connector/paimon/GravitinoPaimonCatalog.java | 25 +++++++++++++++++++++-
.../test/paimon/SparkPaimonCatalogIT.java | 25 ++++++++++++++++++++++
2 files changed, 49 insertions(+), 1 deletion(-)
diff --git
a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/paimon/GravitinoPaimonCatalog.java
b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/paimon/GravitinoPaimonCatalog.java
index 76ed6ed899..1dc2192d28 100644
---
a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/paimon/GravitinoPaimonCatalog.java
+++
b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/paimon/GravitinoPaimonCatalog.java
@@ -27,14 +27,20 @@ import
org.apache.gravitino.spark.connector.PropertiesConverter;
import org.apache.gravitino.spark.connector.SparkTransformConverter;
import org.apache.gravitino.spark.connector.SparkTypeConverter;
import org.apache.gravitino.spark.connector.catalog.BaseCatalog;
+import org.apache.paimon.catalog.Catalog;
import org.apache.paimon.spark.SparkCatalog;
+import org.apache.paimon.spark.SparkProcedures;
import org.apache.paimon.spark.SparkTable;
+import org.apache.paimon.spark.analysis.NoSuchProcedureException;
+import org.apache.paimon.spark.catalog.ProcedureCatalog;
+import org.apache.paimon.spark.procedure.Procedure;
+import org.apache.paimon.spark.procedure.ProcedureBuilder;
import org.apache.spark.sql.connector.catalog.Identifier;
import org.apache.spark.sql.connector.catalog.Table;
import org.apache.spark.sql.connector.catalog.TableCatalog;
import org.apache.spark.sql.util.CaseInsensitiveStringMap;
-public class GravitinoPaimonCatalog extends BaseCatalog {
+public class GravitinoPaimonCatalog extends BaseCatalog implements
ProcedureCatalog {
@Override
protected TableCatalog createAndInitSparkCatalog(
@@ -84,4 +90,21 @@ public class GravitinoPaimonCatalog extends BaseCatalog {
.asTableCatalog()
.purgeTable(NameIdentifier.of(getDatabase(ident), ident.name()));
}
+
+ /**
+ * Procedures will validate the equality of the catalog registered to Spark
catalogManager and the
+ * catalog passed to {@code ProcedureBuilder} which invokes {@code
loadProcedure()}. To meet the
+ * requirement, override the method to pass {@code GravitinoPaimonCatalog}
to the {@code
+ * ProcedureBuilder} instead of the internal spark catalog.
+ */
+ @Override
+ public Procedure loadProcedure(Identifier identifier) throws
NoSuchProcedureException {
+ if (Catalog.SYSTEM_DATABASE_NAME.equals(identifier.namespace()[0])) {
+ ProcedureBuilder builder = SparkProcedures.newBuilder(identifier.name());
+ if (builder != null) {
+ return builder.withTableCatalog(this).build();
+ }
+ }
+ throw new NoSuchProcedureException(identifier);
+ }
}
diff --git
a/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/integration/test/paimon/SparkPaimonCatalogIT.java
b/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/integration/test/paimon/SparkPaimonCatalogIT.java
index 2a395d9a09..1996ccff70 100644
---
a/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/integration/test/paimon/SparkPaimonCatalogIT.java
+++
b/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/integration/test/paimon/SparkPaimonCatalogIT.java
@@ -27,6 +27,7 @@ import
org.apache.gravitino.spark.connector.integration.test.util.SparkTableInfo
import
org.apache.gravitino.spark.connector.integration.test.util.SparkTableInfoChecker;
import org.apache.gravitino.spark.connector.paimon.PaimonPropertiesConstants;
import org.apache.hadoop.fs.Path;
+import org.apache.spark.sql.Row;
import org.apache.spark.sql.types.DataTypes;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
@@ -119,6 +120,30 @@ public abstract class SparkPaimonCatalogIT extends
SparkCommonIT {
checkDirExists(partitionPath);
}
+ @Test
+ void testPaimonCompactProcedure() {
+ String tableName = "test_paimon_compact";
+ dropTableIfExists(tableName);
+ sql(
+ String.format(
+ "CREATE TABLE %s (id INT COMMENT 'id comment', name STRING COMMENT
'', address STRING COMMENT '') USING paimon",
+ tableName));
+
+ sql(String.format("INSERT INTO %s VALUES(1, 'a', 'beijing')", tableName));
+ sql(String.format("INSERT INTO %s VALUES(2, 'b', 'shanghai')", tableName));
+
+ String fullTableName =
+ String.format("%s.%s.%s", getCatalogName(), getDefaultDatabase(),
tableName);
+
+ List<Row> result =
+ getSparkSession()
+ .sql(
+ String.format(
+ "CALL %s.sys.compact(table => '%s')", getCatalogName(),
fullTableName))
+ .collectAsList();
+ Assertions.assertFalse(result.isEmpty());
+ }
+
@Test
void testPaimonPartitionManagement() {
testPaimonListAndDropPartition();