Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions docs/generated/spark_connector_configuration.html
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,12 @@
<td>Boolean</td>
<td>Whether to adjust the target split size based on pruned (projected) columns. If enabled, split size estimation uses only the columns actually being read.</td>
</tr>
<tr>
<td><h5>sql.dynamic-partition-column-order</h5></td>
<td style="word-wrap: break-word;">AUTO</td>
<td><p>Enum</p></td>
<td>Controls how non-BY-NAME Spark SQL dynamic partition writes interpret partition columns. TABLE uses table schema order, HIVE expects dynamic partition columns at the end when they are declared in a PARTITION clause or when a dynamic overwrite query matches that Hive-style output, and AUTO preserves compatible table and Hive order detection.<br /><br />Possible values:<ul><li>"AUTO"</li><li>"TABLE"</li><li>"HIVE"</li></ul></td>
</tr>
<tr>
<td><h5>vector-search.lateral-join.parallelism</h5></td>
<td style="word-wrap: break-word;">16</td>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,15 @@ 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<DynamicPartitionColumnOrder> DYNAMIC_PARTITION_COLUMN_ORDER =
key("sql.dynamic-partition-column-order")
.enumType(DynamicPartitionColumnOrder.class)
.defaultValue(DynamicPartitionColumnOrder.AUTO)
.withDescription(
"Controls how non-BY-NAME Spark SQL dynamic partition writes interpret partition columns. "
+ "TABLE uses table schema order, HIVE expects dynamic partition columns at the end when they are declared in a PARTITION clause or when a dynamic overwrite query matches that Hive-style output, "
+ "and AUTO preserves compatible table and Hive order detection.");

