From a667fbbdaa7cf2f4a78c806d25f3bf74fa257e45 Mon Sep 17 00:00:00 2001 From: Yuming Wang Date: Mon, 20 Jul 2026 00:42:03 +0800 Subject: [PATCH 1/3] [GLUTEN-12566][CORE] Fix BatchScanExecTransformer to propagate keyGroupedPartitioning ScanTransformerFactory.createBatchScanTransformer() does not pass keyGroupedPartitioning from the original BatchScanExec to the BatchScanExecTransformer, leaving it as None even when the source BatchScanExec has SPJ (Storage Partitioned Join) partitioning. This causes Gluten's BatchScanExecTransformer to report UnknownPartitioning instead of KeyGroupedPartitioning, defeating SPJ and inserting redundant shuffles. IcebergScanTransformer and PaimonScanTransformer already pass keyGroupedPartitioning correctly; this fixes the same gap for regular file scans (Parquet, ORC, etc.) created via ScanTransformerFactory. Additionally fix two related issues in BatchScanExecTransformer that were masked by keyGroupedPartitioning always being None: 1. doCanonicalize() does not normalize keyGroupedPartitioning. When two semantically identical scans have different expression IDs in their keyGroupedPartitioning expressions, the canonicalized plans are not equal, preventing AQE exchange reuse. 2. hashCode() does not include keyGroupedPartitioning, violating the equals/hashCode contract (equals compares it via spjParams, but hashCode omits it). Closes #12566. --- .../gluten/execution/BatchScanExecTransformer.scala | 9 +++++++-- .../apache/gluten/execution/ScanTransformerFactory.scala | 2 ++ 2 files changed, 9 insertions(+), 2 deletions(-) diff --git a/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala b/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala index a0c3bb87575..cb70ae2b86d 100644 --- a/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala +++ b/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala @@ -68,6 +68,8 @@ case class BatchScanExecTransformer( runtimeFilters = QueryPlan.normalizePredicates( runtimeFilters.filterNot(_ == DynamicPruningExpression(Literal.TrueLiteral)), output), + keyGroupedPartitioning = keyGroupedPartitioning.map( + _.map(QueryPlan.normalizeExpressions(_, output))), pushDownFilters = pushDownFilters.map(QueryPlan.normalizePredicates(_, output)) ) } @@ -194,12 +196,15 @@ abstract class BatchScanExecTransformerBase( override def equals(other: Any): Boolean = other match { case other: BatchScanExecTransformerBase => - this.pushDownFilters == other.pushDownFilters && super.equals(other) + this.pushDownFilters == other.pushDownFilters && + this.keyGroupedPartitioning == other.keyGroupedPartitioning && + super.equals(other) case _ => false } - override def hashCode(): Int = Objects.hashCode(batch, runtimeFilters, pushDownFilters) + override def hashCode(): Int = + Objects.hashCode(batch, runtimeFilters, keyGroupedPartitioning, pushDownFilters) /** Return a copy of this scan with a new output schema. */ def withOutput(newOutput: Seq[AttributeReference]): BatchScanExecTransformerBase diff --git a/gluten-substrait/src/main/scala/org/apache/gluten/execution/ScanTransformerFactory.scala b/gluten-substrait/src/main/scala/org/apache/gluten/execution/ScanTransformerFactory.scala index 711f8c69086..b2839aa43b1 100644 --- a/gluten-substrait/src/main/scala/org/apache/gluten/execution/ScanTransformerFactory.scala +++ b/gluten-substrait/src/main/scala/org/apache/gluten/execution/ScanTransformerFactory.scala @@ -44,6 +44,8 @@ object ScanTransformerFactory { batchScanExec.output, batchScanExec.scan, batchScanExec.runtimeFilters, + keyGroupedPartitioning = + SparkShimLoader.getSparkShims.getKeyGroupedPartitioning(batchScanExec), table = SparkShimLoader.getSparkShims.getBatchScanExecTable(batchScanExec) ) } From 8ce2eafafa9b720b98ade5b2a1e460d2d3892f3d Mon Sep 17 00:00:00 2001 From: Yuming Wang Date: Mon, 20 Jul 2026 17:15:49 +0800 Subject: [PATCH 2/3] format --- .../apache/gluten/execution/BatchScanExecTransformer.scala | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala b/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala index cb70ae2b86d..525e60cb660 100644 --- a/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala +++ b/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala @@ -196,9 +196,9 @@ abstract class BatchScanExecTransformerBase( override def equals(other: Any): Boolean = other match { case other: BatchScanExecTransformerBase => - this.pushDownFilters == other.pushDownFilters && this.keyGroupedPartitioning == other.keyGroupedPartitioning && - super.equals(other) + this.pushDownFilters == other.pushDownFilters && + super.equals(other) case _ => false } From 87ccde8d339ae9f64a09d1c5caaccc35d5993c76 Mon Sep 17 00:00:00 2001 From: Yuming Wang Date: Mon, 20 Jul 2026 17:43:41 +0800 Subject: [PATCH 3/3] Update BatchScanExecTransformer.scala --- .../apache/gluten/execution/BatchScanExecTransformer.scala | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala b/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala index 525e60cb660..54fd01877a0 100644 --- a/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala +++ b/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala @@ -197,8 +197,8 @@ abstract class BatchScanExecTransformerBase( override def equals(other: Any): Boolean = other match { case other: BatchScanExecTransformerBase => this.keyGroupedPartitioning == other.keyGroupedPartitioning && - this.pushDownFilters == other.pushDownFilters && - super.equals(other) + this.pushDownFilters == other.pushDownFilters && + super.equals(other) case _ => false }