Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -822,7 +822,7 @@ object VeloxConfig extends ConfigRegistry {
.createWithDefault(false)

val CUDF_ENABLE_VALIDATION =
buildStaticConf("spark.gluten.sql.columnar.backend.velox.cudf.enableValidation")
buildConf("spark.gluten.sql.columnar.backend.velox.cudf.enableValidation")
.doc(
"Heuristics you can apply to validate a cuDF/GPU plan and only offload when " +
"the entire stage can be fully and profitably executed on GPU")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,9 @@ case class AppendBatchResizeForShuffleInputAndOutput(isAdaptiveContext: Boolean)
extends Rule[SparkPlan] {
override def apply(plan: SparkPlan): SparkPlan = {
val resizeBatchesShuffleInputEnabled = VeloxConfig.get.veloxResizeBatchesShuffleInput
val resizeBatchesShuffleOutputEnabled = VeloxConfig.get.veloxResizeBatchesShuffleOutput
// TODO: Move cudf resize batches into shuffle reader.
val resizeBatchesShuffleOutputEnabled =
VeloxConfig.get.veloxResizeBatchesShuffleOutput || VeloxConfig.get.enableColumnarCudf
if (!resizeBatchesShuffleInputEnabled && !resizeBatchesShuffleOutputEnabled) {
return plan
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -85,21 +85,21 @@ object AdjustStageExecutionMode extends Logging {
// TODO: support BroadcastQueryStageExec.
case aqeShuffleRead @ AQEShuffleReadExec(s @ ShuffleQueryStageExec(_, _, _), _)
if s.shuffle.isInstanceOf[ColumnarShuffleExchangeExec] =>
ColumnarAQEShuffleReadExec(
Left(aqeShuffleRead),
stageExecutionMode)
ColumnarAQEShuffleReadExec(aqeShuffleRead, stageExecutionMode)
case queryStageExec: ShuffleQueryStageExec
if queryStageExec.shuffle.isInstanceOf[ColumnarShuffleExchangeExec] =>
ColumnarAQEShuffleReadExec(
Right(queryStageExec),
stageExecutionMode)
ColumnarAQEShuffleReadExec(queryStageExec, stageExecutionMode)
case shuffle: ColumnarShuffleExchangeExec =>
shuffle
.copy(mapperStageMode = Some(stageExecutionMode))
.withNewChildren(Seq(adjustExecutionMode(shuffle.child, stageExecutionMode)))
case resizeBatches: VeloxResizeBatchesExec =>
case r: VeloxResizeBatchesExec
// TODO: This should be removed after merging resize into native shuffle read.
// Only change the execution mode for shuffle reader.
if r.child.isInstanceOf[ShuffleQueryStageExec] ||
r.child.isInstanceOf[AQEShuffleReadExec] =>
VeloxResizeBatchesExec(
adjustExecutionMode(resizeBatches.child, stageExecutionMode),
adjustExecutionMode(r.child, stageExecutionMode),
Some(stageExecutionMode))
Comment thread
marin-ma marked this conversation as resolved.
case _ =>
plan.withNewChildren(plan.children.map(adjustExecutionMode(_, stageExecutionMode)))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ import org.apache.spark.SparkConf
import org.apache.spark.shuffle.GlutenShuffleUtils
import org.apache.spark.sql.{DataFrame, Row}
import org.apache.spark.sql.execution._
import org.apache.spark.sql.execution.adaptive.{AdaptiveSparkPlanHelper, AQEShuffleReadExec, ShuffleQueryStageExec}
import org.apache.spark.sql.execution.adaptive.{AdaptiveSparkPlanHelper, AQEShuffleReadExec, ColumnarAQEShuffleReadExec, ShuffleQueryStageExec}
import org.apache.spark.sql.execution.joins.BaseJoinExec
import org.apache.spark.sql.execution.window.WindowExec
import org.apache.spark.sql.functions._
Expand Down Expand Up @@ -2199,6 +2199,49 @@ class MiscOperatorSuite extends VeloxWholeStageTransformerSuite with AdaptiveSpa
})
}

test("Check VeloxResizeBatches is added in ShuffleRead when cuDF is enabled") {
Seq(true, false).foreach(
coalesceEnabled => {
withSQLConf(
GlutenConfig.COLUMNAR_CUDF_ENABLED.key -> "true",
VeloxConfig.CUDF_ENABLE_VALIDATION.key -> "false",
VeloxConfig.COLUMNAR_VELOX_RESIZE_BATCHES_SHUFFLE_OUTPUT.key -> "false",
SQLConf.SHUFFLE_PARTITIONS.key -> "10",
SQLConf.COALESCE_PARTITIONS_ENABLED.key -> coalesceEnabled.toString
) {
runQueryAndCompare(
"SELECT l_orderkey, count(1) from lineitem group by l_orderkey".stripMargin) {
df =>
val executedPlan = getExecutedPlan(df)
if (coalesceEnabled) {
// VeloxResizeBatches(AQEShuffleRead(ShuffleQueryStage(ColumnarShuffleExchange)))
assert(executedPlan.sliding(4).exists {
case Seq(
_: ColumnarShuffleExchangeExec,
_: ShuffleQueryStageExec,
ColumnarAQEShuffleReadExec(AQEShuffleReadExec(_, _), _),
_: VeloxResizeBatchesExec
) =>
true
case _ => false
})
} else {
// VeloxResizeBatches(ShuffleQueryStage(ColumnarShuffleExchange))
assert(executedPlan.sliding(4).exists {
case Seq(
_: ColumnarShuffleExchangeExec,
_: ShuffleQueryStageExec,
ColumnarAQEShuffleReadExec(ShuffleQueryStageExec(_, _, _), _),
_: VeloxResizeBatchesExec) =>
true
case _ => false
})
}
}
}
})
}

test("RowToVeloxColumnar preferredBatchBytes") {
Seq("1", "80", "100000000").foreach(
preferredBatchBytes => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,8 @@ import org.apache.gluten.config.{GlutenConfig, VeloxConfig}
import org.apache.spark.SparkConf
import org.apache.spark.sql.Row
import org.apache.spark.sql.execution.ColumnarShuffleExchangeExec
import org.apache.spark.sql.execution.adaptive.{ColumnarAQEShuffleReadExec, ShuffleQueryStageExec}
import org.apache.spark.sql.execution.adaptive.{AQEShuffleReadExec, ColumnarAQEShuffleReadExec, ShuffleQueryStageExec}
import org.apache.spark.sql.execution.exchange.ShuffleExchangeExec
import org.apache.spark.sql.internal.SQLConf

class StageExecutionModeSuite extends VeloxWholeStageTransformerSuite {
Expand Down Expand Up @@ -108,25 +109,27 @@ class StageExecutionModeSuite extends VeloxWholeStageTransformerSuite {

shuffleReaders.foreach {
reader =>
val canonicalized = reader.canonicalized
// canonicalized plan before applying query stage optimizer rules.
assert(canonicalized.children.forall(_.isInstanceOf[ShuffleExchangeExec]))
assert(
reader.executionMode == MockGPUStageMode,
s"Expected GPU AQE shuffle reader, but got ${reader.executionMode}")
}

val shuffleStages = plan.collect {
case stage: ShuffleQueryStageExec => stage
val shuffleStages: Seq[ShuffleQueryStageExec] = shuffleReaders.map(_.delegate).map {
case a: AQEShuffleReadExec =>
assert(a.child.isInstanceOf[ShuffleQueryStageExec])
a.child.asInstanceOf[ShuffleQueryStageExec]
case s: ShuffleQueryStageExec => s
case _ =>
throw new IllegalArgumentException("Unexpected child of ColumnarAQEShuffleReadExec")
}

val exchanges = shuffleStages.flatMap {
_.plan.collect {
case exchange: ColumnarShuffleExchangeExec => exchange
}
}

assert(exchanges.nonEmpty)

exchanges.foreach {
exchange =>
shuffleStages.foreach {
shuffleStage =>
assert(shuffleStage.shuffle.isInstanceOf[ColumnarShuffleExchangeExec])
val exchange = shuffleStage.shuffle.asInstanceOf[ColumnarShuffleExchangeExec]
assert(
!exchange.mapperStageMode.contains(MockGPUStageMode),
Comment thread
marin-ma marked this conversation as resolved.
s"Expected CPU mapper stage, but got ${exchange.mapperStageMode}")
Expand Down
2 changes: 1 addition & 1 deletion docs/velox-configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ nav_order: 16
| spark.gluten.sql.columnar.backend.velox.cudf.batchSize | 🔄 Dynamic | 2147483647 | Cudf input batch size after shuffle reader |
| spark.gluten.sql.columnar.backend.velox.cudf.concurrentGpuTasks | ⚓ Static | 1 | The number of concurrent GPU tasks to run. |
| spark.gluten.sql.columnar.backend.velox.cudf.enableTableScan | ⚓ Static | false | Enable cudf table scan |
| spark.gluten.sql.columnar.backend.velox.cudf.enableValidation | ⚓ Static | true | Heuristics you can apply to validate a cuDF/GPU plan and only offload when the entire stage can be fully and profitably executed on GPU |
| spark.gluten.sql.columnar.backend.velox.cudf.enableValidation | 🔄 Dynamic | true | Heuristics you can apply to validate a cuDF/GPU plan and only offload when the entire stage can be fully and profitably executed on GPU |
| spark.gluten.sql.columnar.backend.velox.cudf.memoryPercent | ⚓ Static | 50 | The initial percent of GPU memory to allocate for memory resource for one thread. |
| spark.gluten.sql.columnar.backend.velox.cudf.memoryResource | ⚓ Static | async | GPU RMM memory resource. |
| spark.gluten.sql.columnar.backend.velox.cudf.shuffleMaxPrefetchBytes | 🔄 Dynamic | 1028MB | Maximum bytes to prefetch in CPU memory during GPU shuffle read while waiting for GPU available. |
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,49 +31,59 @@ import org.apache.spark.sql.vectorized.ColumnarBatch
* ShuffleQueryStageExec if executionMode is set by the planner.
*
* @param delegate
* The AQEShuffleReadExec or ShuffleQueryStageExec.
* AQEShuffleReadExec, ShuffleQueryStageExec, or (during canonicalization) ShuffleExchange.
* @param executionMode
* The execution mode of the current AQE stage.
*/
case class ColumnarAQEShuffleReadExec(
delegate: Either[AQEShuffleReadExec, ShuffleQueryStageExec],
delegate: SparkPlan,
executionMode: StageExecutionMode) extends UnaryExecNode {

override def nodeName: String = s"ColumnarAQEShuffleRead(${executionMode.name})"

private val isAQEShuffleRead = delegate.isLeft

private val aqeReader: AQEShuffleReadExec = {
if (isAQEShuffleRead) {
delegate.left.get
} else {
// Wrap ShuffleQueryStageExe with dummy PartitionSpecs.
val queryStageExec = delegate.right.get
// Create CoalescedPartitionSpec for each partition.
val partitionSpecs =
Array.tabulate(queryStageExec.shuffle.numPartitions)(i => CoalescedPartitionSpec(i, i + 1))
AQEShuffleReadExec(queryStageExec, partitionSpecs)
}
}

override def supportsColumnar: Boolean = true

override def child: SparkPlan = aqeReader.child
override def child: SparkPlan = delegate match {
case AQEShuffleReadExec(c, _) => c
case _ => delegate
}

override def output: Seq[Attribute] = aqeReader.child.output
override def output: Seq[Attribute] = delegate.output

override lazy val outputPartitioning: Partitioning = aqeReader.outputPartitioning
override lazy val outputPartitioning: Partitioning = delegate.outputPartitioning
Comment thread
marin-ma marked this conversation as resolved.

override def stringArgs: Iterator[Any] = aqeReader.stringArgs
override protected def stringArgs: Iterator[Any] = {
delegate match {
case a: AQEShuffleReadExec => a.stringArgs
case _ => super.stringArgs
}
}

@transient override lazy val metrics: Map[String, SQLMetric] = aqeReader.metrics
override protected def withNewChildInternal(newChild: SparkPlan): ColumnarAQEShuffleReadExec = {
delegate match {
case a: AQEShuffleReadExec => copy(delegate = a.withNewChildren(Seq(newChild)))
case _ => copy(delegate = newChild)
}
}

private def isCoalescedSpec(spec: ShufflePartitionSpec) = {
val method = classOf[AQEShuffleReadExec].getDeclaredMethod("isCoalescedSpec")
method.setAccessible(true)
method.invoke(aqeReader, spec).asInstanceOf[Boolean]
private lazy val aqeReader: AQEShuffleReadExec = {
delegate match {
case a: AQEShuffleReadExec => a
case s: ShuffleQueryStageExec =>
// Wrap ShuffleQueryStageExe with dummy PartitionSpecs by creating CoalescedPartitionSpec
// for each partition.
val partitionSpecs =
Array.tabulate(s.shuffle.numPartitions)(i => CoalescedPartitionSpec(i, i + 1))
AQEShuffleReadExec(s, partitionSpecs)
case _ =>
// The child is Exchange during canonicalization.
throw new IllegalStateException(
s"Cannot get aqeReader from delegate node ${delegate.nodeName}.")
}
}

@transient override lazy val metrics: Map[String, SQLMetric] = aqeReader.metrics

private def shuffleStage = {
val method = classOf[AQEShuffleReadExec].getDeclaredMethod("shuffleStage")
method.setAccessible(true)
Expand All @@ -89,7 +99,8 @@ case class ColumnarAQEShuffleReadExec(
private lazy val shuffleRDD: RDD[_] = {
shuffleStage match {
case Some(stage) =>
if (isAQEShuffleRead) {
// Only send driver metrics if it's a wrapper for AQEShuffleRead.
if (delegate.isInstanceOf[AQEShuffleReadExec]) {
sendDriverMetrics()
}
stage.shuffle match {
Expand All @@ -108,15 +119,4 @@ case class ColumnarAQEShuffleReadExec(
override protected def doExecuteColumnar(): RDD[ColumnarBatch] = {
shuffleRDD.asInstanceOf[RDD[ColumnarBatch]]
}

override protected def withNewChildInternal(newChild: SparkPlan): ColumnarAQEShuffleReadExec = {
if (isAQEShuffleRead) {
copy(delegate =
Left(delegate.left.get.withNewChildren(Seq(newChild)).asInstanceOf[AQEShuffleReadExec]))
} else {
copy(delegate =
Right(
delegate.right.get.withNewChildren(Seq(newChild)).asInstanceOf[ShuffleQueryStageExec]))
}
}
}
Loading