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")
{