diff --git a/docs/source/contributor-guide/jvm_shuffle.md b/docs/source/contributor-guide/jvm_shuffle.md index 1226b4c82b4..c9d798a04db 100644 --- a/docs/source/contributor-guide/jvm_shuffle.md +++ b/docs/source/contributor-guide/jvm_shuffle.md @@ -54,6 +54,20 @@ JVM shuffle (`CometColumnarExchange`) is used instead of native shuffle (`CometE [Supported partition key types](native_shuffle.md#when-native-shuffle-is-used) for the exact rules. Complex types are fully supported as data columns in both implementations. +4. **Partition keys native shuffle cannot serialize**: native shuffle serializes the + partitioning expressions to protobuf, so a key expression Comet has no serde for, or whose + serde reports it incompatible, keeps the exchange off the native path. JVM shuffle has no + such requirement, because it evaluates the key on the JVM through `UnsafeProjection` and + `LazilyGeneratedOrdering`, so these exchanges land here rather than on Spark's shuffle. One + example is the `mapsort(...)` wrapper Spark 4.0 and later adds around a map used as a + shuffle key: Comet cannot serialize it for array or struct map keys, so such an exchange + becomes `CometColumnarExchange`. + +Collated strings are the exception. A hash or range key whose type is a non-`UTF8_BINARY` +collated string is declined by both Comet paths, so the exchange stays a plain Spark +`Exchange`. That fallback is deliberate: it also keeps the rest of the stage off Comet, which +is currently what stops `CometSort` from ordering such a key by raw bytes. + ## Input Handling ### Spark Row-Based Input diff --git a/docs/source/contributor-guide/native_shuffle.md b/docs/source/contributor-guide/native_shuffle.md index e1cabfd6c61..30c4b6a4f1f 100644 --- a/docs/source/contributor-guide/native_shuffle.md +++ b/docs/source/contributor-guide/native_shuffle.md @@ -72,6 +72,16 @@ Native shuffle (`CometExchange`) is selected when all of the following condition depth still disqualifies the key. The config defaults to `false` pending measurement of the nested hashing paths, so by default a complex hash key falls back to JVM shuffle. +5. **Partition key expressions convert to protobuf**: native shuffle serializes the partitioning + into its protobuf plan, so every `HashPartitioning` expression and every `RangePartitioning` + sort order has to convert through `QueryPlanSerde.exprToProto`. An expression Comet has no + serde for, or whose serde reports it incompatible, disqualifies native shuffle even when the + key's type is supported. One example is the `mapsort(...)` wrapper Spark 4.0 and later adds + around a map used as a shuffle key: `CometMapSort` supports scalar map keys only, so a map with + array or struct keys fails this check and the exchange falls back to JVM shuffle even with + nested hash keys enabled. JVM shuffle has no such requirement because it evaluates the key on + the JVM; see [When JVM Shuffle is Used](jvm_shuffle.md#when-jvm-shuffle-is-used). + ## Architecture ``` diff --git a/docs/source/user-guide/latest/understanding-comet-plans.md b/docs/source/user-guide/latest/understanding-comet-plans.md index beb5c6ebd14..95af535f7b8 100644 --- a/docs/source/user-guide/latest/understanding-comet-plans.md +++ b/docs/source/user-guide/latest/understanding-comet-plans.md @@ -363,12 +363,14 @@ use: batches are serialized via the Arrow IPC writer. - **`CometColumnarExchange`** is the **JVM columnar shuffle** path. It accepts either Spark row-based input or Comet columnar input, which makes it the - fallback when the child is not a Comet operator or when a hash/range key + fallback when the child is not a Comet operator, when a hash/range key type is not supported by native shuffle (for example, a struct or array - hash key while nested hash keys are disabled). It is still preferred over - Spark's native shuffle when Comet shuffle is enabled. Keys with a - non-default string collation are supported by neither path and use Spark's - shuffle. + hash key while nested hash keys are disabled), or when native shuffle + cannot serialize the key expression (for example the `mapsort` wrapper + Spark 4.0 and later adds around an array or struct map key). It is still + preferred over Spark's native shuffle when Comet shuffle is enabled. Keys + with a non-default string collation are supported by neither path and use + Spark's shuffle. Both paths support the same set of partitioning schemes (`HashPartitioning`, `RangePartitioning`, `RoundRobinPartitioning`, diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometShuffleExchangeExec.scala b/spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometShuffleExchangeExec.scala index 3ad5527d8cd..84a0cae098b 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometShuffleExchangeExec.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometShuffleExchangeExec.scala @@ -719,11 +719,13 @@ object CometShuffleExchangeExec val partitioning = s.outputPartitioning partitioning match { case HashPartitioning(expressions, _) => - for (expr <- expressions) { - if (QueryPlanSerde.exprToProto(expr, inputs).isEmpty) { - reasons += s"unsupported hash partitioning expression: $expr" - } - } + // Not a shuffle-correctness check: the partition id here comes from Spark's own + // collation-aware Murmur3Hash, evaluated on the JVM. What this rejection does is keep the + // whole stage off Comet, and that is what keeps CometSort away from the collated key -- + // supportedSortType only type-checks single-column sorts, so a multi-column collated sort + // otherwise reaches Comet and compares raw bytes (#1947). This covers hash and range + // partitioning only -- SinglePartition and round robin have no such check, and the sort + // gap itself is tracked in #6158. for (dt <- expressions.map(_.dataType).distinct) { if (isStringCollationType(dt)) { reasons += s"unsupported hash partitioning data type for columnar shuffle: $dt" @@ -734,11 +736,9 @@ object CometShuffleExchangeExec case RoundRobinPartitioning(_) => // we already checked that the input types are supported case RangePartitioning(orderings, _) => - for (o <- orderings) { - if (QueryPlanSerde.exprToProto(o, inputs).isEmpty) { - reasons += s"unsupported range partitioning sort order: $o" - } - } + // Same as the hash branch above: LazilyGeneratedOrdering already orders collated strings + // correctly here, and the fallback is what currently keeps CometSort off the collated + // key (#6158). for (dt <- orderings.map(_.dataType).distinct) { if (isStringCollationType(dt)) { reasons += s"unsupported range partitioning data type for columnar shuffle: $dt" diff --git a/spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala b/spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala index f26ef595a64..bff5be73b82 100644 --- a/spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala @@ -249,7 +249,7 @@ class CometExpressionSuite extends CometTestBase with AdaptiveSparkPlanHelper { CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key -> "true") { checkSparkAnswerAndFallbackReason( "select * from tbl order by 1, 2", - "unsupported range partitioning sort order") + "Sorting on floating-point values nested in arrays, structs, or maps") } } @@ -275,7 +275,7 @@ class CometExpressionSuite extends CometTestBase with AdaptiveSparkPlanHelper { CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key -> "true") { checkSparkAnswerAndFallbackReason( "select * from tbl order by 1, 2", - "unsupported range partitioning sort order") + "Sorting on floating-point values nested in arrays, structs, or maps") } } diff --git a/spark/src/test/scala/org/apache/comet/exec/CometColumnarShuffleSuite.scala b/spark/src/test/scala/org/apache/comet/exec/CometColumnarShuffleSuite.scala index 157cf972d67..dbcedb6eb59 100644 --- a/spark/src/test/scala/org/apache/comet/exec/CometColumnarShuffleSuite.scala +++ b/spark/src/test/scala/org/apache/comet/exec/CometColumnarShuffleSuite.scala @@ -33,7 +33,7 @@ import org.apache.spark.sql.comet.execution.shuffle.{CometShuffleDependency, Com import org.apache.spark.sql.execution.adaptive.{AdaptiveSparkPlanHelper, AQEShuffleReadExec, ShuffleQueryStageExec} import org.apache.spark.sql.execution.exchange.{ReusedExchangeExec, ShuffleExchangeExec} import org.apache.spark.sql.execution.joins.SortMergeJoinExec -import org.apache.spark.sql.functions.col +import org.apache.spark.sql.functions.{col, spark_partition_id, udf} import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types._ @@ -165,10 +165,11 @@ abstract class CometColumnarShuffleSuite extends CometTestBase with AdaptiveSpar } test("columnar shuffle on array/struct map key/value") { - // Spark 4.0 normalizes maps used as shuffle keys with mapsort(...). Comet's map_sort - // relies on Arrow's sort_to_indices, which only supports scalar key types, so a map - // with array or struct keys cannot be sorted natively and the shuffle falls back. - val complexKeyShuffles = if (isSpark40Plus) 0 else 1 + // Spark 4.0 normalizes maps used as shuffle keys with mapsort(...), which Comet cannot + // serialize for array or struct map keys. That used to disqualify columnar shuffle, but the + // columnar path computes partition ids on the JVM from h.partitionIdExpression, mapsort and + // all, so the verdict never applied to it (#5971). Partition assignment is pinned below. + val complexKeyShuffles = 1 Seq("false", "true").foreach { execEnabled => Seq(10, 201).foreach { numPartitions => Seq("1.0", "10.0").foreach { ratio => @@ -816,12 +817,9 @@ abstract class CometColumnarShuffleSuite extends CometTestBase with AdaptiveSpar * exchange operators. */ test("range partitioning on floating-point uses columnar shuffle under strictFloatingPoint") { - // The columnar path partitions on the JVM with Spark's own RangePartitioner, but it still - // probes whether Comet can serialize the sort order. Before #5506 that probe reported - // Incompatible for scalar float and double under strict mode, so this exchange fell back to - // Spark's shuffle for no compatibility reason. Consulting a native-serde gate on a path that - // never goes native is tracked separately in #5971; this test pins the floating-point half of - // it, which #5506 fixed. + // Before #5506 the columnar path reported Incompatible for scalar float and double under + // strict mode, so this exchange fell back to Spark's shuffle for no compatibility reason. + // #5971 removed the native-serde probe that caused it; the nested case is covered below. withSQLConf( CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key -> "true", CometConf.getExprAllowIncompatConfigKey("SortOrder") -> "false") { @@ -837,6 +835,105 @@ abstract class CometColumnarShuffleSuite extends CometTestBase with AdaptiveSpar } } + // A struct containing a double is orderable by the JVM, but CometSortOrder reports Incompatible + // for nested floating point under strict mode. The columnar path never sends the sort order to + // native code, so that verdict must not gate it (#5971). + test("range partitioning on a nested floating-point key uses columnar shuffle") { + withSQLConf( + CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key -> "true", + CometConf.getExprAllowIncompatConfigKey("SortOrder") -> "false") { + withParquetTable((0 until 20).map(i => (i.toDouble, i)), "range_struct_tbl") { + val df = sql("SELECT struct(_1 AS a, _2 AS b) AS c FROM range_struct_tbl") + .repartitionByRange(4, col("c")) + checkShuffleAnswer(df, 1) + } + } + } + + // A Scala UDF has no serde once the codegen dispatcher is off, so exprToProto returns None for + // it. The columnar path evaluates partition keys on the JVM through UnsafeProjection, so that + // verdict must not gate the exchange (#5971). Turning the dispatcher off is only a cheap and + // stable way to obtain an unserializable partition key; any expression without a serde reaches + // the same gate, and unlike the strictFloatingPoint cases above nothing here depends on strict + // mode. + test("range partitioning on an unserializable expression uses columnar shuffle") { + withSQLConf(CometConf.COMET_SCALA_UDF_CODEGEN_ENABLED.key -> "false") { + withParquetTable((0 until 20).map(i => (i, i.toString)), "range_udf_tbl") { + val bump = udf((x: Int) => x + 1) + val df = sql("SELECT _1, _2 FROM range_udf_tbl").repartitionByRange(4, bump(col("_1"))) + checkShuffleAnswer(df, 1) + } + } + } + + test("hash partitioning on an unserializable expression uses columnar shuffle") { + withSQLConf(CometConf.COMET_SCALA_UDF_CODEGEN_ENABLED.key -> "false") { + withParquetTable((0 until 20).map(i => (i, i.toString)), "hash_udf_tbl") { + val bump = udf((x: Int) => x + 1) + val df = sql("SELECT _1, _2 FROM hash_udf_tbl").repartition(4, bump(col("_1"))) + checkShuffleAnswer(df, 1) + } + } + } + + test("collation introduced above the scan still falls back to Spark's shuffle") { + assume(isSpark40Plus, "string collation requires Spark 4.0+") + withParquetTable((0 until 20).map(i => (i, (i % 4).toString)), "coll_above_tbl") { + // The stored column is a plain string, so CometScanRule keeps the scan native and the + // collation is applied in a Project above it. The shuffle itself would be fine -- the + // partition id comes from Spark's collation-aware Murmur3Hash on the JVM -- but declining + // the exchange is what keeps the stage, and so CometSort, off the collated key (#1947). + val df = sql("SELECT _1, _2 COLLATE UTF8_LCASE AS c FROM coll_above_tbl") + .repartition(4, col("c")) + checkShuffleAnswer(df, 0) + } + } + + /** + * checkShuffleAnswer only compares the query answer, which is order-insensitive and so would + * pass even if Comet routed rows to different partitions than Spark. Compare + * spark_partition_id() per row instead, and pin the Comet run to one CometShuffleExchangeExec: + * without that, a future fallback would leave both sides on plain Spark and the comparison + * would pass while testing nothing. + */ + private def checkPartitionAssignmentMatchesSpark(df: => DataFrame, clue: String): Unit = { + checkCometExchange(df, 1, false) + + def pids: Array[(Int, Int)] = + df.select(col("_1"), spark_partition_id().as("pid")) + .collect() + .map(r => (r.getInt(0), r.getInt(1))) + .sorted + + val cometRows = pids + var sparkRows: Array[(Int, Int)] = Array.empty + withSQLConf(CometConf.COMET_ENABLED.key -> "false") { + sparkRows = pids + } + assert(sparkRows.nonEmpty, "Spark produced no rows; the comparison would be vacuous") + assert(cometRows === sparkRows, s"partition assignment differs from Spark for $clue") + } + + test("columnar shuffle partition assignment matches Spark for keys Comet cannot serialize") { + withSQLConf(CometConf.COMET_SCALA_UDF_CODEGEN_ENABLED.key -> "false") { + withParquetTable((0 until 100).map(i => (i, i % 7)), "assign_udf_tbl") { + val bump = udf((x: Int) => x % 5) + checkPartitionAssignmentMatchesSpark( + sql("SELECT * FROM assign_udf_tbl").repartition(8, bump(col("_2"))), + "a UDF hash key") + } + } + } + + test("columnar shuffle partition assignment matches Spark for an array/struct map key") { + assume(isSpark40Plus, "mapsort normalization of map shuffle keys requires Spark 4.0+") + withParquetTable((0 until 100).map(i => (Map(Seq(i % 9, i % 4) -> i), i)), "assign_map_tbl") { + checkPartitionAssignmentMatchesSpark( + sql("SELECT _2 AS _1, _1 AS m FROM assign_map_tbl").repartition(8, col("m")), + "an array/struct map hash key") + } + } + private def checkShuffleAnswer(df: DataFrame, expectedNum: Int): Unit = { checkCometExchange(df, expectedNum, false) checkSparkAnswer(df)