This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new c0c25ccfd0 [spark] Introduce CTAS support for partitioned tables 
(#9378)
c0c25ccfd0 is described below

commit c0c25ccfd0d823817b4de16222b4b815c973138e
Author: Arnav Balyan <[email protected]>
AuthorDate: Wed Aug 26 08:02:46 2026 +0530

    [spark] Introduce CTAS support for partitioned tables (#9378)
---
 .../shim/PaimonCreateTableAsSelectStrategy.scala   | 18 ++------
 .../shim/PaimonCreateTableAsSelectStrategy.scala   | 18 ++------
 .../shim/PaimonCreateTableAsSelectStrategy.scala   | 18 ++------
 .../java/org/apache/paimon/spark/SparkCatalog.java | 19 ++++++++
 .../org/apache/paimon/spark/util/OptionUtils.scala | 14 +++++-
 .../shim/PaimonCreateTableAsSelectStrategy.scala   | 20 +++-----
 .../apache/paimon/spark/util/OptionUtilsTest.scala | 28 +++++++++++
 .../paimon/spark/sql/FormatTableTestBase.scala     | 54 ++++++++++++++++++----
 8 files changed, 127 insertions(+), 62 deletions(-)

diff --git 
a/paimon-spark/paimon-spark-3.2/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
 
b/paimon-spark/paimon-spark-3.2/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
index 4aa8fb7840..37bf307b02 100644
--- 
a/paimon-spark/paimon-spark-3.2/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
+++ 
b/paimon-spark/paimon-spark-3.2/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
@@ -20,7 +20,6 @@ package org.apache.spark.sql.execution.shim
 
 import org.apache.paimon.Snapshot
 import org.apache.paimon.spark.SparkCatalog
-import org.apache.paimon.spark.catalog.FormatTableCatalog
 import org.apache.paimon.spark.write.PaimonWriteOptions
 
 import org.apache.spark.sql.{SparkSession, Strategy}
@@ -48,18 +47,11 @@ case class PaimonCreateTableAsSelectStrategy(spark: 
SparkSession) extends Strate
         splitTableAndWriteOptions(options)
       val newProps = CatalogV2Util.withDefaultOwnership(props) ++ tableOptions
 
-      val isPartitionedFormatTable = {
-        catalog match {
-          case formatCatalog: FormatTableCatalog =>
-            formatCatalog.isFormatTable(newProps.get("provider").orNull) && 
parts.nonEmpty
-          case _ => false
-        }
-      }
-
-      if (isPartitionedFormatTable) {
-        throw new UnsupportedOperationException(
-          "Using CTAS with partitioned format table is not supported yet.")
-      }
+      catalog.checkPartitionedFormatTableCtas(
+        ident,
+        newProps.get("provider").orNull,
+        parts.nonEmpty,
+        newProps.asJava)
 
       CreateTableAsSelectExec(
         catalog,
diff --git 
a/paimon-spark/paimon-spark-3.3/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
 
b/paimon-spark/paimon-spark-3.3/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
index eb3e044459..1f5bcb0209 100644
--- 
a/paimon-spark/paimon-spark-3.3/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
+++ 
b/paimon-spark/paimon-spark-3.3/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
@@ -20,7 +20,6 @@ package org.apache.spark.sql.execution.shim
 
 import org.apache.paimon.Snapshot
 import org.apache.paimon.spark.SparkCatalog
-import org.apache.paimon.spark.catalog.FormatTableCatalog
 import org.apache.paimon.spark.write.PaimonWriteOptions
 
 import org.apache.spark.sql.{SparkSession, Strategy}
@@ -51,18 +50,11 @@ case class PaimonCreateTableAsSelectStrategy(spark: 
SparkSession)
         splitTableAndWriteOptions(options)
       val qualifiedSpec = qualifyTableSpec(tableSpec, tableOptions)
 
-      val isPartitionedFormatTable = {
-        catalog match {
-          case formatCatalog: FormatTableCatalog =>
-            formatCatalog.isFormatTable(qualifiedSpec.provider.orNull) && 
parts.nonEmpty
-          case _ => false
-        }
-      }
-
-      if (isPartitionedFormatTable) {
-        throw new UnsupportedOperationException(
-          "Using CTAS with partitioned format table is not supported yet.")
-      }
+      catalog.checkPartitionedFormatTableCtas(
+        ident.asIdentifier,
+        qualifiedSpec.provider.orNull,
+        parts.nonEmpty,
+        qualifiedSpec.properties.asJava)
 
       CreateTableAsSelectExec(
         catalog.asTableCatalog,
diff --git 
a/paimon-spark/paimon-spark-3.4/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
 
b/paimon-spark/paimon-spark-3.4/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
index 0e0f3037d2..a046d84314 100644
--- 
a/paimon-spark/paimon-spark-3.4/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
+++ 
b/paimon-spark/paimon-spark-3.4/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
@@ -20,7 +20,6 @@ package org.apache.spark.sql.execution.shim
 
 import org.apache.paimon.Snapshot
 import org.apache.paimon.spark.SparkCatalog
-import org.apache.paimon.spark.catalog.FormatTableCatalog
 import org.apache.paimon.spark.write.PaimonWriteOptions
 
 import org.apache.spark.sql.{SparkSession, Strategy}
@@ -53,18 +52,11 @@ case class PaimonCreateTableAsSelectStrategy(spark: 
SparkSession)
         splitTableAndWriteOptions(options)
       val qualifiedSpec = qualifyTableSpec(tableSpec, tableOptions)
 
-      val isPartitionedFormatTable = {
-        catalog match {
-          case formatCatalog: FormatTableCatalog =>
-            formatCatalog.isFormatTable(qualifiedSpec.provider.orNull) && 
parts.nonEmpty
-          case _ => false
-        }
-      }
-
-      if (isPartitionedFormatTable) {
-        throw new UnsupportedOperationException(
-          "Using CTAS with partitioned format table is not supported yet.")
-      }
+      catalog.checkPartitionedFormatTableCtas(
+        ident,
+        qualifiedSpec.provider.orNull,
+        parts.nonEmpty,
+        qualifiedSpec.properties.asJava)
 
       CreateTableAsSelectExec(
         catalog.asTableCatalog,
diff --git 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkCatalog.java
 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkCatalog.java
index e9f737ace3..6fe7ea5033 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkCatalog.java
+++ 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkCatalog.java
@@ -104,6 +104,7 @@ import static 
org.apache.paimon.spark.SparkTypeUtils.CURRENT_DEFAULT_COLUMN_META
 import static org.apache.paimon.spark.SparkTypeUtils.toPaimonType;
 import static 
org.apache.paimon.spark.util.OptionUtils.checkRequiredConfigurations;
 import static org.apache.paimon.spark.util.OptionUtils.copyWithSQLConf;
+import static 
org.apache.paimon.spark.util.OptionUtils.usePaimonFormatTableImplementation;
 import static org.apache.paimon.spark.util.OptionUtils.withBranchFromOptions;
 import static org.apache.paimon.spark.utils.CatalogUtils.checkNamespace;
 import static org.apache.paimon.spark.utils.CatalogUtils.checkNoDefaultValue;
@@ -398,6 +399,24 @@ public class SparkCatalog extends SparkBaseCatalog
         }
     }
 
+    public void checkPartitionedFormatTableCtas(
+            Identifier ident,
+            @Nullable String provider,
+            boolean partitioned,
+            Map<String, String> properties) {
+        if (partitioned
+                && isFormatTable(provider)
+                && !usePaimonFormatTableImplementation(
+                        catalogName,
+                        toIdentifier(ident, catalogName),
+                        catalog.options(),
+                        properties)) {
+            throw new UnsupportedOperationException(
+                    "Using CTAS with a partitioned engine format table is not 
supported. "
+                            + "Set 'format-table.implementation' to 
'paimon'.");
+        }
+    }
+
     @Override
     public StagedTable stageCreate(
             Identifier ident,
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/OptionUtils.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/OptionUtils.scala
index 1649a57ead..c357256b7c 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/OptionUtils.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/OptionUtils.scala
@@ -19,7 +19,7 @@
 package org.apache.paimon.spark.util
 
 import org.apache.paimon.CoreOptions
-import org.apache.paimon.catalog.Identifier
+import org.apache.paimon.catalog.{CatalogUtils, Identifier}
 import org.apache.paimon.options.ConfigOption
 import org.apache.paimon.spark.{SparkCatalogOptions, SparkConnectorOptions}
 import org.apache.paimon.table.Table
@@ -213,6 +213,18 @@ object OptionUtils extends SQLConfHelper with Logging {
     }
   }
 
+  def usePaimonFormatTableImplementation(
+      catalogName: String,
+      ident: Identifier,
+      catalogOptions: JMap[String, String],
+      tableOptions: JMap[String, String]): Boolean = {
+    val mergedOptions =
+      new JHashMap[String, 
String](CatalogUtils.tableDefaultOptions(catalogOptions))
+    mergedOptions.putAll(tableOptions)
+    mergedOptions.putAll(getMergedOptions(catalogName, ident))
+    new CoreOptions(mergedOptions).formatTableImplementationIsPaimon
+  }
+
   def withBranchFromOptions(
       catalogName: String = null,
       identifier: Identifier = null,
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
index 1a8e3ffe4b..ed6193bee1 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala
@@ -20,7 +20,6 @@ package org.apache.spark.sql.execution.shim
 
 import org.apache.paimon.Snapshot
 import org.apache.paimon.spark.SparkCatalog
-import org.apache.paimon.spark.catalog.FormatTableCatalog
 import org.apache.paimon.spark.write.PaimonWriteOptions
 
 import org.apache.spark.sql.SparkSession
@@ -30,6 +29,8 @@ import 
org.apache.spark.sql.execution.{PaimonTableAsSelectHelper, SparkPlan, Spa
 import org.apache.spark.sql.execution.PaimonTableAsSelectHelper._
 import org.apache.spark.sql.execution.datasources.v2.CreateTableAsSelectExec
 
+import scala.collection.JavaConverters._
+
 case class PaimonCreateTableAsSelectStrategy(spark: SparkSession)
   extends SparkStrategy
   with PaimonTableAsSelectHelper {
@@ -49,18 +50,11 @@ case class PaimonCreateTableAsSelectStrategy(spark: 
SparkSession)
         splitTableAndWriteOptions(options)
       val qualifiedSpec = qualifyTableSpec(tableSpec, tableOptions)
 
-      val isPartitionedFormatTable = {
-        catalog match {
-          case formatCatalog: FormatTableCatalog =>
-            formatCatalog.isFormatTable(qualifiedSpec.provider.orNull) && 
parts.nonEmpty
-          case _ => false
-        }
-      }
-
-      if (isPartitionedFormatTable) {
-        throw new UnsupportedOperationException(
-          "Using CTAS with partitioned format table is not supported yet.")
-      }
+      catalog.checkPartitionedFormatTableCtas(
+        ident,
+        qualifiedSpec.provider.orNull,
+        parts.nonEmpty,
+        qualifiedSpec.properties.asJava)
 
       CreateTableAsSelectExec(
         catalog.asTableCatalog,
diff --git 
a/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/util/OptionUtilsTest.scala
 
b/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/util/OptionUtilsTest.scala
index c766a281b9..018c497272 100644
--- 
a/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/util/OptionUtilsTest.scala
+++ 
b/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/util/OptionUtilsTest.scala
@@ -99,6 +99,34 @@ class OptionUtilsTest extends AnyFunSuite {
     assert(exception.getMessage.contains(METASTORE_PARTITIONED_TABLE.key()))
   }
 
+  test("resolve format table implementation option precedence") {
+    val ident = Identifier.create("test_db", "format_table")
+    val catalogOptions =
+      Map(s"table-default.${FORMAT_TABLE_IMPLEMENTATION.key()}" -> 
"engine").asJava
+
+    assert(
+      !OptionUtils.usePaimonFormatTableImplementation(
+        "test_catalog",
+        ident,
+        catalogOptions,
+        Collections.emptyMap()))
+    assert(
+      OptionUtils.usePaimonFormatTableImplementation(
+        "test_catalog",
+        ident,
+        catalogOptions,
+        Map(FORMAT_TABLE_IMPLEMENTATION.key() -> "paimon").asJava))
+
+    SQLConf.withExistingConf(engineSQLConf) {
+      assert(
+        !OptionUtils.usePaimonFormatTableImplementation(
+          "test_catalog",
+          ident,
+          Collections.emptyMap(),
+          Map(FORMAT_TABLE_IMPLEMENTATION.key() -> "paimon").asJava))
+    }
+  }
+
   private def engineSQLConf: SQLConf = {
     val sqlConf = new SQLConf
     sqlConf.setConfString(
diff --git 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/FormatTableTestBase.scala
 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/FormatTableTestBase.scala
index 998c7c591c..a7e3bdce5e 100644
--- 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/FormatTableTestBase.scala
+++ 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/FormatTableTestBase.scala
@@ -197,18 +197,54 @@ abstract class FormatTableTestBase extends 
PaimonHiveTestBase with AdaptiveSpark
 
   test("Format table: CTAS with partitioned table") {
     withTable("t1", "t2") {
-      sql("CREATE TABLE t1 (id INT, p1 INT, p2 INT) USING csv PARTITIONED BY 
(p1, p2)")
-      sql("INSERT INTO t1 VALUES (1, 2, 3)")
+      sql("CREATE TABLE t1 (id INT, p1 INT, p2 INT) USING csv")
+      sql("INSERT INTO t1 VALUES (1, 2, 3), (2, 2, 4), (3, 5, 6)")
 
-      assertThrows[UnsupportedOperationException] {
-        sql("""
-              |CREATE TABLE t2
-              |USING csv
-              |PARTITIONED BY (p1, p2)
-              |AS SELECT * FROM t1
-              |""".stripMargin)
+      sql("""
+            |CREATE TABLE t2
+            |USING parquet
+            |PARTITIONED BY (p1, p2)
+            |AS SELECT * FROM t1
+            |""".stripMargin)
+
+      checkAnswer(
+        sql("SELECT * FROM t2 ORDER BY id"),
+        Seq(Row(1, 2, 3), Row(2, 2, 4), Row(3, 5, 6)))
+      checkAnswer(
+        sql("SHOW PARTITIONS t2"),
+        Seq(Row("p1=2/p2=3"), Row("p1=2/p2=4"), Row("p1=5/p2=6")))
+
+      val filtered = sql("SELECT * FROM t2 WHERE p1 = 2 AND p2 = 4")
+      checkAnswer(filtered, Seq(Row(2, 2, 4)))
+      assert(collectFilteredInputSplits(filtered.queryExecution.executedPlan, 
"t2").size == 1)
+    }
+  }
+
+  test("Format table: CTAS with partitioned engine table") {
+    def checkRejected(tableProperties: String): Unit = {
+      withTable("t1", "t2") {
+        sql("CREATE TABLE t1 (id INT, p1 INT, p2 INT) USING csv")
+        sql("INSERT INTO t1 VALUES (1, 2, 3)")
+
+        val exception = intercept[UnsupportedOperationException] {
+          sql(s"""
+                 |CREATE TABLE t2
+                 |USING parquet
+                 |PARTITIONED BY (p1, p2)
+                 |$tableProperties
+                 |AS SELECT * FROM t1
+                 |""".stripMargin)
+        }
+        assert(exception.getMessage.contains("partitioned engine format 
table"))
+        assert(!spark.catalog.tableExists("t2"))
       }
     }
+
+    checkRejected("TBLPROPERTIES ('format-table.implementation'='engine')")
+    withSparkSQLConf("spark.paimon.format-table.implementation" -> "engine") {
+      checkRejected("")
+      checkRejected("TBLPROPERTIES ('format-table.implementation'='paimon')")
+    }
   }
 
   test("Format table: create or replace as select supports table type change") 
{

Reply via email to