This is an automated email from the ASF dual-hosted git repository.
YannByron 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 ce92c8a196 [spark] support DataSourceV2 dynamic overwrite in
data-partition column out-of-order mode (#8414)
ce92c8a196 is described below
commit ce92c8a1969fab3d67c6440de6a96b43b945743e
Author: Kerwin Zhang <[email protected]>
AuthorDate: Fri Jul 3 17:09:09 2026 +0800
[spark] support DataSourceV2 dynamic overwrite in data-partition column
out-of-order mode (#8414)
---
.../spark/catalyst/analysis/PaimonAnalysis.scala | 262 ++++++++++++++++-----
.../logical/PaimonHiveDynamicPartitionQuery.scala | 32 +++
.../AbstractPaimonSparkSqlExtensionsParser.scala | 34 ++-
.../spark/sql/InsertOverwriteTableTestBase.scala | 105 +++++++++
4 files changed, 375 insertions(+), 58 deletions(-)
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonAnalysis.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonAnalysis.scala
index 87bdfeffba..785fa5cad8 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonAnalysis.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/analysis/PaimonAnalysis.scala
@@ -22,13 +22,14 @@ import org.apache.paimon.options.Options
import org.apache.paimon.spark.SparkTable
import org.apache.paimon.spark.catalyst.Compatibility
import org.apache.paimon.spark.catalyst.analysis.PaimonRelation.isPaimonTable
-import org.apache.paimon.spark.catalyst.plans.logical.PaimonDropPartitions
+import org.apache.paimon.spark.catalyst.plans.logical.{PaimonDropPartitions,
PaimonHiveDynamicPartitionQuery}
import org.apache.paimon.spark.commands.{PaimonAnalyzeTableColumnCommand,
PaimonDynamicPartitionOverwriteCommand, PaimonShowColumnsCommand,
SchemaEvolutionHelper}
import org.apache.paimon.table.FileStoreTable
import org.apache.spark.sql.{PaimonUtils, SparkSession}
-import org.apache.spark.sql.catalyst.analysis.{NamedRelation, ResolvedTable}
-import org.apache.spark.sql.catalyst.plans.logical._
+import org.apache.spark.sql.catalyst.analysis.ResolvedTable
+import org.apache.spark.sql.catalyst.expressions.Attribute
+import org.apache.spark.sql.catalyst.plans.logical.{AnalysisHelper, _}
import org.apache.spark.sql.catalyst.rules.Rule
import org.apache.spark.sql.catalyst.trees.TreeNodeTag
import org.apache.spark.sql.catalyst.util.CharVarcharUtils
@@ -40,57 +41,53 @@ import scala.collection.JavaConverters._
class PaimonAnalysis(session: SparkSession) extends Rule[LogicalPlan] {
import DataSourceV2Implicits._
import PaimonAnalysis._
- override def apply(plan: LogicalPlan): LogicalPlan =
plan.resolveOperatorsDown {
-
- case a @ PaimonV2WriteCommand(table)
- if !paimonWriteResolved(a.query, table) &&
- a.query.getTagValue(PAIMON_WRITE_RESOLVED).isEmpty =>
- val options = Options.fromMap(writeOptions(a).asJava)
- val mergeSchemaEnabled =
SchemaEvolutionHelper.mergeSchemaEnabled(options)
- val expected = SchemaEvolutionHelper.expectedAttrsForCatalogWrite(
- table,
- a.query.schema,
- options,
- a.isByName,
- session)
- val newQuery = PaimonOutputResolver.resolveOutputColumns(
- table.name,
- expected,
- a.query,
- a.isByName,
- mergeSchemaEnabled)
- if (newQuery ne a.query) {
- // Tag to short-circuit the next Analyzer pass; otherwise inline-kept
extras would loop.
- newQuery.setTagValue(PAIMON_WRITE_RESOLVED, ())
- Compatibility.withNewQuery(a, newQuery)
- } else {
- a
+
+ override def apply(plan: LogicalPlan): LogicalPlan = {
+ val withPaimonWrites = AnalysisHelper.allowInvokingTransformsInAnalyzer {
+ plan.transformDown {
+ // Keep fallback conversion before the generic V2 rewrite; otherwise
an already-resolved
+ // query can remain as OverwritePartitionsDynamic until Spark's
capability check rejects it.
+ case o @ PaimonDynamicPartitionOverwrite(r, d) if o.resolved =>
+ PaimonDynamicPartitionOverwriteCommand(r, d, o.query,
o.writeOptions, o.isByName)
+
+ case a @ PaimonV2WriteCommand(table)
+ if a.query.getTagValue(PAIMON_WRITE_RESOLVED).isEmpty =>
+ val options = Options.fromMap(writeOptions(a).asJava)
+ val mergeSchemaEnabled =
SchemaEvolutionHelper.mergeSchemaEnabled(options)
+ val newQuery = resolvePaimonWrite(a, table, options,
mergeSchemaEnabled)
+ if (newQuery ne a.query) {
+ // Tag to short-circuit the next Analyzer pass; otherwise
inline-kept extras would loop.
+ newQuery.setTagValue(PAIMON_WRITE_RESOLVED, ())
+ Compatibility.withNewQuery(a, newQuery)
+ } else {
+ a
+ }
}
+ }
- case o @ PaimonDynamicPartitionOverwrite(r, d) if o.resolved =>
- PaimonDynamicPartitionOverwriteCommand(r, d, o.query, o.writeOptions,
o.isByName)
-
- case merge: MergeIntoTable
- if !merge.resolved && isPaimonTable(merge.targetTable) &&
merge.childrenResolved =>
- PaimonMergeIntoResolver(merge, session)
-
- case s @ ShowColumns(PaimonRelation(table), _, _) if s.resolved =>
- PaimonShowColumnsCommand(table)
-
- case d @ PaimonDropPartitions(ResolvedTable(_, _, table: SparkTable, _),
parts, _, _)
- if d.resolved =>
- PaimonDropPartitions.validate(table, parts.asResolvedPartitionSpecs)
- d
-
- case r: ReplaceColumns if r.resolved && isPaimonTable(r.table) =>
- // Spark rewrites REPLACE COLUMNS into a batch that drops every existing
column and re-adds
- // the new set. Re-adding columns assigns brand-new field ids while
existing data files keep
- // the old ids, so same-named columns are read back as null, silently
corrupting data. Reject
- // it here, before the change batch reaches the catalog where it is
indistinguishable from an
- // ordinary drop+add.
- throw new UnsupportedOperationException(
- "ALTER TABLE ... REPLACE COLUMNS is not supported for Paimon tables. "
+
- "Please use RENAME COLUMN, ALTER COLUMN TYPE, DROP COLUMN, and ADD
COLUMN instead.")
+ withPaimonWrites.resolveOperatorsDown {
+ case merge: MergeIntoTable
+ if !merge.resolved && isPaimonTable(merge.targetTable) &&
merge.childrenResolved =>
+ PaimonMergeIntoResolver(merge, session)
+
+ case s @ ShowColumns(PaimonRelation(table), _, _) if s.resolved =>
+ PaimonShowColumnsCommand(table)
+
+ case d @ PaimonDropPartitions(ResolvedTable(_, _, table: SparkTable, _),
parts, _, _)
+ if d.resolved =>
+ PaimonDropPartitions.validate(table, parts.asResolvedPartitionSpecs)
+ d
+
+ case r: ReplaceColumns if r.resolved && isPaimonTable(r.table) =>
+ // Spark rewrites REPLACE COLUMNS into a batch that drops every
existing column and re-adds
+ // the new set. Re-adding columns assigns brand-new field ids while
existing data files keep
+ // the old ids, so same-named columns are read back as null, silently
corrupting data. Reject
+ // it here, before the change batch reaches the catalog where it is
indistinguishable from an
+ // ordinary drop+add.
+ throw new UnsupportedOperationException(
+ "ALTER TABLE ... REPLACE COLUMNS is not supported for Paimon tables.
" +
+ "Please use RENAME COLUMN, ALTER COLUMN TYPE, DROP COLUMN, and ADD
COLUMN instead.")
+ }
}
private def writeOptions(v2WriteCommand: V2WriteCommand): Map[String,
String] = {
@@ -102,12 +99,159 @@ class PaimonAnalysis(session: SparkSession) extends
Rule[LogicalPlan] {
}
}
+ private def resolvePaimonWrite(
+ v2WriteCommand: V2WriteCommand,
+ table: DataSourceV2Relation,
+ options: Options,
+ mergeSchemaEnabled: Boolean): LogicalPlan = {
+ val query = stripHiveDynamicPartitionMarker(v2WriteCommand.query)
+ hiveDynamicPartitionColumns(v2WriteCommand.query) match {
+ case Some(dynamicPartitionColumns) if !v2WriteCommand.isByName =>
+ resolveDynamicPartitionWrite(
+ query,
+ table,
+ hiveStyleDynamicPartitionOutput(table, dynamicPartitionColumns),
+ options,
+ mergeSchemaEnabled)
+ case _ =>
+ v2WriteCommand match {
+ case o: OverwritePartitionsDynamic if !o.isByName =>
+ resolveDynamicPartitionWrite(
+ query,
+ table,
+ hiveStyleDynamicPartitionOutput(query, table),
+ options,
+ mergeSchemaEnabled)
+ case _ =>
+ val expected =
+ expectedAttrsForWrite(query, table, options,
v2WriteCommand.isByName)
+ resolveWriteOutput(
+ query,
+ table.name,
+ expected,
+ v2WriteCommand.isByName,
+ mergeSchemaEnabled)
+ }
+ }
+ }
+
+ private def hiveDynamicPartitionColumns(query: LogicalPlan):
Option[Seq[String]] = {
+ query.collectFirst {
+ case PaimonHiveDynamicPartitionQuery(dynamicPartitionColumns, _) =>
+ dynamicPartitionColumns
+ }
+ }
+
+ private def stripHiveDynamicPartitionMarker(query: LogicalPlan): LogicalPlan
= {
+ query.transformDown { case PaimonHiveDynamicPartitionQuery(_, child) =>
child }
+ }
+
+ private def resolveDynamicPartitionWrite(
+ query: LogicalPlan,
+ table: DataSourceV2Relation,
+ hiveStyleOutput: Option[Seq[Attribute]],
+ options: Options,
+ mergeSchemaEnabled: Boolean): LogicalPlan = {
+ hiveStyleOutput match {
+ case Some(hiveStyleOutput)
+ if !sameOutputNames(query.output, table.output) &&
+ !sameOutputNames(hiveStyleOutput, table.output) =>
+ val hiveStyleQuery =
+ resolveWriteOutput(query, table.name, hiveStyleOutput, byName =
false, mergeSchemaEnabled)
+ resolveWriteOutput(
+ hiveStyleQuery,
+ table.name,
+ expectedAttrsForWrite(hiveStyleQuery, table, options, byName = true),
+ byName = true,
+ mergeSchemaEnabled)
+ case _ =>
+ resolveWriteOutput(
+ query,
+ table.name,
+ expectedAttrsForWrite(query, table, options, byName = false),
+ byName = false,
+ mergeSchemaEnabled)
+ }
+ }
+
+ private def expectedAttrsForWrite(
+ query: LogicalPlan,
+ table: DataSourceV2Relation,
+ options: Options,
+ byName: Boolean): Seq[Attribute] = {
+ SchemaEvolutionHelper.expectedAttrsForCatalogWrite(
+ table,
+ query.schema,
+ options,
+ byName,
+ session)
+ }
+
+ private def resolveWriteOutput(
+ query: LogicalPlan,
+ tableName: String,
+ expectedOutput: Seq[Attribute],
+ byName: Boolean,
+ mergeSchemaEnabled: Boolean): LogicalPlan = {
+ if (paimonWriteResolved(query, expectedOutput)) {
+ query
+ } else {
+ PaimonOutputResolver.resolveOutputColumns(
+ tableName,
+ expectedOutput,
+ query,
+ byName,
+ mergeSchemaEnabled)
+ }
+ }
+
+ private def hiveStyleDynamicPartitionOutput(
+ query: LogicalPlan,
+ table: DataSourceV2Relation): Option[Seq[Attribute]] = {
+ val dynamicPartitionColumns =
+
table.table.asInstanceOf[SparkTable].getTable.partitionKeys().asScala.toSeq
+ hiveStyleDynamicPartitionOutput(table, dynamicPartitionColumns).filter {
+ hiveStyleOutput => sameOutputNames(query.output, hiveStyleOutput)
+ }
+ }
+
+ private def hiveStyleDynamicPartitionOutput(
+ table: DataSourceV2Relation,
+ dynamicPartitionColumns: Seq[String]): Option[Seq[Attribute]] = {
+ val partitionKeys =
table.table.asInstanceOf[SparkTable].getTable.partitionKeys().asScala.toSeq
+ if (partitionKeys.isEmpty || dynamicPartitionColumns.isEmpty) {
+ None
+ } else {
+ val dynamicPartitionAttrs = partitionKeys
+ .filter {
+ partition => dynamicPartitionColumns.exists(dynamic =>
conf.resolver(dynamic, partition))
+ }
+ .flatMap {
+ dynamicPartition => table.output.find(attr =>
conf.resolver(attr.name, dynamicPartition))
+ }
+ val dataAttrs = table.output.filterNot {
+ attr => dynamicPartitionColumns.exists(partition =>
conf.resolver(attr.name, partition))
+ }
+ val hiveStyleOutput = dataAttrs ++ dynamicPartitionAttrs
+ if (dynamicPartitionAttrs.size == dynamicPartitionColumns.size) {
+ Some(hiveStyleOutput)
+ } else {
+ None
+ }
+ }
+ }
+
+ private def sameOutputNames(left: Seq[Attribute], right: Seq[Attribute]):
Boolean = {
+ left.length == right.length &&
+ left.zip(right).forall { case (l, r) => conf.resolver(l.name, r.name) }
+ }
+
// Mirrors Spark's V2WriteCommand `outputResolved` strictness: query and
table outputs must match
// by name, position, type (ignoring nullable compatibility), and
nullability. Any nested
// structural differences also have to be reconciled before we declare the
write resolved.
- private def paimonWriteResolved(query: LogicalPlan, table: NamedRelation):
Boolean = {
- query.output.size == table.output.size &&
- query.output.zip(table.output).forall {
+ private def paimonWriteResolved(query: LogicalPlan, expectedOutput:
Seq[Attribute]): Boolean = {
+ query.output.size == expectedOutput.size &&
+ query.output.zip(expectedOutput).forall {
case (inAttr, outAttr) =>
val inType =
CharVarcharUtils.getRawType(inAttr.metadata).getOrElse(inAttr.dataType)
val outType =
CharVarcharUtils.getRawType(outAttr.metadata).getOrElse(outAttr.dataType)
@@ -126,7 +270,11 @@ object PaimonAnalysis {
case class PaimonPostHocResolutionRules(session: SparkSession) extends
Rule[LogicalPlan] {
override def apply(plan: LogicalPlan): LogicalPlan = {
- plan match {
+ val withoutHiveDynamicPartitionMarkers =
AnalysisHelper.allowInvokingTransformsInAnalyzer {
+ plan.transformDown { case PaimonHiveDynamicPartitionQuery(_, child) =>
child }
+ }
+
+ withoutHiveDynamicPartitionMarkers match {
case a @ AnalyzeTable(
ResolvedTable(catalog, identifier, table: SparkTable, _),
partitionSpec,
@@ -150,7 +298,7 @@ case class PaimonPostHocResolutionRules(session:
SparkSession) extends Rule[Logi
allColumns) if a.resolved =>
PaimonAnalyzeTableColumnCommand(catalog, identifier, table,
columnNames, allColumns)
- case _ => plan
+ case _ => withoutHiveDynamicPartitionMarkers
}
}
}
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/plans/logical/PaimonHiveDynamicPartitionQuery.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/plans/logical/PaimonHiveDynamicPartitionQuery.scala
new file mode 100644
index 0000000000..ca161f808d
--- /dev/null
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/plans/logical/PaimonHiveDynamicPartitionQuery.scala
@@ -0,0 +1,32 @@
+/*
+ * 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.paimon.spark.catalyst.plans.logical
+
+import org.apache.spark.sql.catalyst.expressions.Attribute
+import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, UnaryNode}
+
+case class PaimonHiveDynamicPartitionQuery(dynamicPartitionColumns:
Seq[String], child: LogicalPlan)
+ extends UnaryNode {
+
+ override def output: Seq[Attribute] = child.output
+
+ override protected def withNewChildInternal(newChild: LogicalPlan):
LogicalPlan = {
+ copy(child = newChild)
+ }
+}
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/AbstractPaimonSparkSqlExtensionsParser.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/AbstractPaimonSparkSqlExtensionsParser.scala
index 7108f0715e..7770143c80 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/AbstractPaimonSparkSqlExtensionsParser.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/AbstractPaimonSparkSqlExtensionsParser.scala
@@ -19,6 +19,7 @@
package org.apache.spark.sql.catalyst.parser.extensions
import org.apache.paimon.spark.SparkProcedures
+import
org.apache.paimon.spark.catalyst.plans.logical.PaimonHiveDynamicPartitionQuery
import org.antlr.v4.runtime._
import org.antlr.v4.runtime.atn.PredictionMode
@@ -30,7 +31,7 @@ import org.apache.spark.sql.catalyst.{FunctionIdentifier,
TableIdentifier}
import org.apache.spark.sql.catalyst.expressions.Expression
import org.apache.spark.sql.catalyst.parser.{ParseException, ParserInterface}
import
org.apache.spark.sql.catalyst.parser.extensions.PaimonSqlExtensionsParser.{NonReservedContext,
QuotedIdentifierContext}
-import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan
+import org.apache.spark.sql.catalyst.plans.logical.{AnalysisHelper,
InsertIntoStatement, LogicalPlan}
import org.apache.spark.sql.catalyst.rules.Rule
import org.apache.spark.sql.internal.VariableSubstitution
import org.apache.spark.sql.paimon.shims.SparkShimLoader
@@ -111,6 +112,7 @@ abstract class AbstractPaimonSparkSqlExtensionsParser(val
delegate: ParserInterf
private def parserRules(sparkSession: SparkSession): Seq[Rule[LogicalPlan]]
= {
Seq(
+ MarkHiveDynamicPartitionWrite,
RewritePaimonViewCommands(sparkSession),
RewritePaimonFunctionCommands(sparkSession),
SparkShimLoader.shim.rewritePaimonSQLFunctionCommands(sparkSession),
@@ -368,6 +370,36 @@ class UpperCaseCharStream(wrapped: CodePointCharStream)
extends CharStream {
// scalastyle:on
}
+object MarkHiveDynamicPartitionWrite extends Rule[LogicalPlan] {
+
+ override def apply(plan: LogicalPlan): LogicalPlan = {
+ AnalysisHelper.allowInvokingTransformsInAnalyzer {
+ plan.transformDown {
+ case insert: InsertIntoStatement
+ if insert.userSpecifiedCols.isEmpty && !isByName(insert) &&
+ insert.partitionSpec.exists(_._2.isEmpty) =>
+ val dynamicPartitionColumns =
+ insert.partitionSpec.collect { case (name, None) => name }.toSeq
+ withNewQuery(
+ insert,
+ PaimonHiveDynamicPartitionQuery(dynamicPartitionColumns,
insert.query))
+ }
+ }
+ }
+
+ private def withNewQuery(insert: InsertIntoStatement, query: LogicalPlan):
InsertIntoStatement = {
+ insert.withNewChildren(Seq(query)).asInstanceOf[InsertIntoStatement]
+ }
+
+ private def isByName(insert: InsertIntoStatement): Boolean = {
+ try {
+ insert.getClass.getMethod("byName").invoke(insert).asInstanceOf[Boolean]
+ } catch {
+ case _: NoSuchMethodException => false
+ }
+ }
+}
+
/** The post-processor validates & cleans-up the parse tree during the parse
process. */
case object PaimonSqlExtensionsPostProcessor extends
PaimonSqlExtensionsBaseListener {
diff --git
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/InsertOverwriteTableTestBase.scala
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/InsertOverwriteTableTestBase.scala
index 2780a49e53..ad68360114 100644
---
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/InsertOverwriteTableTestBase.scala
+++
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/InsertOverwriteTableTestBase.scala
@@ -706,6 +706,111 @@ abstract class InsertOverwriteTableTestBase extends
PaimonSparkTestBase {
}
}
+ test("Paimon Insert: overwrite format(parquet) table in static mode") {
+ try {
+ sql("USE spark_catalog.default")
+ withTable("t_parquet") {
+ sql("""
+ |CREATE TABLE t_parquet (id INT, dt STRING)
+ |USING parquet PARTITIONED BY (dt)
+ |""".stripMargin)
+
+ sql("""
+ |INSERT OVERWRITE t_parquet PARTITION (dt)
+ |SELECT 1 AS id, '2026-07-01' AS dt
+ |""".stripMargin)
+
+ checkAnswer(sql("SELECT id, dt FROM t_parquet"), Row(1, "2026-07-01"))
+ }
+ } finally {
+ sql(s"USE paimon.$dbName0")
+ }
+ }
+
+ test("Paimon Insert: V2 dynamic overwrite accepts Hive partition column
order") {
+ if (gteqSpark3_4) {
+ withSparkSQLConf(
+ "spark.sql.sources.partitionOverwriteMode" -> "dynamic",
+ "spark.paimon.write.use-v2-write" -> "true") {
+ withTable("my_table") {
+ sql("""
+ |CREATE TABLE my_table (
+ | id INT,
+ | dt STRING,
+ | name STRING,
+ | hr STRING
+ |) PARTITIONED BY (dt, hr)
+ |TBLPROPERTIES (
+ | 'primary-key' = 'dt,hr,id',
+ | 'bucket' = '2',
+ | 'bucket-key' = 'id'
+ |)
+ |""".stripMargin)
+
+ sql("""
+ |INSERT INTO my_table VALUES
+ | (1, '2026-06-29', 'old-00', '00'),
+ | (2, '2026-06-29', 'old-01', '01')
+ |""".stripMargin)
+
+ sql("""
+ |INSERT OVERWRITE my_table PARTITION (dt, hr)
+ |SELECT
+ | 3 AS id,
+ | 'new-10' AS name,
+ | '2026-06-30' AS dt,
+ | '10' AS hr
+ |""".stripMargin)
+
+ checkAnswer(
+ sql("SELECT id, dt, name, hr FROM my_table ORDER BY id"),
+ Seq(
+ Row(1, "2026-06-29", "old-00", "00"),
+ Row(2, "2026-06-29", "old-01", "01"),
+ Row(3, "2026-06-30", "new-10", "10"))
+ )
+
+ sql("""
+ |INSERT OVERWRITE my_table PARTITION (dt, hr)
+ |SELECT
+ | 4 AS id,
+ | '2026-07-01' AS dt,
+ | 'table-order-11' AS name,
+ | '11' AS hr
+ |""".stripMargin)
+
+ checkAnswer(
+ sql("SELECT id, dt, name, hr FROM my_table ORDER BY id"),
+ Seq(
+ Row(1, "2026-06-29", "old-00", "00"),
+ Row(2, "2026-06-29", "old-01", "01"),
+ Row(3, "2026-06-30", "new-10", "10"),
+ Row(4, "2026-07-01", "table-order-11", "11"))
+ )
+
+ sql("""
+ |INSERT OVERWRITE my_table PARTITION (dt = '2026-07-02', hr)
+ |SELECT
+ | 5 AS id,
+ | 'mixed-12' AS name,
+ | '12' AS hr
+ |""".stripMargin)
+
+ checkAnswer(
+ sql("SELECT id, dt, name, hr FROM my_table ORDER BY id"),
+ Seq(
+ Row(1, "2026-06-29", "old-00", "00"),
+ Row(2, "2026-06-29", "old-01", "01"),
+ Row(3, "2026-06-30", "new-10", "10"),
+ Row(4, "2026-07-01", "table-order-11", "11"),
+ Row(5, "2026-07-02", "mixed-12", "12")
+ )
+ )
+ }
+ }
+ }
+ }
+
test("Paimon Insert: dynamic insert into table with partition columns
contain primary key") {
withSparkSQLConf("spark.sql.shuffle.partitions" -> "10") {
withTable("pk_pt") {