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


Reply via email to