Skip to content
Open
Show file tree
Hide file tree
Changes from 4 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
Original file line number Diff line number Diff line change
Expand Up @@ -170,10 +170,18 @@ object GpuRangePartitioner {

case class GpuRangePartitioner(
rangeBounds: Array[InternalRow],
sorter: GpuSorter) extends GpuExpression with ShimExpression with GpuPartitioning {
sorter: GpuSorter,
boundarySorter: Option[GpuSorter] = None,
boundaryInputProjection: Option[Seq[Expression]] = None)
extends GpuExpression with ShimExpression with GpuPartitioning {

require(boundarySorter.isDefined == boundaryInputProjection.isDefined,
"boundary sorter and input projection must be specified together")

private lazy val rangeBoundsSorter = boundarySorter.getOrElse(sorter)

private lazy val converters = new GpuRowToColumnConverter(
TrampolineUtil.fromAttributes(sorter.projectedBatchSchema))
TrampolineUtil.fromAttributes(rangeBoundsSorter.projectedBatchSchema))

override def nullable: Boolean = false
override def dataType: DataType = IntegerType
Expand All @@ -189,9 +197,19 @@ case class GpuRangePartitioner(
// Don't make this retry-block avoiding nested try-blocks
// from computeBoundsAndCloseWithRetry
withResource(converters.convertBatch(rangeBounds,
TrampolineUtil.fromAttributes(sorter.projectedBatchSchema))) { ranges =>
withResource(sorter.appendProjectedColumns(cb)) { withExtraColumns =>
sorter.lowerBound(ranges, withExtraColumns)
TrampolineUtil.fromAttributes(rangeBoundsSorter.projectedBatchSchema))) { ranges =>
boundaryInputProjection match {
case Some(projectList) =>
withResource(GpuProjectExec.project(cb, projectList)) { boundaryInput =>
withResource(rangeBoundsSorter.appendProjectedColumns(boundaryInput)) {
withExtraColumns =>
rangeBoundsSorter.lowerBound(ranges, withExtraColumns)
}
}
case None =>
withResource(sorter.appendProjectedColumns(cb)) { withExtraColumns =>
sorter.lowerBound(ranges, withExtraColumns)
}
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -840,6 +840,15 @@ val GPU_COREDUMP_PIPE_PATTERN = conf("spark.rapids.gpu.coreDump.pipePattern")
.checkValue(v => v >= 0.0 && v <= 2.0, "The threshold must be between 0.0 and 2.0.")
.createWithDefault(0.75)

val RANGE_PARTITIONING_SAMPLE_KEYS_ONLY =
conf("spark.rapids.sql.rangePartitioning.sampleKeysOnly")
.doc("When enabled, range partitioning collects range boundaries from only the columns " +
"needed to compute the range keys when the input plan can be safely pruned. Disable " +
"this to collect boundaries from the original full-width GPU input.")
.internal()
.booleanConf
.createWithDefault(true)

val SHUFFLED_HASH_JOIN_OPTIMIZE_SHUFFLE =
conf("spark.rapids.sql.shuffledHashJoin.optimizeShuffle")
.doc("Enable or disable an optimization where shuffled build side batches are kept " +
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,165 @@
/*
* Copyright (c) 2026, NVIDIA CORPORATION.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package org.apache.spark.sql.rapids.execution

import scala.annotation.tailrec

import com.nvidia.spark.rapids.{GpuCoalesceBatches, GpuExec, GpuFilterExec, GpuProjectExec}
import com.nvidia.spark.rapids.shims.ShimUnaryExecNode

import org.apache.spark.rdd.RDD
import org.apache.spark.sql.catalyst.InternalRow
import org.apache.spark.sql.catalyst.expressions.{
Attribute, Expression, ExprId, NamedExpression, SortOrder}
import org.apache.spark.sql.execution.SparkPlan
import org.apache.spark.sql.rapids.GpuFileSourceScanExec
import org.apache.spark.sql.types.StructType
import org.apache.spark.sql.vectorized.ColumnarBatch

/**
* An auxiliary physical plan used to collect range-partition boundaries.
*
* The range exchange exposes this node as a subquery so Spark includes the narrow scan and its
* normal GPU metrics in the SQL physical plan. The exchange executes it once while constructing
* its range partitioner; the full-width exchange child remains the source of shuffled rows.
*
* This is an internal auxiliary node, not a CPU-to-GPU replacement registered with GpuOverrides.
* If the auxiliary plan cannot be built, the range exchange samples its original full-width GPU
* child instead.
*/
private[rapids] case class GpuRangeBoundaryExec(child: SparkPlan)
extends ShimUnaryExecNode with GpuExec {
override def output: Seq[Attribute] = child.output
override def nodeName: String = "GpuRangeBoundaryCollect"

override protected def doExecute(): RDD[InternalRow] =
throw new IllegalStateException(s"Row-based execution should not occur for $this")

override protected def internalDoExecuteColumnar(): RDD[ColumnarBatch] =
child.executeColumnar()
}

/**
* Builds a key-only physical plan for range-boundary collection when it is safe to do so.
* Unsupported plans fall back to sampling the full exchange input.
Comment thread
greptile-apps[bot] marked this conversation as resolved.
*/
private[rapids] object GpuRangeBoundaryPlan {
private def referencedExprIds(expressions: Seq[Expression]): Set[ExprId] =
expressions.flatMap(_.references).map(_.exprId).toSet

def build(plan: SparkPlan, ordering: Seq[SortOrder]): Option[GpuRangeBoundaryExec] = {
val required = referencedExprIds(ordering)
if (required.isEmpty) {
None
} else {
prune(plan, required).flatMap { pruned =>
val selected = pruned.output.filter(attr => required.contains(attr.exprId))
if (selected.map(_.exprId).toSet != required) {
None
} else {
val keyOnly = if (selected.length == pruned.output.length) {
pruned
} else {
GpuProjectExec(selected.toList, pruned)
}
Some(GpuRangeBoundaryExec(keyOnly))
}
}
}
}

private def selectProjectExpressions(
projectList: List[NamedExpression],
localOutputIds: Set[ExprId],
required: Set[ExprId]): Option[List[NamedExpression]] = {
val projectOutputIds = projectList.map(_.exprId).toSet
if (!required.subsetOf(projectOutputIds)) {
None
} else {
@tailrec
def dependencyClosure(needed: Set[ExprId]): List[NamedExpression] = {
val selected = projectList.filter(ne => needed.contains(ne.exprId))
val localDependencies = referencedExprIds(selected).intersect(localOutputIds)
val expanded = needed ++ localDependencies
if (expanded == needed) selected else dependencyClosure(expanded)
}

val selected = dependencyClosure(required)
// Boundary collection and shuffle input are separate executions of the source plan.
// Nondeterministic expressions can produce different keys if the executions use different
// batch boundaries, even when their seeds and partition IDs match.
if (selected.forall(_.deterministic)) Some(selected) else None
}
}

private def prune(plan: SparkPlan, required: Set[ExprId]): Option[SparkPlan] = plan match {
case project: GpuProjectExec =>
val childOutputIds = project.child.output.map(_.exprId).toSet
val projectOutputIds = project.projectList.map(_.exprId).toSet
val localOutputIds = projectOutputIds -- childOutputIds
selectProjectExpressions(project.projectList, localOutputIds, required).flatMap { selected =>
// Keep project-local dependencies here instead of requesting aliases from the child scan.
val childRequired = referencedExprIds(selected) -- localOutputIds
prune(project.child, childRequired).map { child =>
project.copy(projectList = selected.toList, child = child)
}
}

case filter: GpuFilterExec if filter.condition.deterministic =>
val childRequired = required ++ referencedExprIds(Seq(filter.condition))
prune(filter.child, childRequired).map { child =>
filter.withNewChildren(Seq(child))
}

case coalesce: GpuCoalesceBatches =>
prune(coalesce.child, required).map { child =>
coalesce.withNewChildren(Seq(child))
}

case scan: GpuFileSourceScanExec =>
pruneFileScan(scan, required)

case _ =>
None
}

private def pruneFileScan(
scan: GpuFileSourceScanExec,
required: Set[ExprId]): Option[GpuFileSourceScanExec] = {
if (!required.subsetOf(scan.output.map(_.exprId).toSet)) {
return None
}

val dataColumnCount = scan.requiredSchema.length
val dataPairs = scan.output.take(dataColumnCount).zip(scan.requiredSchema.fields)
val partitionPairs = scan.output.drop(dataColumnCount).zip(scan.readPartitionSchema.fields)
val selectedData = dataPairs.filter { case (attr, _) => required.contains(attr.exprId) }
val selectedPartitions =
partitionPairs.filter { case (attr, _) => required.contains(attr.exprId) }

val selectedOutput = selectedData.map(_._1) ++ selectedPartitions.map(_._1)
if (selectedOutput.map(_.exprId).toSet != required) {
None
} else {
val originalPartitionOutput = scan.originalOutput.drop(dataColumnCount)
Some(scan.copy(
originalOutput = selectedData.map(_._1) ++ originalPartitionOutput,
requiredSchema = StructType(selectedData.map(_._2)),
requiredPartitionSchema = Some(StructType(selectedPartitions.map(_._2))))(scan.rapidsConf))
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -31,10 +31,10 @@ import org.apache.spark.rapids.shims.GpuShuffleExchangeExec
import org.apache.spark.rdd.RDD
import org.apache.spark.serializer.Serializer
import org.apache.spark.sql.catalyst.InternalRow
import org.apache.spark.sql.catalyst.expressions.{Ascending, Attribute, SortOrder}
import org.apache.spark.sql.catalyst.expressions.{Ascending, Attribute, Expression, SortOrder}
import org.apache.spark.sql.catalyst.plans.physical.RoundRobinPartitioning
import org.apache.spark.sql.catalyst.trees.TreeNodeTag
import org.apache.spark.sql.execution.SparkPlan
import org.apache.spark.sql.execution.{ExecSubqueryExpression, SparkPlan}
import org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanExec
import org.apache.spark.sql.execution.exchange.{Exchange, ShuffleExchangeExec}
import org.apache.spark.sql.execution.metric._
Expand Down Expand Up @@ -245,6 +245,28 @@ abstract class GpuShuffleExchangeExecBase(

@transient lazy val inputBatchRDD: RDD[ColumnarBatch] = child.executeColumnar()

@transient private lazy val rangeBoundaryPlan: Option[GpuRangeBoundaryExec] =

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Non-blocking: unconditional for every GPU range exchange the pruner accepts, with no off switch, so a field regression needs a patched release to mitigate.

Suggest an internal boolean defaulting to true, shaped like spark.rapids.sql.shuffledHashJoin.optimizeShuffle. It also lets a test select between the two paths, which the suite comment needs.

@sdrp713 sdrp713 Sep 14, 2026

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. I added the default on internal setting spark.rapids.sql.rangePartitioning.sampleKeysOnly. Setting it to false restores the original full width GPU sampling path, and the disabled path is covered by a correctness test.

gpuOutputPartitioning match {
case range: GpuRangePartitioning
if RapidsConf.RANGE_PARTITIONING_SAMPLE_KEYS_ONLY.get(child.conf) =>
GpuRangeBoundaryPlan.build(child, range.gpuOrdering)
case _ =>
None
}

private def expressionSubqueries(expression: Expression): Seq[SparkPlan] = {
val nested = expression.children.flatMap(expressionSubqueries)
expression match {
case subquery: ExecSubqueryExpression => nested :+ subquery.plan
case _ => nested
}
}

// Boundary collection is an auxiliary query of the exchange. Exposing it through Spark's
// subquery mechanism makes its physical operators and native metrics part of the SQL plan.
@transient override lazy val subqueries: Seq[SparkPlan] =
expressions.flatMap(expressionSubqueries) ++ rangeBoundaryPlan.toSeq

/**
* Returns the GPU partitioning used to build the shuffle dependency. Distributions whose
* physical partition count depends on the input RDD can override this hook.
Expand Down Expand Up @@ -278,7 +300,8 @@ abstract class GpuShuffleExchangeExecBase(
additionalMetrics,
opTimeNewShuffleWrite,
descendantOpTimeMetrics,
enableOpTimeTrackingRdd)
enableOpTimeTrackingRdd,
rangeBoundaryPlan)
}

/**
Expand Down Expand Up @@ -409,7 +432,8 @@ object GpuShuffleExchangeExecBase {
additionalMetrics: Map[String, GpuMetric],
opTimeNewShuffleWrite: Option[GpuMetric] = None,
descendantOpTimeMetrics: Seq[GpuMetric] = Seq.empty,
enableOpTimeTrackingRdd: Boolean = true)
enableOpTimeTrackingRdd: Boolean = true,
rangeBoundaryPlan: Option[GpuRangeBoundaryExec] = None)
: ShuffleDependency[Int, ColumnarBatch, ColumnarBatch] = {
val isRoundRobin = newPartitioning match {
case _: GpuRoundRobinPartitioning => true
Expand Down Expand Up @@ -438,7 +462,7 @@ object GpuShuffleExchangeExecBase {
rdd
}
val partitioner: GpuExpression = getPartitioner(newRdd, outputAttributes,
newPartitioning, metrics)
newPartitioning, metrics, rangeBoundaryPlan)
val useRangeInputBatching = rangeInputBatchingEnabled &&
newPartitioning.isInstanceOf[GpuRangePartitioning]
// Inject debugging subMetrics, such as D2HTime before SliceOnCpu
Expand Down Expand Up @@ -557,16 +581,32 @@ object GpuShuffleExchangeExecBase {
rdd: RDD[ColumnarBatch],
outputAttributes: Seq[Attribute],
newPartitioning: GpuPartitioning,
metrics: Map[String, GpuMetric]): GpuExpression with GpuPartitioning = {
metrics: Map[String, GpuMetric],
rangeBoundaryPlan: Option[GpuRangeBoundaryExec]): GpuExpression with GpuPartitioning = {
newPartitioning match {
case h: GpuHashPartitioning =>
GpuBindReferences.bindReference(h, outputAttributes, metrics)
case r: GpuRangePartitioning =>
val sorter = new GpuSorter(r.gpuOrdering, outputAttributes, metrics)
val bounds = GpuRangePartitioner.createRangeBounds(r.numPartitions, sorter,
rdd, SQLConf.get.rangeExchangeSampleSizePerPartition)
val (boundaryRdd, boundarySorter, boundaryProjection) = rangeBoundaryPlan match {
case Some(plan) =>
val projectList = plan.output.map { boundaryAttr =>
val ordinal = outputAttributes.indexWhere(_.exprId == boundaryAttr.exprId)
require(ordinal >= 0,
s"Range boundary attribute $boundaryAttr is missing from the shuffle input")
GpuBoundReference(ordinal, boundaryAttr.dataType, boundaryAttr.nullable)(
boundaryAttr.exprId, boundaryAttr.name)
}
(plan.executeColumnar(),
new GpuSorter(r.gpuOrdering, plan.output, metrics), Some(projectList))
case None =>
(rdd, sorter, None)
}
val bounds = GpuRangePartitioner.createRangeBounds(r.numPartitions, boundarySorter,
boundaryRdd, SQLConf.get.rangeExchangeSampleSizePerPartition)
// No need to bind arguments for the GpuRangePartitioner. The Sorter has already done it
new GpuRangePartitioner(bounds, sorter)
new GpuRangePartitioner(bounds, sorter,
boundaryProjection.map(_ => boundarySorter), boundaryProjection)
case GpuSinglePartitioning =>
GpuSinglePartitioning
case rrp: GpuRoundRobinPartitioning =>
Expand Down
Loading
Loading