public static final ConfigOption<Integer> DATA_EVOLUTION_UPDATE_CONFLICT_RETRY_MAX_ATTEMPTS =
key("write.data-evolution.update-conflict-retry.max-attempts")
.intType()
Expand Down Expand Up @@ -152,4 +161,11 @@ public class SparkConnectorOptions {
.withDescription(
"Whether to adjust the target split size based on pruned (projected) columns. "
+ "If enabled, split size estimation uses only the columns actually being read.");

/** Column order policy for non-BY-NAME dynamic partition writes. */
public enum DynamicPartitionColumnOrder {
AUTO,
TABLE,
HIVE
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,13 @@
package org.apache.paimon.spark.catalyst.analysis

import org.apache.paimon.options.Options
import org.apache.paimon.spark.SparkConnectorOptions.DynamicPartitionColumnOrder
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, 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}
Expand Down Expand Up @@ -111,19 +113,40 @@ class PaimonAnalysis(session: SparkSession) extends Rule[LogicalPlan] {
val query = stripHiveDynamicPartitionMarker(v2WriteCommand.query)
hiveDynamicPartitionColumns(v2WriteCommand.query) match {
case Some(dynamicPartitionColumns) if !v2WriteCommand.isByName =>
val configuredColumnOrder = OptionUtils.dynamicPartitionColumnOrder()
val columnOrder = configuredColumnOrder match {
case DynamicPartitionColumnOrder.AUTO
if dynamicPartitionColumnsUseTableOrder(query, table, dynamicPartitionColumns) =>
DynamicPartitionColumnOrder.TABLE
case order => order
}
resolveDynamicPartitionWrite(
query,
table,
hiveStyleDynamicPartitionOutput(table, dynamicPartitionColumns),
columnOrder,
hiveStyleDynamicPartitionOutput(query, table, dynamicPartitionColumns),
options,
mergeSchemaEnabled)
case _ =>
v2WriteCommand match {
case o: OverwritePartitionsDynamic if !o.isByName =>
val hiveStyleCandidate = hiveStyleDynamicPartitionOutput(query, table)
val configuredColumnOrder =
if (hiveStyleCandidate.isDefined) {
OptionUtils.dynamicPartitionColumnOrder()
} else {
DynamicPartitionColumnOrder.AUTO
}
val hiveStyleOutput = configuredColumnOrder match {
case DynamicPartitionColumnOrder.TABLE => None
case DynamicPartitionColumnOrder.HIVE | DynamicPartitionColumnOrder.AUTO =>
hiveStyleCandidate
}
resolveDynamicPartitionWrite(
query,
table,
hiveStyleDynamicPartitionOutput(query, table),
configuredColumnOrder,
hiveStyleOutput,
options,
mergeSchemaEnabled)
case _ =>
Expand Down Expand Up @@ -153,13 +176,17 @@ class PaimonAnalysis(session: SparkSession) extends Rule[LogicalPlan] {
private def resolveDynamicPartitionWrite(
query: LogicalPlan,
table: DataSourceV2Relation,
columnOrder: DynamicPartitionColumnOrder,
hiveStyleOutput: Option[Seq[Attribute]],
options: Options,
mergeSchemaEnabled: Boolean): LogicalPlan = {
hiveStyleOutput match {
case Some(hiveStyleOutput)
if !sameOutputNames(query.output, table.output) &&
!sameOutputNames(hiveStyleOutput, table.output) =>
if (columnOrder == DynamicPartitionColumnOrder.HIVE &&
!sameOutputNames(query.output, table.output)) ||
(columnOrder == DynamicPartitionColumnOrder.AUTO &&
!sameOutputNames(query.output, table.output) &&
!sameOutputNames(hiveStyleOutput, table.output)) =>
val hiveStyleQuery =
resolveWriteOutput(query, table.name, hiveStyleOutput, byName = false, mergeSchemaEnabled)
resolveWriteOutput(
Expand Down Expand Up @@ -214,12 +241,31 @@ class PaimonAnalysis(session: SparkSession) extends Rule[LogicalPlan] {
table: DataSourceV2Relation): Option[Seq[Attribute]] = {
val dynamicPartitionColumns =
table.table.asInstanceOf[SparkTable].getTable.partitionKeys().asScala.toSeq
hiveStyleDynamicPartitionOutput(table, dynamicPartitionColumns).filter {
hiveStyleDynamicPartitionOutput(query, table, dynamicPartitionColumns).filter {
hiveStyleOutput => sameOutputNames(query.output, hiveStyleOutput)
}
}

private def dynamicPartitionColumnsUseTableOrder(
query: LogicalPlan,
table: DataSourceV2Relation,
dynamicPartitionColumns: Seq[String]): Boolean = {
if (query.output.size != table.output.size) {
false
} else {
val dynamicPartitionAttrs = table.output.zipWithIndex.filter {
case (attr, _) =>
dynamicPartitionColumns.exists(partition => conf.resolver(attr.name, partition))
}
dynamicPartitionAttrs.size == dynamicPartitionColumns.size &&
dynamicPartitionAttrs.forall {
case (attr, index) => conf.resolver(query.output(index).name, attr.name)
}
}
}

private def hiveStyleDynamicPartitionOutput(
query: LogicalPlan,
table: DataSourceV2Relation,
dynamicPartitionColumns: Seq[String]): Option[Seq[Attribute]] = {
val partitionKeys = table.table.asInstanceOf[SparkTable].getTable.partitionKeys().asScala.toSeq
Expand All @@ -238,7 +284,23 @@ class PaimonAnalysis(session: SparkSession) extends Rule[LogicalPlan] {
}
val hiveStyleOutput = dataAttrs ++ dynamicPartitionAttrs
if (dynamicPartitionAttrs.size == dynamicPartitionColumns.size) {
Some(hiveStyleOutput)
val staticPartitionAttrsByIndex = table.output.zipWithIndex.collect {
case (attr, index)
if partitionKeys.exists(partition => conf.resolver(attr.name, partition)) &&
!dynamicPartitionColumns.exists(partition => conf.resolver(attr.name, partition)) &&
query.output
.lift(index)
.exists(queryAttr => conf.resolver(queryAttr.name, attr.name)) =>
index -> attr
}.toMap
val staticPartitionAttrs = staticPartitionAttrsByIndex.values.toSeq
val remainingAttrs = hiveStyleOutput.filterNot {
attr =>
staticPartitionAttrs.exists(staticAttr => conf.resolver(attr.name, staticAttr.name))
}.iterator
Some(table.output.indices.map {
index => staticPartitionAttrsByIndex.getOrElse(index, remainingAttrs.next())
})
} else {
None
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,14 +22,15 @@ import org.apache.paimon.CoreOptions
import org.apache.paimon.catalog.Identifier
import org.apache.paimon.options.ConfigOption
import org.apache.paimon.spark.{SparkCatalogOptions, SparkConnectorOptions}
import org.apache.paimon.spark.SparkConnectorOptions.DynamicPartitionColumnOrder
import org.apache.paimon.table.Table

import org.apache.spark.internal.Logging
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.catalyst.SQLConfHelper
import org.apache.spark.sql.internal.StaticSQLConf

import java.util.{HashMap => JHashMap, Map => JMap}
import java.util.{HashMap => JHashMap, Locale, Map => JMap}
import java.util.regex.Pattern

import scala.collection.JavaConverters._
Expand Down Expand Up @@ -106,6 +107,19 @@ object OptionUtils extends SQLConfHelper with Logging {
configuredValue && isVersionSupported
}

def dynamicPartitionColumnOrder(): DynamicPartitionColumnOrder = {
val configuredValue = getOptionString(SparkConnectorOptions.DYNAMIC_PARTITION_COLUMN_ORDER)
try {
DynamicPartitionColumnOrder.valueOf(configuredValue.trim.toUpperCase(Locale.ROOT))
} catch {
case _: IllegalArgumentException =>
throw new IllegalArgumentException(
s"Invalid value '$configuredValue' for " +
s"spark.paimon.${SparkConnectorOptions.DYNAMIC_PARTITION_COLUMN_ORDER.key()}. " +
"Supported values are AUTO, TABLE, and HIVE.")
}
}

def writeMergeSchemaEnabled(): Boolean = {
getOptionString(SparkConnectorOptions.MERGE_SCHEMA).toBoolean
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -811,6 +811,38 @@ abstract class InsertOverwriteTableTestBase extends PaimonSparkTestBase {
}
}

test("Paimon Insert: table dynamic partition order survives UNION output aliases") {
if (gteqSpark3_4) {
withSparkSQLConf(
"spark.sql.sources.partitionOverwriteMode" -> "dynamic",
"spark.paimon.write.use-v2-write" -> "true",
"spark.paimon.sql.dynamic-partition-column-order" -> "table"
) {
withTable("dynamic_union") {
sql("""
|CREATE TABLE dynamic_union (
| ds STRING,
| part STRING,
| uid STRING,
| value STRING
|) PARTITIONED BY (ds, part)
|""".stripMargin)

sql("""
|INSERT OVERWRITE dynamic_union PARTITION (ds, part)
|SELECT '2026-08-10' AS ds, 'p1' AS part, 'u1' AS uid, 'v1' AS detail_ratio
|UNION ALL
|SELECT '2026-08-10' AS ds, 'p2' AS part, 'u2' AS uid, 'v2' AS value
|""".stripMargin)

checkAnswer(
sql("SELECT ds, part, uid, value FROM dynamic_union ORDER BY part"),
Seq(Row("2026-08-10", "p1", "u1", "v1"), Row("2026-08-10", "p2", "u2", "v2")))
}
}
}
}

test("Paimon Insert: dynamic insert into table with partition columns contain primary key") {
withSparkSQLConf("spark.sql.shuffle.partitions" -> "10") {
withTable("pk_pt") {
Expand Down
Loading