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 (