Skip to content
Open
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
4 changes: 3 additions & 1 deletion docs/source/contributor-guide/jvm_shuffle.md
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,9 @@ JVM shuffle (`CometColumnarExchange`) is used instead of native shuffle (`CometE
range key always falls back here. `HashPartitioning` keys must be primitive only by default:
setting `spark.comet.shuffle.native.partitioning.hash.nested.enabled` to `true` keeps struct and
array keys, and map keys on Spark 4.0 and later, on the native path. The config defaults to
`false`, so a complex hash key falls back to JVM columnar shuffle unless it is enabled. See
`false`, so a complex hash key falls back to JVM columnar shuffle unless it is enabled.
Decimal hash keys with precision greater than 18, including nested decimal leaves, also use
JVM shuffle when there is more than one partition and the shuffle mode allows it. See
[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.

Expand Down
5 changes: 4 additions & 1 deletion docs/source/contributor-guide/native_shuffle.md
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,9 @@ Native shuffle (`CometExchange`) is selected when all of the following condition
Spark's `mapsort` normalization makes physical entry order irrelevant. A collated string at any
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.
Decimal leaves with precision greater than 18 also disqualify a hash key when there is more
than one output partition: native decimal hashing differs from Spark's. A one-partition hash
exchange stays native because the writer serializes it as `SinglePartition` without hashing.

## Architecture

Expand Down Expand Up @@ -505,7 +508,7 @@ independently compressed, allowing parallel decompression during reads.
| Input format | Columnar (direct from Comet operators) | Row-based (via ColumnarToRowExec) |
| Partitioning logic | Rust implementation | Spark's partitioner |
| Supported schemes | Hash, Range, Single, RoundRobin | Hash, Range, Single, RoundRobin |
| Partition key types | Primitives only (Hash, Range) | Any type |
| Partition key types | See the hash and range rules above | Any type |
| Performance | Higher (no format conversion) | Lower (columnar→row→columnar) |
| Writer variants | Single path | Bypass (hash) and sort-based |

Expand Down
11 changes: 11 additions & 0 deletions docs/source/user-guide/latest/tuning/shuffle.md
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,17 @@ scalar types. Hash partitioning keys must be scalar types unless
and later) map keys. That setting is disabled by default until the performance of the nested hashing paths has been
measured. Columns that are not partitioning keys may contain complex types like maps, structs, and arrays.

Hash partitioning into multiple partitions on decimal keys with precision greater than 18 falls back because native
hashing does not match Spark's partition assignments. Mixing the two partitioners can silently lose join rows.
The restriction also applies to decimal leaves in nested hash keys. Wider decimals remain supported as payload
columns, range partitioning keys, and in single-partition shuffles, including one-partition hash exchanges.

With `spark.comet.shuffle.mode=auto`, Comet uses Columnar Shuffle when eligible, adding a columnar-to-row-to-columnar
conversion. With `native`, it uses Spark shuffle, which can also restore aggregates such as `collect_list` and
`collect_set` to Spark to preserve buffer compatibility. Celeborn uses Spark/Celeborn shuffle for unsupported native
keys because it has no Comet columnar fallback. Spark-compatible native decimal hashing is tracked in
[#5994](https://github.com/apache/datafusion-comet/issues/5994).

### Columnar (JVM) Shuffle

Comet Columnar shuffle is JVM-based and supports `HashPartitioning`, `RoundRobinPartitioning`, `RangePartitioning`, and
Expand Down
1 change: 1 addition & 0 deletions spark/src/main/scala/org/apache/comet/serde/hash.scala
Original file line number Diff line number Diff line change
Expand Up @@ -134,6 +134,7 @@ private object HashUtils {
}

private def unsupportedReasonFor(dt: DataType): Option[String] = dt match {
// Keep in sync with CometShuffleExchangeExec's hash-key restriction until #5994 is fixed.
case d: DecimalType if d.precision > 18 => Some(unsupportedDecimalReason)
case s: StructType =>
s.fields.iterator.flatMap(f => unsupportedReasonFor(f.dataType).iterator).toSeq.headOption
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -555,12 +555,11 @@ object CometShuffleExchangeExec
_: FloatType | _: DoubleType | _: StringType | _: BinaryType | _: TimestampType |
_: TimestampNTZType | _: DateType =>
true
case _: DecimalType =>
// TODO enforce this check
// https://github.com/apache/datafusion-comet/issues/3079
// Decimals with precision > 18 require Java BigDecimal conversion before hashing
// d.precision <= 18
true
case d: DecimalType =>
// Match the SQL hash restriction in serde/HashUtils until #5994 fixes native encoding.
// Different partition assignments break mixed native/Spark joins. A single partition
// does not hash the key: CometNativeShuffleWriter serializes it as SinglePartition.
d.precision <= 18 || s.outputPartitioning.numPartitions == 1
case dt if isTimeType(dt) =>
true
case StructType(fields) if nestedHashPartitioningEnabled =>
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
CometColumnarToRow
+- CometTakeOrderedAndProject
+- CometHashAggregate
+- CometExchange
+- CometColumnarExchange
+- CometHashAggregate
+- CometUnion
:- CometProject
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ CometColumnarToRow
+- CometSort
+- CometExchange
+- CometHashAggregate
+- CometExchange
+- CometColumnarExchange
+- CometHashAggregate
+- CometUnion
:- CometHashAggregate
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
CometColumnarToRow
+- CometTakeOrderedAndProject
+- CometHashAggregate
+- CometExchange
+- CometColumnarExchange
+- CometHashAggregate
+- CometUnion
:- CometHashAggregate
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
CometColumnarToRow
+- CometTakeOrderedAndProject
+- CometHashAggregate
+- CometExchange
+- CometColumnarExchange
+- CometHashAggregate
+- CometUnion
:- CometHashAggregate
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ CometColumnarToRow
+- CometSort
+- CometExchange
+- CometHashAggregate
+- CometExchange
+- CometColumnarExchange
+- CometHashAggregate
+- CometUnion
:- CometHashAggregate
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
CometColumnarToRow
+- CometTakeOrderedAndProject
+- CometHashAggregate
+- CometExchange
+- CometColumnarExchange
+- CometHashAggregate
+- CometUnion
:- CometProject
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
CometColumnarToRow
+- CometTakeOrderedAndProject
+- CometHashAggregate
+- CometExchange
+- CometColumnarExchange
+- CometHashAggregate
+- CometUnion
:- CometHashAggregate
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ CometColumnarToRow
+- CometSort
+- CometExchange
+- CometHashAggregate
+- CometExchange
+- CometColumnarExchange
+- CometHashAggregate
+- CometUnion
:- CometHashAggregate
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
CometColumnarToRow
+- CometTakeOrderedAndProject
+- CometHashAggregate
+- CometExchange
+- CometColumnarExchange
+- CometHashAggregate
+- CometUnion
:- CometHashAggregate
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
CometColumnarToRow
+- CometTakeOrderedAndProject
+- CometHashAggregate
+- CometExchange
+- CometColumnarExchange
+- CometHashAggregate
+- CometUnion
:- CometHashAggregate
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ CometColumnarToRow
+- CometSort
+- CometExchange
+- CometHashAggregate
+- CometExchange
+- CometColumnarExchange
+- CometHashAggregate
+- CometUnion
:- CometHashAggregate
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,10 @@

package org.apache.comet

import org.apache.spark.sql.execution.aggregate.HashAggregateExec
import org.apache.spark.sql.execution.exchange.ShuffleExchangeExec
import org.apache.spark.sql.types.DecimalType

import org.apache.comet.DataTypeSupport.isComplexType

class CometFuzzAggregateSuite extends CometFuzzTestBase {
Expand All @@ -31,7 +35,16 @@ class CometFuzzAggregateSuite extends CometFuzzTestBase {
val (_, cometPlan) = checkSparkAnswer(sql)
assert(1 == collectNativeScans(cometPlan).length)

checkSparkAnswerAndOperator(sql)
val hasWideDecimalKey = df.schema(col).dataType match {
case d: DecimalType => d.precision > 18
case _ => false
}
// Wide decimal hash keys require Spark shuffle when columnar shuffle is disabled.
if (CometConf.COMET_SHUFFLE_MODE.get() == "native" && hasWideDecimalKey) {
checkSparkAnswerAndOperator(sql, classOf[HashAggregateExec], classOf[ShuffleExchangeExec])
} else {
checkSparkAnswerAndOperator(sql)
}
}
}

Expand All @@ -55,7 +68,18 @@ class CometFuzzAggregateSuite extends CometFuzzTestBase {
val (_, cometPlan) = checkSparkAnswer(sql)
assert(1 == collectNativeScans(cometPlan).length)

checkSparkAnswerAndOperator(sql)
val hasWideDecimalKey = Seq("c1", "c2", "c3", col).exists { key =>
df.schema(key).dataType match {
case d: DecimalType => d.precision > 18
case _ => false
}
}
// Check both GROUP BY and DISTINCT keys, not unrelated decimal payload columns.
if (CometConf.COMET_SHUFFLE_MODE.get() == "native" && hasWideDecimalKey) {
checkSparkAnswerAndOperator(sql, classOf[HashAggregateExec], classOf[ShuffleExchangeExec])
} else {
checkSparkAnswerAndOperator(sql)
}
}
}

Expand Down
24 changes: 20 additions & 4 deletions spark/src/test/scala/org/apache/comet/CometFuzzTestSuite.scala
Original file line number Diff line number Diff line change
Expand Up @@ -160,6 +160,17 @@ class CometFuzzTestSuite extends CometFuzzTestBase {
}

test("distribute by single column (complex types)") {
// Inspect only the key schema: any wide decimal leaf requires Spark's hash partition
// assignments, even when native hashing of nested keys is enabled.
def hasWideDecimal(dataType: DataType): Boolean = dataType match {
case decimal: DecimalType => decimal.precision > 18
case StructType(fields) => fields.exists(field => hasWideDecimal(field.dataType))
case ArrayType(elementType, _) => hasWideDecimal(elementType)
case MapType(keyType, valueType, _) =>
hasWideDecimal(keyType) || hasWideDecimal(valueType)
case _ => false
}

val df = spark.read.parquet(filename)
df.createOrReplaceTempView("t1")
val columns = df.schema.fields.filter(f => isComplexType(f.dataType)).map(_.name)
Expand All @@ -180,15 +191,20 @@ class CometFuzzTestSuite extends CometFuzzTestBase {
}
assert(cometShuffleExchanges.length == expectedNumCometShuffles)

// With the config enabled these keys do run through native shuffle. This is the widest
// nested-type coverage in the repo, so it is worth asserting that they are admitted rather
// than only that they fall back.
// Enabling nested keys admits supported types, but wide decimal leaves still require
// Spark's hash partition assignments. JVM shuffle supports both kinds of key.
withSQLConf(CometConf.COMET_SHUFFLE_NATIVE_HASH_PARTITIONING_NESTED_ENABLED.key -> "true") {
val enabledDf = spark.sql(sql)
enabledDf.collect()
val enabledPlan =
enabledDf.queryExecution.executedPlan.asInstanceOf[AdaptiveSparkPlanExec].executedPlan
assert(collectCometShuffleExchanges(enabledPlan).length == 1)
val expectedEnabledShuffles =
if (CometConf.COMET_SHUFFLE_MODE.get() == "native" &&
hasWideDecimal(df.schema(col).dataType)) 0
else 1
assert(
collectCometShuffleExchanges(enabledPlan).length == expectedEnabledShuffles,
s"Unexpected shuffle for ${df.schema(col)}")
}
}
}
Expand Down
Loading
Loading