From 3716fcffa5ec34735a45b03d622f60ee7a91a79a Mon Sep 17 00:00:00 2001 From: wforget <643348094@qq.com> Date: Wed, 15 Jul 2026 15:13:11 +0800 Subject: [PATCH 1/6] test GLUTEN-12474 --- .../GlutenV1WriteCommandSuite.scala | 35 +++++++++++++++++++ 1 file changed, 35 insertions(+) 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 5fc887d8d41..6291ed3ac74 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") { From 19d6bb665f4fbdd4de688dbf25f3c93b98c29811 Mon Sep 17 00:00:00 2001 From: wforget <643348094@qq.com> Date: Wed, 15 Jul 2026 15:44:11 +0800 Subject: [PATCH 2/6] Preserve V1 write ordering when inserting local sorts --- .../EnsureLocalSortRequirements.scala | 24 +++++++++++++++++-- .../apache/gluten/sql/shims/SparkShims.scala | 13 +++++++++- .../sql/shims/spark34/Spark34Shims.scala | 15 ++++++++++++ .../sql/shims/spark35/Spark35Shims.scala | 15 ++++++++++++ .../sql/shims/spark40/Spark40Shims.scala | 15 ++++++++++++ .../sql/shims/spark41/Spark41Shims.scala | 15 ++++++++++++ 6 files changed, 94 insertions(+), 3 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 e17a8e74603..515442cde9a 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 @@ -17,10 +17,12 @@ package org.apache.gluten.extension.columnar 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 /** * This rule is similar with `EnsureRequirements` but only handle local `SortExec`. @@ -33,6 +35,24 @@ import org.apache.spark.sql.execution.{SortExec, SparkPlan} object EnsureLocalSortRequirements extends Rule[SparkPlan] { private lazy val transform: HeuristicTransform = HeuristicTransform.static() + 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, + writeFiles.staticPartitions.size)) + case _ => plan.requiredChildOrdering + } + } + private def addLocalSort( originalChild: SparkPlan, requiredOrdering: Seq[SortOrder]): SparkPlan = { @@ -44,7 +64,7 @@ object EnsureLocalSortRequirements extends Rule[SparkPlan] { 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. 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 750bea0c41a..068ccdda6f7 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 97ff19a84a1..ae93d067a5f 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 d62fdfea193..08047c57570 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 5847e62c106..1e50984e193 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 bb0b94cc022..b0cd31be0ef 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) } From 0bbfcd567de9422b6a0bcf35bb4efbccd9ae9824 Mon Sep 17 00:00:00 2001 From: wforget <643348094@qq.com> Date: Wed, 15 Jul 2026 16:58:00 +0800 Subject: [PATCH 3/6] add unit tests for other spark version --- .../GlutenV1WriteCommandSuite.scala | 35 +++++++++++++++++++ .../GlutenV1WriteCommandSuite.scala | 35 +++++++++++++++++++ .../GlutenV1WriteCommandSuite.scala | 35 +++++++++++++++++++ 3 files changed, 105 insertions(+) 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 eb6794bba81..efd4105e689 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/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 a287f5fffb6..b8d9a1156ee 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 a287f5fffb6..b8d9a1156ee 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") { From 4d776e12956e9be1313b9222830d294d1aaa55c4 Mon Sep 17 00:00:00 2001 From: wforget <643348094@qq.com> Date: Wed, 15 Jul 2026 17:46:13 +0800 Subject: [PATCH 4/6] fix --- .../columnar/EnsureLocalSortRequirements.scala | 14 +++++++++++++- 1 file changed, 13 insertions(+), 1 deletion(-) 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 515442cde9a..036d059c367 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 @@ -23,6 +23,7 @@ import org.apache.spark.sql.catalyst.expressions.SortOrder import org.apache.spark.sql.catalyst.rules.Rule 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`. @@ -35,6 +36,17 @@ import org.apache.spark.sql.execution.datasources.WriteFilesExec 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, @@ -48,7 +60,7 @@ object EnsureLocalSortRequirements extends Rule[SparkPlan] { writeFiles.partitionColumns, writeFiles.bucketSpec, writeFiles.options, - writeFiles.staticPartitions.size)) + numStaticPartitionCols(writeFiles))) case _ => plan.requiredChildOrdering } } From 9cfff276464e221d9d8ce383dde8349d5e92446b Mon Sep 17 00:00:00 2001 From: wforget <643348094@qq.com> Date: Wed, 15 Jul 2026 17:59:46 +0800 Subject: [PATCH 5/6] fix spotless check --- .../extension/columnar/EnsureLocalSortRequirements.scala | 3 +-- 1 file changed, 1 insertion(+), 2 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 036d059c367..fb1e7be7ef0 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 @@ -42,8 +42,7 @@ object EnsureLocalSortRequirements extends Rule[SparkPlan] { val resolver = SQLConf.get.resolver val staticPartitionNames = writeFiles.staticPartitions.keys writeFiles.partitionColumns.takeWhile { - partitionColumn => - staticPartitionNames.exists(resolver(_, partitionColumn.name)) + partitionColumn => staticPartitionNames.exists(resolver(_, partitionColumn.name)) }.size } From bc63bb3be2bfaff2d339a0743b07af7ebe7808f6 Mon Sep 17 00:00:00 2001 From: wforget <643348094@qq.com> Date: Fri, 17 Jul 2026 18:06:07 +0800 Subject: [PATCH 6/6] Enhance local sort handling for GlutenPlan support --- .../columnar/EnsureLocalSortRequirements.scala | 13 +++++++++++-- 1 file changed, 11 insertions(+), 2 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 fb1e7be7ef0..d22a71e7e92 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,6 +16,7 @@ */ 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 @@ -65,11 +66,19 @@ object EnsureLocalSortRequirements extends Rule[SparkPlan] { } 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 = { @@ -82,7 +91,7 @@ object EnsureLocalSortRequirements extends Rule[SparkPlan] { if (SortOrder.orderingSatisfies(child.outputOrdering, requiredOrdering)) { child } else { - addLocalSort(child, requiredOrdering) + addLocalSort(p, child, requiredOrdering) } } p.withNewChildren(newChildren)