This is an automated email from the ASF dual-hosted git repository.

wForget pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git


The following commit(s) were added to refs/heads/main by this push:
     new 1b545dfac8 [GLUTEN-12474][CORE] Preserve V1 write ordering for dynamic 
partition writes (#12514)
1b545dfac8 is described below

commit 1b545dfac89cd5cd21e6d2ff621bdf957d39a74d
Author: Zhen Wang <[email protected]>
AuthorDate: Thu Jul 23 09:24:14 2026 +0800

    [GLUTEN-12474][CORE] Preserve V1 write ordering for dynamic partition 
writes (#12514)
    
    * test GLUTEN-12474
    
    * Preserve V1 write ordering when inserting local sorts
    
    * add unit tests for other spark version
    
    * fix
    
    * fix spotless check
    
    * Enhance local sort handling for GlutenPlan support
---
 .../columnar/EnsureLocalSortRequirements.scala     | 48 ++++++++++++++++++++--
 .../datasources/GlutenV1WriteCommandSuite.scala    | 35 ++++++++++++++++
 .../datasources/GlutenV1WriteCommandSuite.scala    | 35 ++++++++++++++++
 .../datasources/GlutenV1WriteCommandSuite.scala    | 35 ++++++++++++++++
 .../datasources/GlutenV1WriteCommandSuite.scala    | 35 ++++++++++++++++
 .../org/apache/gluten/sql/shims/SparkShims.scala   | 13 +++++-
 .../gluten/sql/shims/spark34/Spark34Shims.scala    | 15 +++++++
 .../gluten/sql/shims/spark35/Spark35Shims.scala    | 15 +++++++
 .../gluten/sql/shims/spark40/Spark40Shims.scala    | 15 +++++++
 .../gluten/sql/shims/spark41/Spark41Shims.scala    | 15 +++++++
 10 files changed, 256 insertions(+), 5 deletions(-)

diff --git 
a/gluten-substrait/src/main/scala/org/apache/gluten/extension/columnar/EnsureLocalSortRequirements.scala
 
b/gluten-substrait/src/main/scala/org/apache/gluten/extension/columnar/EnsureLocalSortRequirements.scala
index e17a8e7460..d22a71e7e9 100644
--- 
a/gluten-substrait/src/main/scala/org/apache/gluten/extension/columnar/EnsureLocalSortRequirements.scala
+++ 
b/gluten-substrait/src/main/scala/org/apache/gluten/extension/columnar/EnsureLocalSortRequirements.scala
@@ -16,11 +16,15 @@
  */
 package org.apache.gluten.extension.columnar
 
+import org.apache.gluten.execution.GlutenPlan
 import org.apache.gluten.extension.columnar.heuristic.HeuristicTransform
+import org.apache.gluten.sql.shims.SparkShimLoader
 
 import org.apache.spark.sql.catalyst.expressions.SortOrder
 import org.apache.spark.sql.catalyst.rules.Rule
-import org.apache.spark.sql.execution.{SortExec, SparkPlan}
+import org.apache.spark.sql.execution.{ColumnarWriteFilesExec, SortExec, 
SparkPlan}
+import org.apache.spark.sql.execution.datasources.WriteFilesExec
+import org.apache.spark.sql.internal.SQLConf
 
 /**
  * This rule is similar with `EnsureRequirements` but only handle local 
`SortExec`.
@@ -33,25 +37,61 @@ import org.apache.spark.sql.execution.{SortExec, SparkPlan}
 object EnsureLocalSortRequirements extends Rule[SparkPlan] {
   private lazy val transform: HeuristicTransform = HeuristicTransform.static()
 
+  private def numStaticPartitionCols(writeFiles: WriteFilesExec): Int = {
+    // HadoopFs writes include static partition columns in partitionColumns, 
while Hive writes may
+    // only include the partition columns that are present in the write query.
+    val resolver = SQLConf.get.resolver
+    val staticPartitionNames = writeFiles.staticPartitions.keys
+    writeFiles.partitionColumns.takeWhile {
+      partitionColumn => staticPartitionNames.exists(resolver(_, 
partitionColumn.name))
+    }.size
+  }
+
+  private def requiredChildOrdering(plan: SparkPlan): Seq[Seq[SortOrder]] = {
+    plan match {
+      // V1Writes assumes that the logical ordering it prepared is preserved 
in the physical plan,
+      // so WriteFilesExec does not expose requiredChildOrdering itself. 
Gluten may invalidate that
+      // ordering when it replaces a SortAggregateExec with a hash aggregate.
+      case writeFiles: WriteFilesExec
+          if ColumnarWriteFilesExec.OnNoopLeafPath.unapply(writeFiles).isEmpty 
=>
+        Seq(
+          SparkShimLoader.getSparkShims.getV1WriteRequiredOrdering(
+            writeFiles.child.output,
+            writeFiles.partitionColumns,
+            writeFiles.bucketSpec,
+            writeFiles.options,
+            numStaticPartitionCols(writeFiles)))
+      case _ => plan.requiredChildOrdering
+    }
+  }
+
   private def addLocalSort(
+      plan: SparkPlan,
       originalChild: SparkPlan,
       requiredOrdering: Seq[SortOrder]): SparkPlan = {
     // FIXME: HeuristicTransform is costly. Re-applying it may cause 
performance issues.
     val newChild = SortExec(requiredOrdering, global = false, child = 
originalChild)
-    transform.apply(newChild)
+    (plan, originalChild) match {
+      case (_, child: GlutenPlan) if child.supportsColumnar =>
+        transform.apply(newChild)
+      case (parent: GlutenPlan, _) if parent.supportsColumnar =>
+        transform.apply(newChild)
+      case _ =>
+        newChild
+    }
   }
 
   override def apply(plan: SparkPlan): SparkPlan = {
     plan.transformUp {
       case p =>
-        val newChildren = p.children.zip(p.requiredChildOrdering).map {
+        val newChildren = p.children.zip(requiredChildOrdering(p)).map {
           case (child, requiredOrdering) =>
             // If child.outputOrdering already satisfies the requiredOrdering,
             // we do not need to sort.
             if (SortOrder.orderingSatisfies(child.outputOrdering, 
requiredOrdering)) {
               child
             } else {
-              addLocalSort(child, requiredOrdering)
+              addLocalSort(p, child, requiredOrdering)
             }
         }
         p.withNewChildren(newChildren)
diff --git 
a/gluten-ut/spark34/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
 
b/gluten-ut/spark34/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
index eb6794bba8..efd4105e68 100644
--- 
a/gluten-ut/spark34/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
+++ 
b/gluten-ut/spark34/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
@@ -98,6 +98,41 @@ class GlutenV1WriteCommandSuite
   with GlutenV1WriteCommandSuiteBase
   with GlutenSQLTestsBaseTrait {
 
+  testGluten("GLUTEN-12474: preserve ordering for dynamic partition writes") {
+    withSQLConf(
+      "spark.sql.maxConcurrentOutputFileWriters" -> "0",
+      "spark.sql.sources.partitionOverwriteMode" -> "DYNAMIC") {
+      withTable("gluten_12474_src", "gluten_12474_tgt") {
+        sql(
+          """
+            |CREATE TABLE gluten_12474_src USING ORC AS
+            |SELECT concat('k', id) AS k, format_string('v%02d', id) AS v,
+            |  if(id % 2 = 1, '2026-06-01', '2026-06-02') AS day
+            |FROM range(0, 10)
+            |""".stripMargin)
+
+        sql(
+          """
+            |CREATE TABLE gluten_12474_tgt (k STRING, m STRING)
+            |USING ORC
+            |PARTITIONED BY (day STRING)
+            |""".stripMargin)
+
+        sql(
+          """
+            |INSERT OVERWRITE TABLE gluten_12474_tgt PARTITION (day)
+            |SELECT k, max(v) AS m, day
+            |FROM gluten_12474_src
+            |GROUP BY day, k
+            |""".stripMargin)
+
+        checkAnswer(
+          sql("SELECT k, m, day FROM gluten_12474_tgt"),
+          sql("SELECT k, v, day FROM gluten_12474_src"))
+      }
+    }
+  }
+
   testGluten(
     "SPARK-41914: v1 write with AQE and in-partition sorted - non-string 
partition column") {
     withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "true") {
diff --git 
a/gluten-ut/spark35/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
 
b/gluten-ut/spark35/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
index 5fc887d8d4..6291ed3ac7 100644
--- 
a/gluten-ut/spark35/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
+++ 
b/gluten-ut/spark35/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
@@ -100,6 +100,41 @@ class GlutenV1WriteCommandSuite
   with GlutenSQLTestsBaseTrait
   with GlutenColumnarWriteTestSupport {
 
+  testGluten("GLUTEN-12474: preserve ordering for dynamic partition writes") {
+    withSQLConf(
+      "spark.sql.maxConcurrentOutputFileWriters" -> "0",
+      "spark.sql.sources.partitionOverwriteMode" -> "DYNAMIC") {
+      withTable("gluten_12474_src", "gluten_12474_tgt") {
+        sql(
+          """
+            |CREATE TABLE gluten_12474_src USING ORC AS
+            |SELECT concat('k', id) AS k, format_string('v%02d', id) AS v,
+            |  if(id % 2 = 1, '2026-06-01', '2026-06-02') AS day
+            |FROM range(0, 10)
+            |""".stripMargin)
+
+        sql(
+          """
+            |CREATE TABLE gluten_12474_tgt (k STRING, m STRING)
+            |USING ORC
+            |PARTITIONED BY (day STRING)
+            |""".stripMargin)
+
+        sql(
+          """
+            |INSERT OVERWRITE TABLE gluten_12474_tgt PARTITION (day)
+            |SELECT k, max(v) AS m, day
+            |FROM gluten_12474_src
+            |GROUP BY day, k
+            |""".stripMargin)
+
+        checkAnswer(
+          sql("SELECT k, m, day FROM gluten_12474_tgt"),
+          sql("SELECT k, v, day FROM gluten_12474_src"))
+      }
+    }
+  }
+
   testGluten(
     "SPARK-41914: v1 write with AQE and in-partition sorted - non-string 
partition column") {
     withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "true") {
diff --git 
a/gluten-ut/spark40/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
 
b/gluten-ut/spark40/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
index a287f5fffb..b8d9a1156e 100644
--- 
a/gluten-ut/spark40/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
+++ 
b/gluten-ut/spark40/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
@@ -100,6 +100,41 @@ class GlutenV1WriteCommandSuite
   with GlutenSQLTestsBaseTrait
   with GlutenColumnarWriteTestSupport {
 
+  testGluten("GLUTEN-12474: preserve ordering for dynamic partition writes") {
+    withSQLConf(
+      "spark.sql.maxConcurrentOutputFileWriters" -> "0",
+      "spark.sql.sources.partitionOverwriteMode" -> "DYNAMIC") {
+      withTable("gluten_12474_src", "gluten_12474_tgt") {
+        sql(
+          """
+            |CREATE TABLE gluten_12474_src USING ORC AS
+            |SELECT concat('k', id) AS k, format_string('v%02d', id) AS v,
+            |  if(id % 2 = 1, '2026-06-01', '2026-06-02') AS day
+            |FROM range(0, 10)
+            |""".stripMargin)
+
+        sql(
+          """
+            |CREATE TABLE gluten_12474_tgt (k STRING, m STRING)
+            |USING ORC
+            |PARTITIONED BY (day STRING)
+            |""".stripMargin)
+
+        sql(
+          """
+            |INSERT OVERWRITE TABLE gluten_12474_tgt PARTITION (day)
+            |SELECT k, max(v) AS m, day
+            |FROM gluten_12474_src
+            |GROUP BY day, k
+            |""".stripMargin)
+
+        checkAnswer(
+          sql("SELECT k, m, day FROM gluten_12474_tgt"),
+          sql("SELECT k, v, day FROM gluten_12474_src"))
+      }
+    }
+  }
+
   // TODO: fix in Spark-4.0
   ignoreGluten(
     "SPARK-41914: v1 write with AQE and in-partition sorted - non-string 
partition column") {
diff --git 
a/gluten-ut/spark41/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
 
b/gluten-ut/spark41/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
index a287f5fffb..b8d9a1156e 100644
--- 
a/gluten-ut/spark41/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
+++ 
b/gluten-ut/spark41/src/test/scala/org/apache/spark/sql/execution/datasources/GlutenV1WriteCommandSuite.scala
@@ -100,6 +100,41 @@ class GlutenV1WriteCommandSuite
   with GlutenSQLTestsBaseTrait
   with GlutenColumnarWriteTestSupport {
 
+  testGluten("GLUTEN-12474: preserve ordering for dynamic partition writes") {
+    withSQLConf(
+      "spark.sql.maxConcurrentOutputFileWriters" -> "0",
+      "spark.sql.sources.partitionOverwriteMode" -> "DYNAMIC") {
+      withTable("gluten_12474_src", "gluten_12474_tgt") {
+        sql(
+          """
+            |CREATE TABLE gluten_12474_src USING ORC AS
+            |SELECT concat('k', id) AS k, format_string('v%02d', id) AS v,
+            |  if(id % 2 = 1, '2026-06-01', '2026-06-02') AS day
+            |FROM range(0, 10)
+            |""".stripMargin)
+
+        sql(
+          """
+            |CREATE TABLE gluten_12474_tgt (k STRING, m STRING)
+            |USING ORC
+            |PARTITIONED BY (day STRING)
+            |""".stripMargin)
+
+        sql(
+          """
+            |INSERT OVERWRITE TABLE gluten_12474_tgt PARTITION (day)
+            |SELECT k, max(v) AS m, day
+            |FROM gluten_12474_src
+            |GROUP BY day, k
+            |""".stripMargin)
+
+        checkAnswer(
+          sql("SELECT k, m, day FROM gluten_12474_tgt"),
+          sql("SELECT k, v, day FROM gluten_12474_src"))
+      }
+    }
+  }
+
   // TODO: fix in Spark-4.0
   ignoreGluten(
     "SPARK-41914: v1 write with AQE and in-partition sorted - non-string 
partition column") {
diff --git 
a/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala 
b/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
index 750bea0c41..068ccdda6f 100644
--- a/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
+++ b/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
@@ -24,7 +24,8 @@ import org.apache.spark.broadcast.Broadcast
 import org.apache.spark.internal.io.FileCommitProtocol
 import org.apache.spark.sql.{AnalysisException, SparkSession}
 import org.apache.spark.sql.catalyst.InternalRow
-import org.apache.spark.sql.catalyst.expressions.{Attribute, BinaryArithmetic, 
Expression, InputFileBlockLength, InputFileBlockStart, InputFileName, 
RaiseError, UnBase64}
+import org.apache.spark.sql.catalyst.catalog.BucketSpec
+import org.apache.spark.sql.catalyst.expressions.{Attribute, BinaryArithmetic, 
Expression, InputFileBlockLength, InputFileBlockStart, InputFileName, 
RaiseError, SortOrder, UnBase64}
 import org.apache.spark.sql.catalyst.plans.JoinType
 import org.apache.spark.sql.catalyst.plans.QueryPlan
 import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan
@@ -129,6 +130,16 @@ trait SparkShims {
 
   def enableNativeWriteFilesByDefault(): Boolean = false
 
+  // Planned V1 writes were introduced in Spark 3.4. Older versions do not 
expose a required
+  // ordering utility and keep the default empty ordering.
+  // TODO: Remove this shim after dropping Spark 3.3 support.
+  def getV1WriteRequiredOrdering(
+      outputColumns: Seq[Attribute],
+      partitionColumns: Seq[Attribute],
+      bucketSpec: Option[BucketSpec],
+      options: Map[String, String],
+      numStaticPartitionCols: Int): Seq[SortOrder] = Seq.empty
+
   def broadcastInternal[T: ClassTag](sc: SparkContext, value: T): Broadcast[T] 
= {
     // Since Spark 3.4, the `sc.broadcast` has been optimized to use 
`sc.broadcastInternal`.
     // More details see SPARK-39983.
diff --git 
a/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
 
b/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
index 97ff19a84a..ae93d067a5 100644
--- 
a/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
+++ 
b/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
@@ -27,6 +27,7 @@ import org.apache.spark.paths.SparkPath
 import org.apache.spark.sql.{AnalysisException, SparkSession}
 import org.apache.spark.sql.catalyst.InternalRow
 import org.apache.spark.sql.catalyst.analysis.DecimalPrecision
+import org.apache.spark.sql.catalyst.catalog.BucketSpec
 import org.apache.spark.sql.catalyst.expressions._
 import org.apache.spark.sql.catalyst.expressions.aggregate._
 import org.apache.spark.sql.catalyst.plans.QueryPlan
@@ -208,6 +209,20 @@ class Spark34Shims extends SparkShims {
 
   override def enableNativeWriteFilesByDefault(): Boolean = true
 
+  override def getV1WriteRequiredOrdering(
+      outputColumns: Seq[Attribute],
+      partitionColumns: Seq[Attribute],
+      bucketSpec: Option[BucketSpec],
+      options: Map[String, String],
+      numStaticPartitionCols: Int): Seq[SortOrder] = {
+    V1WritesUtils.getSortOrder(
+      outputColumns,
+      partitionColumns,
+      bucketSpec,
+      options,
+      numStaticPartitionCols)
+  }
+
   override def broadcastInternal[T: ClassTag](sc: SparkContext, value: T): 
Broadcast[T] = {
     SparkContextUtils.broadcastInternal(sc, value)
   }
diff --git 
a/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
 
b/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
index d62fdfea19..08047c5757 100644
--- 
a/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
+++ 
b/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
@@ -27,6 +27,7 @@ import org.apache.spark.paths.SparkPath
 import org.apache.spark.sql.{AnalysisException, SparkSession}
 import org.apache.spark.sql.catalyst.{ExtendedAnalysisException, InternalRow}
 import org.apache.spark.sql.catalyst.analysis.DecimalPrecision
+import org.apache.spark.sql.catalyst.catalog.BucketSpec
 import org.apache.spark.sql.catalyst.expressions._
 import org.apache.spark.sql.catalyst.expressions.aggregate._
 import org.apache.spark.sql.catalyst.plans.QueryPlan
@@ -249,6 +250,20 @@ class Spark35Shims extends SparkShims {
 
   override def enableNativeWriteFilesByDefault(): Boolean = true
 
+  override def getV1WriteRequiredOrdering(
+      outputColumns: Seq[Attribute],
+      partitionColumns: Seq[Attribute],
+      bucketSpec: Option[BucketSpec],
+      options: Map[String, String],
+      numStaticPartitionCols: Int): Seq[SortOrder] = {
+    V1WritesUtils.getSortOrder(
+      outputColumns,
+      partitionColumns,
+      bucketSpec,
+      options,
+      numStaticPartitionCols)
+  }
+
   override def broadcastInternal[T: ClassTag](sc: SparkContext, value: T): 
Broadcast[T] = {
     SparkContextUtils.broadcastInternal(sc, value)
   }
diff --git 
a/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
 
b/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
index 5847e62c10..1e50984e19 100644
--- 
a/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
+++ 
b/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
@@ -27,6 +27,7 @@ import org.apache.spark.paths.SparkPath
 import org.apache.spark.sql.{AnalysisException, SparkSession}
 import org.apache.spark.sql.catalyst.{ExtendedAnalysisException, InternalRow}
 import org.apache.spark.sql.catalyst.analysis.DecimalPrecisionTypeCoercion
+import org.apache.spark.sql.catalyst.catalog.BucketSpec
 import org.apache.spark.sql.catalyst.expressions._
 import org.apache.spark.sql.catalyst.expressions.aggregate._
 import org.apache.spark.sql.catalyst.plans.{JoinType, LeftSingle}
@@ -254,6 +255,20 @@ class Spark40Shims extends SparkShims {
 
   override def enableNativeWriteFilesByDefault(): Boolean = true
 
+  override def getV1WriteRequiredOrdering(
+      outputColumns: Seq[Attribute],
+      partitionColumns: Seq[Attribute],
+      bucketSpec: Option[BucketSpec],
+      options: Map[String, String],
+      numStaticPartitionCols: Int): Seq[SortOrder] = {
+    V1WritesUtils.getSortOrder(
+      outputColumns,
+      partitionColumns,
+      bucketSpec,
+      options,
+      numStaticPartitionCols)
+  }
+
   override def broadcastInternal[T: ClassTag](sc: SparkContext, value: T): 
Broadcast[T] = {
     SparkContextUtils.broadcastInternal(sc, value)
   }
diff --git 
a/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
 
b/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
index bb0b94cc02..b0cd31be0e 100644
--- 
a/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
+++ 
b/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
@@ -27,6 +27,7 @@ import org.apache.spark.paths.SparkPath
 import org.apache.spark.sql.{AnalysisException, SparkSession}
 import org.apache.spark.sql.catalyst.{ExtendedAnalysisException, InternalRow}
 import org.apache.spark.sql.catalyst.analysis.DecimalPrecisionTypeCoercion
+import org.apache.spark.sql.catalyst.catalog.BucketSpec
 import org.apache.spark.sql.catalyst.expressions._
 import org.apache.spark.sql.catalyst.expressions.aggregate._
 import org.apache.spark.sql.catalyst.plans.{JoinType, LeftSingle}
@@ -253,6 +254,20 @@ class Spark41Shims extends SparkShims {
 
   override def enableNativeWriteFilesByDefault(): Boolean = true
 
+  override def getV1WriteRequiredOrdering(
+      outputColumns: Seq[Attribute],
+      partitionColumns: Seq[Attribute],
+      bucketSpec: Option[BucketSpec],
+      options: Map[String, String],
+      numStaticPartitionCols: Int): Seq[SortOrder] = {
+    V1WritesUtils.getSortOrder(
+      outputColumns,
+      partitionColumns,
+      bucketSpec,
+      options,
+      numStaticPartitionCols)
+  }
+
   override def broadcastInternal[T: ClassTag](sc: SparkContext, value: T): 
Broadcast[T] = {
     SparkContextUtils.broadcastInternal(sc, value)
   }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to