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 @@ -27,7 +27,7 @@ import org.apache.spark.serializer.Serializer
import org.apache.spark.sql.catalyst.InternalRow
import org.apache.spark.sql.catalyst.expressions.Attribute
import org.apache.spark.sql.catalyst.plans.logical.Statistics
import org.apache.spark.sql.catalyst.plans.physical._
import org.apache.spark.sql.catalyst.plans.physical.{SinglePartition, _}
import org.apache.spark.sql.catalyst.util.truncatedString
import org.apache.spark.sql.execution.exchange._
import org.apache.spark.sql.execution.metric.SQLShuffleWriteMetricsReporter
Expand Down Expand Up @@ -93,14 +93,21 @@ abstract class ColumnarShuffleExchangeExecBase(
var cachedShuffleRDD: ShuffledColumnarBatchRDD = _

override protected def doValidateInternal(): ValidationResult = {
BackendsApiManager.getValidatorApiInstance
val validation = BackendsApiManager.getValidatorApiInstance
.doColumnarShuffleExchangeExecValidate(output, outputPartitioning, child)
.map {
reason =>
ValidationResult.failed(
s"Found schema check failure for schema ${child.schema} due to: $reason")
}
.getOrElse(ValidationResult.succeeded)
if (validation.nonEmpty) {
return ValidationResult.failed(
s"Found schema check failure for schema ${child.schema} due to: ${validation.get}")
}
outputPartitioning match {
case _: HashPartitioning => ValidationResult.succeeded
case _: RangePartitioning => ValidationResult.succeeded
case SinglePartition => ValidationResult.succeeded
case _: RoundRobinPartitioning => ValidationResult.succeeded
case _ =>
ValidationResult.failed(
s"Unsupported partitioning ${outputPartitioning.getClass.getSimpleName}")
}
}

override def numMappers: Int = inputColumnarRDD.getNumPartitions
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -62,16 +62,20 @@ class VeloxTestSettings extends BackendTestSettings {
enableSuite[GlutenFileDataSourceV2FallBackSuite]
// Rewritten
.exclude("Fallback Parquet V2 to V1")
// TODO: fix in Spark-4.0
// enableSuite[GlutenKeyGroupedPartitioningSuite]
// // NEW SUITE: disable as they check vanilla spark plan
// .exclude("partitioned join: number of buckets mismatch should trigger shuffle")
// .exclude("partitioned join: only one side reports partitioning")
// .exclude("partitioned join: join with two partition keys and different # of partition keys")
// // disable due to check for SMJ node
// .excludeByPrefix("SPARK-41413: partitioned join:")
// .excludeByPrefix("SPARK-42038: partially clustered:")
// .exclude("SPARK-44641: duplicated records when SPJ is not triggered")
enableSuite[GlutenKeyGroupedPartitioningSuite]
// NEW SUITE: disable as they check vanilla spark plan
.exclude("partitioned join: number of buckets mismatch should trigger shuffle")
.exclude("partitioned join: only one side reports partitioning")
.exclude("partitioned join: join with two partition keys and different # of partition keys")
.excludeByPrefix("SPARK-47094")
.excludeByPrefix("SPARK-48655")
.excludeByPrefix("SPARK-48012")
.excludeByPrefix("SPARK-44647")
.excludeByPrefix("SPARK-41471")
// disable due to check for SMJ node
.excludeByPrefix("SPARK-41413: partitioned join:")
.excludeByPrefix("SPARK-42038: partially clustered:")
.exclude("SPARK-44641: duplicated records when SPJ is not triggered")
enableSuite[GlutenLocalScanSuite]
enableSuite[GlutenMetadataColumnSuite]
enableSuite[GlutenSupportsCatalogOptionsSuite]
Expand Down
Loading