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();

Reply via email to