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 367eddb43f [spark] Fix self-referencing RTAS data loss (#8232)
367eddb43f is described below
commit 367eddb43fa2ec938ac7cd9176d4ddfea00cab99
Author: Zouxxyy <[email protected]>
AuthorDate: Sun Jun 14 19:06:27 2026 +0800
[spark] Fix self-referencing RTAS data loss (#8232)
A self-referencing RTAS like `CREATE OR REPLACE TABLE t AS SELECT * FROM
t` previously read the table *after* it was truncated, losing all data.
This PR pins the query to the pre-truncation snapshot for
self-referencing RTAS, while leaving relations with user-specified time
travel (e.g. `VERSION AS OF`) untouched. The shared logic is extracted
into `PaimonTableAsSelectHelper` to avoid duplication across Spark
versions.
---
.../shim/PaimonCreateTableAsSelectStrategy.scala | 6 +-
.../shim/PaimonCreateTableAsSelectStrategy.scala | 8 +-
.../shim/PaimonCreateTableAsSelectStrategy.scala | 8 +-
.../shim/PaimonReplaceTableAsSelectStrategy.scala | 114 +++----------
.../spark/sql/execution/PaimonStrategyHelper.scala | 62 -------
...ategy.scala => PaimonTableAsSelectHelper.scala} | 183 +++++++++------------
.../shim/PaimonCreateTableAsSelectStrategy.scala | 8 +-
.../shim/PaimonReplaceTableAsSelectStrategy.scala | 119 ++------------
.../org/apache/paimon/spark/sql/DDLTestBase.scala | 35 ++++
9 files changed, 172 insertions(+), 371 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 96f194668f..44095b7d2e 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
@@ -24,7 +24,8 @@ import org.apache.paimon.spark.catalog.FormatTableCatalog
import org.apache.spark.sql.{SparkSession, Strategy}
import org.apache.spark.sql.catalyst.plans.logical.{CreateTableAsSelect,
LogicalPlan}
import org.apache.spark.sql.connector.catalog.CatalogV2Util
-import org.apache.spark.sql.execution.{PaimonStrategyHelper, SparkPlan}
+import org.apache.spark.sql.execution.{PaimonTableAsSelectHelper, SparkPlan}
+import org.apache.spark.sql.execution.PaimonTableAsSelectHelper._
import org.apache.spark.sql.execution.datasources.v2.CreateTableAsSelectExec
import org.apache.spark.sql.util.CaseInsensitiveStringMap
@@ -41,7 +42,8 @@ case class PaimonCreateTableAsSelectStrategy(spark:
SparkSession) extends Strate
props,
options,
ifNotExists) =>
- val (tableOptions, writeOptions) =
PaimonStrategyHelper.splitTableAndWriteOptions(options)
+ val (tableOptions, writeOptions) =
+ splitTableAndWriteOptions(options)
val newProps = CatalogV2Util.withDefaultOwnership(props) ++ tableOptions
val isPartitionedFormatTable = {
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 42d6fe2816..ca10cb259f 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
@@ -24,7 +24,8 @@ import org.apache.paimon.spark.catalog.FormatTableCatalog
import org.apache.spark.sql.{SparkSession, Strategy}
import org.apache.spark.sql.catalyst.analysis.ResolvedDBObjectName
import org.apache.spark.sql.catalyst.plans.logical.{CreateTableAsSelect,
LogicalPlan, TableSpec}
-import org.apache.spark.sql.execution.{PaimonStrategyHelper, SparkPlan}
+import org.apache.spark.sql.execution.{PaimonTableAsSelectHelper, SparkPlan}
+import org.apache.spark.sql.execution.PaimonTableAsSelectHelper._
import org.apache.spark.sql.execution.datasources.v2.CreateTableAsSelectExec
import org.apache.spark.sql.util.CaseInsensitiveStringMap
@@ -32,7 +33,7 @@ import scala.collection.JavaConverters._
case class PaimonCreateTableAsSelectStrategy(spark: SparkSession)
extends Strategy
- with PaimonStrategyHelper {
+ with PaimonTableAsSelectHelper {
import org.apache.spark.sql.connector.catalog.CatalogV2Implicits._
@@ -44,7 +45,8 @@ case class PaimonCreateTableAsSelectStrategy(spark:
SparkSession)
tableSpec: TableSpec,
options,
ifNotExists) =>
- val (tableOptions, writeOptions) =
PaimonStrategyHelper.splitTableAndWriteOptions(options)
+ val (tableOptions, writeOptions) =
+ splitTableAndWriteOptions(options)
val qualifiedSpec = qualifyTableSpec(tableSpec, tableOptions)
val isPartitionedFormatTable = {
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 14f540f3e1..be377c17d1 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
@@ -24,7 +24,8 @@ import org.apache.paimon.spark.catalog.FormatTableCatalog
import org.apache.spark.sql.{SparkSession, Strategy}
import org.apache.spark.sql.catalyst.analysis.ResolvedIdentifier
import org.apache.spark.sql.catalyst.plans.logical.{CreateTableAsSelect,
LogicalPlan, TableSpec}
-import org.apache.spark.sql.execution.{PaimonStrategyHelper, SparkPlan}
+import org.apache.spark.sql.execution.{PaimonTableAsSelectHelper, SparkPlan}
+import org.apache.spark.sql.execution.PaimonTableAsSelectHelper._
import org.apache.spark.sql.execution.datasources.v2.CreateTableAsSelectExec
import org.apache.spark.sql.util.CaseInsensitiveStringMap
@@ -32,7 +33,7 @@ import scala.collection.JavaConverters._
case class PaimonCreateTableAsSelectStrategy(spark: SparkSession)
extends Strategy
- with PaimonStrategyHelper {
+ with PaimonTableAsSelectHelper {
import org.apache.spark.sql.connector.catalog.CatalogV2Implicits._
@@ -46,7 +47,8 @@ case class PaimonCreateTableAsSelectStrategy(spark:
SparkSession)
ifNotExists,
analyzedQuery) =>
assert(analyzedQuery.isDefined)
- val (tableOptions, writeOptions) =
PaimonStrategyHelper.splitTableAndWriteOptions(options)
+ val (tableOptions, writeOptions) =
+ splitTableAndWriteOptions(options)
val qualifiedSpec = qualifyTableSpec(tableSpec, tableOptions)
val isPartitionedFormatTable = {
diff --git
a/paimon-spark/paimon-spark-3.4/src/main/scala/org/apache/spark/sql/execution/shim/PaimonReplaceTableAsSelectStrategy.scala
b/paimon-spark/paimon-spark-3.4/src/main/scala/org/apache/spark/sql/execution/shim/PaimonReplaceTableAsSelectStrategy.scala
index 12ec5ee1d0..ec4c0498e1 100644
---
a/paimon-spark/paimon-spark-3.4/src/main/scala/org/apache/spark/sql/execution/shim/PaimonReplaceTableAsSelectStrategy.scala
+++
b/paimon-spark/paimon-spark-3.4/src/main/scala/org/apache/spark/sql/execution/shim/PaimonReplaceTableAsSelectStrategy.scala
@@ -18,20 +18,15 @@
package org.apache.spark.sql.execution.shim
-import org.apache.paimon.CoreOptions.TYPE
-import org.apache.paimon.options.Options
-import org.apache.paimon.spark.{SparkCatalog, SparkGenericCatalog,
SparkSource, SparkTable}
import org.apache.paimon.spark.catalog.SparkBaseCatalog
import org.apache.spark.sql.{SparkSession, Strategy}
-import org.apache.spark.sql.catalyst.analysis.{NoSuchTableException,
ResolvedIdentifier}
-import org.apache.spark.sql.catalyst.expressions.Literal
-import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan,
OverwriteByExpression, OverwritePartitionsDynamic, ReplaceTable,
ReplaceTableAsSelect, TableSpec}
+import org.apache.spark.sql.catalyst.analysis.ResolvedIdentifier
+import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, ReplaceTable,
ReplaceTableAsSelect, TableSpec}
import org.apache.spark.sql.connector.catalog.{Identifier,
StagingTableCatalog, Table, TableCatalog}
-import org.apache.spark.sql.connector.expressions.Transform
-import org.apache.spark.sql.execution.{PaimonStrategyHelper, SparkPlan}
-import
org.apache.spark.sql.execution.datasources.v2.{AtomicReplaceTableAsSelectExec,
DataSourceV2Relation, ReplaceTableAsSelectExec}
-import org.apache.spark.sql.internal.SQLConf.PartitionOverwriteMode
+import org.apache.spark.sql.execution.{PaimonTableAsSelectHelper, SparkPlan}
+import org.apache.spark.sql.execution.PaimonTableAsSelectHelper._
+import
org.apache.spark.sql.execution.datasources.v2.{AtomicReplaceTableAsSelectExec,
ReplaceTableAsSelectExec}
import org.apache.spark.sql.paimon.shims.SparkShimLoader
import org.apache.spark.sql.util.CaseInsensitiveStringMap
@@ -39,7 +34,7 @@ import scala.collection.JavaConverters._
case class PaimonReplaceTableAsSelectStrategy(spark: SparkSession)
extends Strategy
- with PaimonStrategyHelper {
+ with PaimonTableAsSelectHelper {
import org.apache.spark.sql.connector.catalog.CatalogV2Implicits._
@@ -51,28 +46,31 @@ case class PaimonReplaceTableAsSelectStrategy(spark:
SparkSession)
tableSpec: TableSpec,
options,
orCreate,
- analyzedQuery) if
PaimonReplaceTableStrategyHelper.supportsCatalog(catalog, tableSpec) =>
+ analyzedQuery) if supportsCatalog(catalog, tableSpec) =>
assert(analyzedQuery.isDefined)
// For V1 saveAsTable + overwrite on an existing table, rewrite to
// OverwriteByExpression to preserve table definition.
- if (PaimonReplaceTableStrategyHelper.isV1SaveAsTableOverwrite) {
- val overwrite = PaimonReplaceTableStrategyHelper
- .rewriteToOverwrite(spark, catalog, ident, analyzedQuery.get,
options)
+ if (isV1SaveAsTableOverwrite) {
+ val overwrite =
+ rewriteToOverwrite(spark, catalog, ident, analyzedQuery.get, options)
if (overwrite.isDefined) {
val qe = spark.sessionState.executePlan(overwrite.get)
return qe.sparkPlan :: Nil
}
}
- val (tableOptions, writeOptions) =
PaimonStrategyHelper.splitTableAndWriteOptions(options)
+ val (tableOptions, writeOptions) =
+ splitTableAndWriteOptions(options)
val qualifiedSpec = qualifyTableSpec(tableSpec, tableOptions)
val writeOpts = new CaseInsensitiveStringMap(writeOptions.asJava)
- if (PaimonReplaceTableStrategyHelper.canAtomicReplace(catalog, ident,
qualifiedSpec, parts)) {
+ val pinnedQuery =
+ pinSnapshotInQuery(catalog, ident, analyzedQuery.get)
+ if (canAtomicReplace(catalog, ident, qualifiedSpec, parts)) {
AtomicReplaceTableAsSelectExec(
catalog.asInstanceOf[StagingTableCatalog],
ident,
parts,
- analyzedQuery.get,
+ pinnedQuery,
planLater(query),
qualifiedSpec,
writeOpts,
@@ -84,7 +82,7 @@ case class PaimonReplaceTableAsSelectStrategy(spark:
SparkSession)
catalog,
ident,
parts,
- analyzedQuery.get,
+ pinnedQuery,
planLater(query),
qualifiedSpec,
writeOpts,
@@ -102,7 +100,7 @@ case class PaimonReplaceTableAsSelectStrategy(spark:
SparkSession)
case class PaimonReplaceTableStrategy(spark: SparkSession)
extends Strategy
- with PaimonStrategyHelper {
+ with PaimonTableAsSelectHelper {
import org.apache.spark.sql.connector.catalog.CatalogV2Implicits._
@@ -112,7 +110,7 @@ case class PaimonReplaceTableStrategy(spark: SparkSession)
schemaOrColumns,
parts,
tableSpec: TableSpec,
- orCreate) if
PaimonReplaceTableStrategyHelper.supportsCatalog(catalog, tableSpec) =>
+ orCreate) if supportsCatalog(catalog, tableSpec) =>
val columns =
SparkShimLoader.shim.toReplaceTableColumns(
replace.tableSchema,
@@ -120,7 +118,7 @@ case class PaimonReplaceTableStrategy(spark: SparkSession)
catalog,
ident)
val qualifiedSpec = qualifyTableSpec(tableSpec, Map.empty)
- if (PaimonReplaceTableStrategyHelper.canAtomicReplace(catalog, ident,
qualifiedSpec, parts)) {
+ if (canAtomicReplace(catalog, ident, qualifiedSpec, parts)) {
SparkShimLoader.shim.createAtomicReplaceTableExec(
catalog.asInstanceOf[StagingTableCatalog],
ident,
@@ -140,75 +138,3 @@ case class PaimonReplaceTableStrategy(spark: SparkSession)
case _ => Nil
}
}
-
-private[shim] object PaimonReplaceTableStrategyHelper {
-
- def supportsCatalog(catalog: SparkBaseCatalog, tableSpec: TableSpec):
Boolean = catalog match {
- case _: SparkCatalog => true
- case _: SparkGenericCatalog =>
- tableSpec.provider.exists(_.equalsIgnoreCase(SparkSource.NAME))
- case _ => false
- }
-
- /** @see PaimonReplaceTableStrategyHelper in paimon-spark-common for full
documentation. */
- def isV1SaveAsTableOverwrite: Boolean = {
- Thread.currentThread().getStackTrace.exists {
- e =>
- val cls = e.getClassName
- cls.contains("DataFrameWriter") && !cls.contains("DataFrameWriterV2")
- }
- }
-
- /** @see PaimonReplaceTableStrategyHelper in paimon-spark-common for full
documentation. */
- def rewriteToOverwrite(
- spark: SparkSession,
- catalog: SparkBaseCatalog,
- ident: Identifier,
- query: LogicalPlan,
- writeOptions: Map[String, String]): Option[LogicalPlan] = {
- try {
- val existing = catalog.loadTable(ident)
- if (!existing.isInstanceOf[SparkTable]) return None
- val relation =
- DataSourceV2Relation.create(existing,
Some(catalog.asInstanceOf[TableCatalog]), Some(ident))
- val dynamicOverwrite = existing.partitioning().nonEmpty &&
- spark.sessionState.conf.partitionOverwriteMode ==
PartitionOverwriteMode.DYNAMIC
- if (dynamicOverwrite) {
- Some(OverwritePartitionsDynamic.byName(relation, query, writeOptions))
- } else {
- Some(OverwriteByExpression.byName(relation, query, Literal(true),
writeOptions))
- }
- } catch {
- case _: NoSuchTableException => None
- }
- }
-
- /**
- * Whether replace can use Spark's staged replace path. Paimon's
replaceTable is not a
- * rollbackable atomic replace; it swaps the current schema and truncates
current data while
- * preserving old snapshots. Return false for cases replaceTable would
reject so Spark falls back
- * to drop+create.
- */
- def canAtomicReplace(
- catalog: SparkBaseCatalog,
- ident: Identifier,
- tableSpec: TableSpec,
- parts: Seq[Transform]): Boolean = {
- try {
- val existing = catalog.loadTable(ident)
- if (!existing.isInstanceOf[SparkTable]) return false
- val existingProvider =
-
Option(existing.properties().get(TableCatalog.PROP_PROVIDER)).getOrElse(SparkSource.NAME)
- val targetProvider = tableSpec.provider.getOrElse(SparkSource.NAME)
- if (!existingProvider.equalsIgnoreCase(targetProvider)) return false
- val existingType = Options.fromMap(existing.properties()).get(TYPE)
- val targetType = Options.fromMap(tableSpec.properties.asJava).get(TYPE)
- if (existingType != targetType) return false
- val existingParts = existing.partitioning().toSeq
- existingParts.size == parts.size &&
- existingParts.zip(parts).forall { case (a, b) => a.toString ==
b.toString }
- } catch {
- case _: NoSuchTableException => true
- }
- }
-}
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/PaimonStrategyHelper.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/PaimonStrategyHelper.scala
deleted file mode 100644
index 2ee33a1829..0000000000
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/PaimonStrategyHelper.scala
+++ /dev/null
@@ -1,62 +0,0 @@
-/*
- * 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.spark.sql.execution
-
-import org.apache.paimon.CoreOptions
-import org.apache.paimon.iceberg.IcebergOptions
-
-import org.apache.spark.sql.SparkSession
-import org.apache.spark.sql.catalyst.catalog.CatalogUtils
-import org.apache.spark.sql.catalyst.plans.logical.TableSpec
-import org.apache.spark.sql.internal.StaticSQLConf.WAREHOUSE_PATH
-import org.apache.spark.sql.paimon.shims.SparkShimLoader
-
-import scala.collection.JavaConverters._
-
-trait PaimonStrategyHelper {
-
- def spark: SparkSession
-
- protected def makeQualifiedDBObjectPath(location: String): String = {
- CatalogUtils.makeQualifiedDBObjectPath(
- spark.sharedState.conf.get(WAREHOUSE_PATH),
- location,
- spark.sharedState.hadoopConf)
- }
-
- protected def qualifyTableSpec(
- tableSpec: TableSpec,
- tableOptions: Map[String, String]): TableSpec = {
- SparkShimLoader.shim.copyTableSpec(
- tableSpec,
- tableOptions,
- tableSpec.location.map(makeQualifiedDBObjectPath))
- }
-}
-
-object PaimonStrategyHelper {
- private val tableOptionKeys: Set[String] =
- (CoreOptions.getOptions.asScala.map(_.key()) ++
IcebergOptions.getOptions.asScala.map(
- _.key())).toSet
-
- def splitTableAndWriteOptions(
- options: Map[String, String]): (Map[String, String], Map[String,
String]) = {
- options.partition { case (key, _) => tableOptionKeys.contains(key) }
- }
-}
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/shim/PaimonReplaceTableAsSelectStrategy.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/PaimonTableAsSelectHelper.scala
similarity index 50%
copy from
paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/shim/PaimonReplaceTableAsSelectStrategy.scala
copy to
paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/PaimonTableAsSelectHelper.scala
index 9e156c3e26..a6774200e9 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/shim/PaimonReplaceTableAsSelectStrategy.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/PaimonTableAsSelectHelper.scala
@@ -16,121 +16,62 @@
* limitations under the License.
*/
-package org.apache.spark.sql.execution.shim
+package org.apache.spark.sql.execution
+import org.apache.paimon.CoreOptions
import org.apache.paimon.CoreOptions.TYPE
+import org.apache.paimon.iceberg.IcebergOptions
import org.apache.paimon.options.Options
import org.apache.paimon.spark.{SparkCatalog, SparkGenericCatalog,
SparkSource, SparkTable}
import org.apache.paimon.spark.catalog.SparkBaseCatalog
+import org.apache.paimon.table.FileStoreTable
+import org.apache.paimon.table.source.snapshot.TimeTravelUtil
import org.apache.spark.sql.SparkSession
-import org.apache.spark.sql.catalyst.analysis.{NoSuchTableException,
ResolvedIdentifier}
+import org.apache.spark.sql.catalyst.analysis.NoSuchTableException
+import org.apache.spark.sql.catalyst.catalog.CatalogUtils
import org.apache.spark.sql.catalyst.expressions.Literal
-import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan,
OverwriteByExpression, OverwritePartitionsDynamic, ReplaceTable,
ReplaceTableAsSelect, TableSpec}
-import org.apache.spark.sql.connector.catalog.{Identifier,
StagingTableCatalog, TableCatalog}
+import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan,
OverwriteByExpression, OverwritePartitionsDynamic, TableSpec}
+import org.apache.spark.sql.connector.catalog.{CatalogPlugin, Identifier,
TableCatalog}
import org.apache.spark.sql.connector.expressions.Transform
-import org.apache.spark.sql.execution.{PaimonStrategyHelper, SparkPlan,
SparkStrategy}
import org.apache.spark.sql.execution.datasources.v2.DataSourceV2Relation
import org.apache.spark.sql.internal.SQLConf.PartitionOverwriteMode
+import org.apache.spark.sql.internal.StaticSQLConf.WAREHOUSE_PATH
import org.apache.spark.sql.paimon.shims.SparkShimLoader
import scala.collection.JavaConverters._
-case class PaimonReplaceTableAsSelectStrategy(spark: SparkSession)
- extends SparkStrategy
- with PaimonStrategyHelper {
-
- override def apply(plan: LogicalPlan): Seq[SparkPlan] = plan match {
- case ReplaceTableAsSelect(
- ResolvedIdentifier(catalog: SparkBaseCatalog, ident),
- parts,
- query,
- tableSpec: TableSpec,
- options,
- orCreate,
- true) if PaimonReplaceTableStrategyHelper.supportsCatalog(catalog,
tableSpec) =>
- // For V1 saveAsTable + overwrite on an existing table, rewrite to
- // OverwriteByExpression to preserve table definition.
- if (PaimonReplaceTableStrategyHelper.isV1SaveAsTableOverwrite) {
- val overwrite = PaimonReplaceTableStrategyHelper
- .rewriteToOverwrite(spark, catalog, ident, query, options)
- if (overwrite.isDefined) {
- val qe = spark.sessionState.executePlan(overwrite.get)
- return qe.sparkPlan :: Nil
- }
- }
+/** Shared trait for CreateTableAsSelect and ReplaceTableAsSelect strategies.
*/
+trait PaimonTableAsSelectHelper {
- val (tableOptions, writeOptions) =
PaimonStrategyHelper.splitTableAndWriteOptions(options)
- val qualifiedSpec = qualifyTableSpec(tableSpec, tableOptions)
- if (PaimonReplaceTableStrategyHelper.canAtomicReplace(catalog, ident,
qualifiedSpec, parts)) {
- SparkShimLoader.shim.createAtomicReplaceTableAsSelectExec(
- catalog.asInstanceOf[StagingTableCatalog],
- ident,
- parts,
- query,
- qualifiedSpec,
- writeOptions,
- orCreate = orCreate) :: Nil
- } else {
- SparkShimLoader.shim.createReplaceTableAsSelectExec(
- catalog,
- ident,
- parts,
- query,
- qualifiedSpec,
- writeOptions,
- orCreate = orCreate) :: Nil
- }
- case _ => Nil
+ def spark: SparkSession
+
+ protected def makeQualifiedDBObjectPath(location: String): String = {
+ CatalogUtils.makeQualifiedDBObjectPath(
+ spark.sharedState.conf.get(WAREHOUSE_PATH),
+ location,
+ spark.sharedState.hadoopConf)
}
-}
-case class PaimonReplaceTableStrategy(spark: SparkSession)
- extends SparkStrategy
- with PaimonStrategyHelper {
-
- override def apply(plan: LogicalPlan): Seq[SparkPlan] = plan match {
- case replace @ ReplaceTable(
- ResolvedIdentifier(catalog: SparkBaseCatalog, ident),
- schemaOrColumns,
- parts,
- tableSpec: TableSpec,
- orCreate) if
PaimonReplaceTableStrategyHelper.supportsCatalog(catalog, tableSpec) =>
- val columns =
- SparkShimLoader.shim.toReplaceTableColumns(
- replace.tableSchema,
- schemaOrColumns,
- catalog,
- ident)
- val qualifiedSpec = qualifyTableSpec(tableSpec, Map.empty)
- if (PaimonReplaceTableStrategyHelper.canAtomicReplace(catalog, ident,
qualifiedSpec, parts)) {
- SparkShimLoader.shim.createAtomicReplaceTableExec(
- catalog.asInstanceOf[StagingTableCatalog],
- ident,
- columns,
- parts,
- qualifiedSpec,
- orCreate = orCreate) :: Nil
- } else {
- SparkShimLoader.shim.createReplaceTableExec(
- catalog,
- ident,
- columns,
- parts,
- qualifiedSpec,
- orCreate = orCreate) :: Nil
- }
- case _ => Nil
+ protected def qualifyTableSpec(
+ tableSpec: TableSpec,
+ tableOptions: Map[String, String]): TableSpec = {
+ SparkShimLoader.shim.copyTableSpec(
+ tableSpec,
+ tableOptions,
+ tableSpec.location.map(makeQualifiedDBObjectPath))
}
}
-private[shim] object PaimonReplaceTableStrategyHelper {
+object PaimonTableAsSelectHelper {
- def supportsCatalog(catalog: SparkBaseCatalog, tableSpec: TableSpec):
Boolean = catalog match {
- case _: SparkCatalog => true
- case _: SparkGenericCatalog =>
- tableSpec.provider.exists(_.equalsIgnoreCase(SparkSource.NAME))
- case _ => false
+ private val tableOptionKeys: Set[String] =
+ (CoreOptions.getOptions.asScala.map(_.key()) ++
IcebergOptions.getOptions.asScala.map(
+ _.key())).toSet
+
+ def splitTableAndWriteOptions(
+ options: Map[String, String]): (Map[String, String], Map[String,
String]) = {
+ options.partition { case (key, _) => tableOptionKeys.contains(key) }
}
/** Whether the current call originates from V1
DataFrameWriter.saveAsTable(). */
@@ -142,13 +83,49 @@ private[shim] object PaimonReplaceTableStrategyHelper {
}
}
+ /**
+ * Pin snapshot for self-referencing RTAS queries. When the query reads from
the same table being
+ * replaced, this rewrites the DataSourceV2Relation to pin it to the
pre-truncation snapshot.
+ * Relations with user-specified time travel options are left unchanged.
+ */
+ def pinSnapshotInQuery(
+ catalog: TableCatalog,
+ ident: Identifier,
+ query: LogicalPlan): LogicalPlan = {
+ val snapshotId: java.lang.Long =
+ try {
+ val existing = catalog.loadTable(ident)
+ if (!existing.isInstanceOf[SparkTable]) return query
+ val paimonTable = existing.asInstanceOf[SparkTable].getTable
+ if (!paimonTable.isInstanceOf[FileStoreTable]) return query
+
paimonTable.asInstanceOf[FileStoreTable].snapshotManager().latestSnapshotId()
+ } catch {
+ case _: Exception => return query
+ }
+ if (snapshotId == null) return query
+
+ query.transformDown {
+ case r: DataSourceV2Relation if r.catalog.contains(catalog) &&
r.identifier.contains(ident) =>
+ r.table match {
+ case sparkTable: SparkTable
+ if !TimeTravelUtil.hasTimeTravelOptions(
+ Options.fromMap(sparkTable.getTable.options())) =>
+ val pinnedTable = sparkTable.getTable.copy(
+ java.util.Collections
+ .singletonMap(CoreOptions.SCAN_SNAPSHOT_ID.key(),
snapshotId.toString))
+ SparkShimLoader.shim.copyDataSourceV2Relation(r,
SparkTable.of(pinnedTable), r.output)
+ case _ => r
+ }
+ }
+ }
+
/**
* Rewrite to OverwriteByExpression or OverwritePartitionsDynamic for an
existing table,
* preserving table definition. Returns None if the table does not exist.
*/
def rewriteToOverwrite(
spark: SparkSession,
- catalog: SparkBaseCatalog,
+ catalog: TableCatalog,
ident: Identifier,
query: LogicalPlan,
writeOptions: Map[String, String]): Option[LogicalPlan] = {
@@ -156,7 +133,7 @@ private[shim] object PaimonReplaceTableStrategyHelper {
val existing = catalog.loadTable(ident)
if (!existing.isInstanceOf[SparkTable]) return None
val relation =
- DataSourceV2Relation.create(existing,
Some(catalog.asInstanceOf[TableCatalog]), Some(ident))
+ DataSourceV2Relation.create(existing, Some(catalog), Some(ident))
val dynamicOverwrite = existing.partitioning().nonEmpty &&
spark.sessionState.conf.partitionOverwriteMode ==
PartitionOverwriteMode.DYNAMIC
if (dynamicOverwrite) {
@@ -165,16 +142,17 @@ private[shim] object PaimonReplaceTableStrategyHelper {
Some(OverwriteByExpression.byName(relation, query, Literal(true),
writeOptions))
}
} catch {
- case _: NoSuchTableException => None
+ case _: Exception => None
}
}
- /**
- * Whether replace can use Spark's staged replace path. Paimon's
replaceTable is not a
- * rollbackable atomic replace; it swaps the current schema and truncates
current data while
- * preserving old snapshots. Return false for cases replaceTable would
reject so Spark falls back
- * to drop+create.
- */
+ def supportsCatalog(catalog: SparkBaseCatalog, tableSpec: TableSpec):
Boolean = catalog match {
+ case _: SparkCatalog => true
+ case _: SparkGenericCatalog =>
+ tableSpec.provider.exists(_.equalsIgnoreCase(SparkSource.NAME))
+ case _ => false
+ }
+
def canAtomicReplace(
catalog: SparkBaseCatalog,
ident: Identifier,
@@ -197,5 +175,4 @@ private[shim] object PaimonReplaceTableStrategyHelper {
case _: NoSuchTableException => true
}
}
-
}
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 cc0c12960c..fbfb6e4eca 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
@@ -24,12 +24,13 @@ import org.apache.paimon.spark.catalog.FormatTableCatalog
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.catalyst.analysis.ResolvedIdentifier
import org.apache.spark.sql.catalyst.plans.logical.{CreateTableAsSelect,
LogicalPlan, TableSpec}
-import org.apache.spark.sql.execution.{PaimonStrategyHelper, SparkPlan,
SparkStrategy}
+import org.apache.spark.sql.execution.{PaimonTableAsSelectHelper, SparkPlan,
SparkStrategy}
+import org.apache.spark.sql.execution.PaimonTableAsSelectHelper._
import org.apache.spark.sql.execution.datasources.v2.CreateTableAsSelectExec
case class PaimonCreateTableAsSelectStrategy(spark: SparkSession)
extends SparkStrategy
- with PaimonStrategyHelper {
+ with PaimonTableAsSelectHelper {
import org.apache.spark.sql.connector.catalog.CatalogV2Implicits._
@@ -42,7 +43,8 @@ case class PaimonCreateTableAsSelectStrategy(spark:
SparkSession)
options,
ifNotExists,
true) =>
- val (tableOptions, writeOptions) =
PaimonStrategyHelper.splitTableAndWriteOptions(options)
+ val (tableOptions, writeOptions) =
+ splitTableAndWriteOptions(options)
val qualifiedSpec = qualifyTableSpec(tableSpec, tableOptions)
val isPartitionedFormatTable = {
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/shim/PaimonReplaceTableAsSelectStrategy.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/shim/PaimonReplaceTableAsSelectStrategy.scala
index 9e156c3e26..4c21a448e9 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/shim/PaimonReplaceTableAsSelectStrategy.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/shim/PaimonReplaceTableAsSelectStrategy.scala
@@ -18,27 +18,19 @@
package org.apache.spark.sql.execution.shim
-import org.apache.paimon.CoreOptions.TYPE
-import org.apache.paimon.options.Options
-import org.apache.paimon.spark.{SparkCatalog, SparkGenericCatalog,
SparkSource, SparkTable}
import org.apache.paimon.spark.catalog.SparkBaseCatalog
import org.apache.spark.sql.SparkSession
-import org.apache.spark.sql.catalyst.analysis.{NoSuchTableException,
ResolvedIdentifier}
-import org.apache.spark.sql.catalyst.expressions.Literal
-import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan,
OverwriteByExpression, OverwritePartitionsDynamic, ReplaceTable,
ReplaceTableAsSelect, TableSpec}
-import org.apache.spark.sql.connector.catalog.{Identifier,
StagingTableCatalog, TableCatalog}
-import org.apache.spark.sql.connector.expressions.Transform
-import org.apache.spark.sql.execution.{PaimonStrategyHelper, SparkPlan,
SparkStrategy}
-import org.apache.spark.sql.execution.datasources.v2.DataSourceV2Relation
-import org.apache.spark.sql.internal.SQLConf.PartitionOverwriteMode
+import org.apache.spark.sql.catalyst.analysis.ResolvedIdentifier
+import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, ReplaceTable,
ReplaceTableAsSelect, TableSpec}
+import org.apache.spark.sql.connector.catalog.StagingTableCatalog
+import org.apache.spark.sql.execution.{PaimonTableAsSelectHelper, SparkPlan,
SparkStrategy}
+import org.apache.spark.sql.execution.PaimonTableAsSelectHelper._
import org.apache.spark.sql.paimon.shims.SparkShimLoader
-import scala.collection.JavaConverters._
-
case class PaimonReplaceTableAsSelectStrategy(spark: SparkSession)
extends SparkStrategy
- with PaimonStrategyHelper {
+ with PaimonTableAsSelectHelper {
override def apply(plan: LogicalPlan): Seq[SparkPlan] = plan match {
case ReplaceTableAsSelect(
@@ -48,26 +40,27 @@ case class PaimonReplaceTableAsSelectStrategy(spark:
SparkSession)
tableSpec: TableSpec,
options,
orCreate,
- true) if PaimonReplaceTableStrategyHelper.supportsCatalog(catalog,
tableSpec) =>
+ true) if supportsCatalog(catalog, tableSpec) =>
// For V1 saveAsTable + overwrite on an existing table, rewrite to
// OverwriteByExpression to preserve table definition.
- if (PaimonReplaceTableStrategyHelper.isV1SaveAsTableOverwrite) {
- val overwrite = PaimonReplaceTableStrategyHelper
- .rewriteToOverwrite(spark, catalog, ident, query, options)
+ if (isV1SaveAsTableOverwrite) {
+ val overwrite = rewriteToOverwrite(spark, catalog, ident, query,
options)
if (overwrite.isDefined) {
val qe = spark.sessionState.executePlan(overwrite.get)
return qe.sparkPlan :: Nil
}
}
- val (tableOptions, writeOptions) =
PaimonStrategyHelper.splitTableAndWriteOptions(options)
+ val (tableOptions, writeOptions) = splitTableAndWriteOptions(options)
val qualifiedSpec = qualifyTableSpec(tableSpec, tableOptions)
- if (PaimonReplaceTableStrategyHelper.canAtomicReplace(catalog, ident,
qualifiedSpec, parts)) {
+ // Pin snapshot in query to prevent self-referencing RTAS from reading
truncated data
+ val pinnedQuery = pinSnapshotInQuery(catalog, ident, query)
+ if (canAtomicReplace(catalog, ident, qualifiedSpec, parts)) {
SparkShimLoader.shim.createAtomicReplaceTableAsSelectExec(
catalog.asInstanceOf[StagingTableCatalog],
ident,
parts,
- query,
+ pinnedQuery,
qualifiedSpec,
writeOptions,
orCreate = orCreate) :: Nil
@@ -76,7 +69,7 @@ case class PaimonReplaceTableAsSelectStrategy(spark:
SparkSession)
catalog,
ident,
parts,
- query,
+ pinnedQuery,
qualifiedSpec,
writeOptions,
orCreate = orCreate) :: Nil
@@ -87,7 +80,7 @@ case class PaimonReplaceTableAsSelectStrategy(spark:
SparkSession)
case class PaimonReplaceTableStrategy(spark: SparkSession)
extends SparkStrategy
- with PaimonStrategyHelper {
+ with PaimonTableAsSelectHelper {
override def apply(plan: LogicalPlan): Seq[SparkPlan] = plan match {
case replace @ ReplaceTable(
@@ -95,7 +88,7 @@ case class PaimonReplaceTableStrategy(spark: SparkSession)
schemaOrColumns,
parts,
tableSpec: TableSpec,
- orCreate) if
PaimonReplaceTableStrategyHelper.supportsCatalog(catalog, tableSpec) =>
+ orCreate) if supportsCatalog(catalog, tableSpec) =>
val columns =
SparkShimLoader.shim.toReplaceTableColumns(
replace.tableSchema,
@@ -103,7 +96,7 @@ case class PaimonReplaceTableStrategy(spark: SparkSession)
catalog,
ident)
val qualifiedSpec = qualifyTableSpec(tableSpec, Map.empty)
- if (PaimonReplaceTableStrategyHelper.canAtomicReplace(catalog, ident,
qualifiedSpec, parts)) {
+ if (canAtomicReplace(catalog, ident, qualifiedSpec, parts)) {
SparkShimLoader.shim.createAtomicReplaceTableExec(
catalog.asInstanceOf[StagingTableCatalog],
ident,
@@ -123,79 +116,3 @@ case class PaimonReplaceTableStrategy(spark: SparkSession)
case _ => Nil
}
}
-
-private[shim] object PaimonReplaceTableStrategyHelper {
-
- def supportsCatalog(catalog: SparkBaseCatalog, tableSpec: TableSpec):
Boolean = catalog match {
- case _: SparkCatalog => true
- case _: SparkGenericCatalog =>
- tableSpec.provider.exists(_.equalsIgnoreCase(SparkSource.NAME))
- case _ => false
- }
-
- /** Whether the current call originates from V1
DataFrameWriter.saveAsTable(). */
- def isV1SaveAsTableOverwrite: Boolean = {
- Thread.currentThread().getStackTrace.exists {
- e =>
- val cls = e.getClassName
- cls.contains("DataFrameWriter") && !cls.contains("DataFrameWriterV2")
- }
- }
-
- /**
- * Rewrite to OverwriteByExpression or OverwritePartitionsDynamic for an
existing table,
- * preserving table definition. Returns None if the table does not exist.
- */
- def rewriteToOverwrite(
- spark: SparkSession,
- catalog: SparkBaseCatalog,
- ident: Identifier,
- query: LogicalPlan,
- writeOptions: Map[String, String]): Option[LogicalPlan] = {
- try {
- val existing = catalog.loadTable(ident)
- if (!existing.isInstanceOf[SparkTable]) return None
- val relation =
- DataSourceV2Relation.create(existing,
Some(catalog.asInstanceOf[TableCatalog]), Some(ident))
- val dynamicOverwrite = existing.partitioning().nonEmpty &&
- spark.sessionState.conf.partitionOverwriteMode ==
PartitionOverwriteMode.DYNAMIC
- if (dynamicOverwrite) {
- Some(OverwritePartitionsDynamic.byName(relation, query, writeOptions))
- } else {
- Some(OverwriteByExpression.byName(relation, query, Literal(true),
writeOptions))
- }
- } catch {
- case _: NoSuchTableException => None
- }
- }
-
- /**
- * Whether replace can use Spark's staged replace path. Paimon's
replaceTable is not a
- * rollbackable atomic replace; it swaps the current schema and truncates
current data while
- * preserving old snapshots. Return false for cases replaceTable would
reject so Spark falls back
- * to drop+create.
- */
- def canAtomicReplace(
- catalog: SparkBaseCatalog,
- ident: Identifier,
- tableSpec: TableSpec,
- parts: Seq[Transform]): Boolean = {
- try {
- val existing = catalog.loadTable(ident)
- if (!existing.isInstanceOf[SparkTable]) return false
- val existingProvider =
-
Option(existing.properties().get(TableCatalog.PROP_PROVIDER)).getOrElse(SparkSource.NAME)
- val targetProvider = tableSpec.provider.getOrElse(SparkSource.NAME)
- if (!existingProvider.equalsIgnoreCase(targetProvider)) return false
- val existingType = Options.fromMap(existing.properties()).get(TYPE)
- val targetType = Options.fromMap(tableSpec.properties.asJava).get(TYPE)
- if (existingType != targetType) return false
- val existingParts = existing.partitioning().toSeq
- existingParts.size == parts.size &&
- existingParts.zip(parts).forall { case (a, b) => a.toString ==
b.toString }
- } catch {
- case _: NoSuchTableException => true
- }
- }
-
-}
diff --git
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DDLTestBase.scala
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DDLTestBase.scala
index 98cac73f4a..53e4594578 100644
---
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DDLTestBase.scala
+++
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DDLTestBase.scala
@@ -436,6 +436,41 @@ abstract class DDLTestBase extends PaimonSparkTestBase {
}
}
+ test("Paimon DDL: REPLACE TABLE AS SELECT from same table preserves data") {
+ assume(gteqSpark3_4)
+ withTable("t") {
+ sql("""
+ |CREATE TABLE t (id INT, data STRING)
+ |USING paimon
+ |TBLPROPERTIES ('bucket' = '-1')
+ |""".stripMargin)
+ sql("INSERT INTO t VALUES (1, 'a'), (2, 'b')")
+
+ // Self-referencing RTAS: should read old data, not the truncated data
+ sql("CREATE OR REPLACE TABLE t TBLPROPERTIES ('bucket' = '-1') AS SELECT
* FROM t")
+ checkAnswer(sql("SELECT * FROM t ORDER BY id"), Row(1, "a") :: Row(2,
"b") :: Nil)
+ }
+ }
+
+ test("Paimon DDL: REPLACE TABLE AS SELECT with time travel reads specified
snapshot") {
+ assume(gteqSpark3_4)
+ withTable("t") {
+ sql("""
+ |CREATE TABLE t (id INT, data STRING)
+ |USING paimon
+ |TBLPROPERTIES ('bucket' = '-1')
+ |""".stripMargin)
+ sql("INSERT INTO t VALUES (1, 'v1')")
+ val snapshotId1 = loadTable("t").snapshotManager().latestSnapshotId()
+ sql("INSERT INTO t VALUES (2, 'v2')")
+
+ // RTAS with VERSION AS OF should read the specified snapshot, not the
latest
+ sql(
+ s"CREATE OR REPLACE TABLE t TBLPROPERTIES ('bucket' = '-1') AS SELECT
* FROM t VERSION AS OF $snapshotId1")
+ checkAnswer(sql("SELECT * FROM t ORDER BY id"), Row(1, "v1") :: Nil)
+ }
+ }
+
fileFormats.foreach {
format =>
test(s"Paimon DDL: create table with char/varchar/string, file.format:
$format") {