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 dc52c87f5f [spark] Add option for Hive-style dynamic partition writes 
(#9222)
dc52c87f5f is described below

commit dc52c87f5fe4f6de0cb7df064728dfbeaee1cbe9
Author: Kerwin Zhang <[email protected]>
AuthorDate: Sat Aug 15 20:15:45 2026 +0800

    [spark] Add option for Hive-style dynamic partition writes (#9222)
    
    ### Purpose
    
    Closes #9156.
    
    PR #8414 made positional writes with explicit dynamic partitions
    automatically use Hive's column order. This can silently misalign a
    query that is already in table-schema order when one of its output names
    differs from the target column name.
    
    This change adds
    `spark.paimon.write.hive-style-dynamic-partition.enabled`:
    
    - `false` (default) uses table-schema order, matching the behavior
    before #8414.
    - `true` preserves the Hive-style dynamic partition handling introduced
    by #8414.
    
    ### Tests
    
    CI
---
 docs/generated/spark_connector_configuration.html  |  6 +++
 .../apache/paimon/spark/SparkConnectorOptions.java | 10 +++++
 .../spark/catalyst/analysis/PaimonAnalysis.scala   |  7 +++-
 .../org/apache/paimon/spark/util/OptionUtils.scala |  4 ++
 .../spark/sql/InsertOverwriteTableTestBase.scala   | 43 +++++++++++++++++++++-
 5 files changed, 66 insertions(+), 4 deletions(-)

diff --git a/docs/generated/spark_connector_configuration.html 
b/docs/generated/spark_connector_configuration.html
index cd95fd5fd4..875b3d5639 100644
--- a/docs/generated/spark_connector_configuration.html
+++ b/docs/generated/spark_connector_configuration.html
@@ -104,6 +104,12 @@ under the License.
             <td>Long</td>
             <td>Wait time in milliseconds between retry attempts for Spark V1 
UPDATE on data-evolution tables after row-id range update conflicts.</td>
         </tr>
+        <tr>
+            <td><h5>write.hive-style-dynamic-partition.enabled</h5></td>
+            <td style="word-wrap: break-word;">false</td>
+            <td>Boolean</td>
+            <td>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.</td>
+        </tr>
         <tr>
             <td><h5>write.merge-schema</h5></td>
             <td style="word-wrap: break-word;">false</td>
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 2f315b8df0..108fe12dac 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<Boolean> 
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<Integer> 
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 d888401c25..1ecb417abe 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 37521ceda6..10d403248b 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 ad68360114..40f5d1e120 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 (

Reply via email to