diff --git a/docs/generated/spark_connector_configuration.html b/docs/generated/spark_connector_configuration.html index cd95fd5fd41a..875b3d563994 100644 --- a/docs/generated/spark_connector_configuration.html +++ b/docs/generated/spark_connector_configuration.html @@ -104,6 +104,12 @@ Long Wait time in milliseconds between retry attempts for Spark V1 UPDATE on data-evolution tables after row-id range update conflicts. + +
write.hive-style-dynamic-partition.enabled
+ false + Boolean + If true, positional SQL inserts with explicit dynamic partitions use Hive's column order, with non-dynamic columns followed by dynamic partition columns. If false, the query output follows the table schema order. +
write.merge-schema
false diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkConnectorOptions.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkConnectorOptions.java index 2f315b8df0f5..108fe12dac79 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkConnectorOptions.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkConnectorOptions.java @@ -82,6 +82,16 @@ public class SparkConnectorOptions { .withDescription( "If true, v2 write will be used. Currently, only HASH_FIXED and BUCKET_UNAWARE bucket modes are supported. Will fall back to v1 write for other bucket modes. Currently, Spark V2 write does not support TableCapability.STREAMING_WRITE."); + public static final ConfigOption HIVE_STYLE_DYNAMIC_PARTITION_ENABLED = + key("write.hive-style-dynamic-partition.enabled") + .booleanType() + .defaultValue(false) + .withDescription( + "If true, positional SQL inserts with explicit dynamic partitions " + + "use Hive's column order, with non-dynamic columns followed by " + + "dynamic partition columns. If false, the query output follows " + + "the table schema order."); + public static final ConfigOption DATA_EVOLUTION_UPDATE_CONFLICT_RETRY_MAX_ATTEMPTS = key("write.data-evolution.update-conflict-retry.max-attempts") .intType() 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 d888401c25d0..1ecb417abe4d 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 @@ -24,6 +24,7 @@ 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, PaimonHiveDynamicPartitionQuery} import org.apache.paimon.spark.commands.{PaimonAnalyzeTableColumnCommand, PaimonDynamicPartitionOverwriteCommand, PaimonShowColumnsCommand, SchemaEvolutionHelper} +import org.apache.paimon.spark.util.OptionUtils import org.apache.paimon.table.FileStoreTable import org.apache.spark.sql.{PaimonUtils, SparkSession} @@ -109,8 +110,10 @@ class PaimonAnalysis(session: SparkSession) extends Rule[LogicalPlan] { options: Options, mergeSchemaEnabled: Boolean): LogicalPlan = { val query = stripHiveDynamicPartitionMarker(v2WriteCommand.query) + val hiveStyleDynamicPartitionEnabled = OptionUtils.hiveStyleDynamicPartitionEnabled() hiveDynamicPartitionColumns(v2WriteCommand.query) match { - case Some(dynamicPartitionColumns) if !v2WriteCommand.isByName => + case Some(dynamicPartitionColumns) + if hiveStyleDynamicPartitionEnabled && !v2WriteCommand.isByName => resolveDynamicPartitionWrite( query, table, @@ -119,7 +122,7 @@ class PaimonAnalysis(session: SparkSession) extends Rule[LogicalPlan] { mergeSchemaEnabled) case _ => v2WriteCommand match { - case o: OverwritePartitionsDynamic if !o.isByName => + case o: OverwritePartitionsDynamic if hiveStyleDynamicPartitionEnabled && !o.isByName => resolveDynamicPartitionWrite( query, table, diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/OptionUtils.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/OptionUtils.scala index 37521ceda653..10d403248bd3 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/OptionUtils.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/OptionUtils.scala @@ -110,6 +110,10 @@ object OptionUtils extends SQLConfHelper with Logging { getOptionString(SparkConnectorOptions.MERGE_SCHEMA).toBoolean } + def hiveStyleDynamicPartitionEnabled(): Boolean = { + getOptionString(SparkConnectorOptions.HIVE_STYLE_DYNAMIC_PARTITION_ENABLED).toBoolean + } + def writeMergeSchemaExplicitCastEnabled(): Boolean = { getOptionString(SparkConnectorOptions.EXPLICIT_CAST).toBoolean } 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 ad6836011411..40f5d1e120db 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 @@ -727,11 +727,50 @@ abstract class InsertOverwriteTableTestBase extends PaimonSparkTestBase { } } - test("Paimon Insert: V2 dynamic overwrite accepts Hive partition column order") { + test("Paimon Insert: [table-order-default] dynamic partition follows table order") { + for (useV2Write <- Seq("true", "false")) { + withSparkSQLConf( + "spark.sql.sources.partitionOverwriteMode" -> "dynamic", + "spark.paimon.write.use-v2-write" -> useV2Write) { + withTable("target_table") { + sql(""" + |CREATE TABLE target_table ( + | ds STRING, + | part STRING, + | uid STRING, + | value STRING + |) PARTITIONED BY (ds, part) + |TBLPROPERTIES ( + | 'primary-key' = 'ds,part,uid', + | 'bucket' = '2', + | 'bucket-key' = 'uid' + |) + |""".stripMargin) + + sql(""" + |INSERT OVERWRITE target_table PARTITION (ds, part) + |SELECT + | '20260808' AS ds, + | 'p1' AS part, + | '1001' AS uid, + | '0.8' AS metric + |""".stripMargin) + + checkAnswer( + sql("SELECT ds, part, uid, value FROM target_table"), + Row("20260808", "p1", "1001", "0.8")) + } + } + } + } + + test("Paimon Insert: [hive-tail-enabled] dynamic overwrite accepts Hive partition order") { if (gteqSpark3_4) { withSparkSQLConf( "spark.sql.sources.partitionOverwriteMode" -> "dynamic", - "spark.paimon.write.use-v2-write" -> "true") { + "spark.paimon.write.use-v2-write" -> "true", + "spark.paimon.write.hive-style-dynamic-partition.enabled" -> "true" + ) { withTable("my_table") { sql(""" |CREATE TABLE my_table (