comphead opened a new issue, #6528:
URL: https://github.com/apache/datafusion-comet/issues/6528
### Describe the bug
A skewed-join benchmark that reads Spark **task time from the event log**
(not wall-clock time) shows two hash-join cases where Comet uses more task time
than Spark:
1. **Shuffled hash join, nested payload, hot key, AQE off.** The single task
that receives the hot key is slower in Comet than in Spark: 679 vs 598 ms with
90% of rows on one key (0.88x), and 414 vs 366 ms at 50% (0.88x). With flat
payloads or uniform keys, Comet's slowest task is 1.6x to 2.8x faster than
Spark's. A skewed stage cannot finish before its slowest task, so on a large
cluster this query finishes later on Comet.
2. **Broadcast hash join, every key distribution.** Comet uses more total
task time than Spark in all six cases, uniform keys included: 0.64x to 0.77x
with a flat payload, 0.86x to 0.94x with a nested payload. Skew makes no
difference to either engine here, since the probe side is never shuffled by
key. Related to #6013, see the metrics below.
Numbers are the median of 5 runs, in ms. "Total" is `executorRunTime` summed
over every task of the query, "longest" is the single slowest task. Speedup is
Spark divided by Comet, so below 1.00x means Comet needs more time.
`executorRunTime` is used rather than CPU time because Comet runs part of its
native work on Tokio threads, which task CPU time does not count.
**Total task time (ms)**
| case | Spark AQE on | Spark AQE off | Comet AQE on | Comet AQE off |
speedup on | speedup off |
|---|---:|---:|---:|---:|---:|---:|
| SHJ, flat, uniform keys | 2275 | 2062 | 790 | 848 | 2.88x | 2.43x |
| SHJ, flat, hot key 50% | 2612 | 2374 | 831 | 779 | 3.14x | 3.05x |
| SHJ, flat, hot key 90% | 2027 | 1804 | 681 | 721 | 2.98x | 2.50x |
| SHJ, nested, uniform keys | 4328 | 4345 | 2281 | 2460 | 1.90x | 1.77x |
| SHJ, nested, hot key 50% | 4375 | 4065 | 2256 | 2147 | 1.94x | 1.89x |
| SHJ, nested, hot key 90% | 4252 | 3853 | 1870 | 1891 | 2.27x | 2.04x |
| BHJ, flat, uniform keys | 339 | 327 | 461 | 443 | **0.74x** | **0.74x** |
| BHJ, flat, hot key 50% | 321 | 322 | 438 | 419 | **0.73x** | **0.77x** |
| BHJ, flat, hot key 90% | 256 | 268 | 403 | 388 | **0.64x** | **0.69x** |
| BHJ, nested, uniform keys | 1096 | 1046 | 1173 | 1195 | **0.93x** |
**0.88x** |
| BHJ, nested, hot key 50% | 1089 | 1036 | 1246 | 1134 | **0.87x** |
**0.91x** |
| BHJ, nested, hot key 90% | 950 | 1052 | 1102 | 1118 | **0.86x** |
**0.94x** |
**Longest task (ms)**
| case | Spark AQE on | Spark AQE off | Comet AQE on | Comet AQE off |
speedup on | speedup off |
|---|---:|---:|---:|---:|---:|---:|
| SHJ, flat, uniform keys | 58 | 53 | 16 | 19 | 3.63x | 2.79x |
| SHJ, flat, hot key 50% | 77 | 150 | 14 | 85 | 5.50x | 1.76x |
| SHJ, flat, hot key 90% | 51 | 219 | 13 | 135 | 3.92x | 1.62x |
| SHJ, nested, uniform keys | 99 | 104 | 46 | 46 | 2.15x | 2.26x |
| SHJ, nested, hot key 50% | 110 | 366 | 36 | 414 | 3.06x | **0.88x** |
| SHJ, nested, hot key 90% | 110 | 598 | 33 | 679 | 3.33x | **0.88x** |
| BHJ, flat, uniform keys | 12 | 12 | 16 | 16 | 0.75x | 0.75x |
| BHJ, flat, hot key 50% | 12 | 12 | 15 | 15 | 0.80x | 0.80x |
| BHJ, flat, hot key 90% | 9 | 11 | 14 | 14 | 0.64x | 0.79x |
| BHJ, nested, uniform keys | 39 | 35 | 41 | 42 | 0.95x | 0.83x |
| BHJ, nested, hot key 50% | 42 | 34 | 41 | 38 | 1.02x | 0.89x |
| BHJ, nested, hot key 90% | 33 | 38 | 38 | 37 | 0.87x | 1.03x |
The SHJ "Comet AQE on" numbers with a hot key are not Comet's join, see
"Additional context". Every other plan, including both BHJ settings, is fully
native.
#### Metrics: SHJ hot task, nested payload, 90% hot key, AQE off
One run, per-task SQL metric updates from the event log.
**Spark, 584 ms** (CPU 585 ms). One fused `WholeStageCodegen` pipeline (583
ms). `ShuffledHashJoin` builds 1.5 MiB in 1 ms. Shuffle read is 153.2 MiB,
1,893,939 rows.
**Comet, 661 ms** (CPU 657 ms):
| operator | metric | value |
|---|---|---:|
| `CometExchange` (reduce-side read) | local bytes read / records | 101.1
MiB / 1,893,939 |
| `CometExchange` | `native shuffle writer time` (as reported in the reduce
task) | 107.1 ms |
| `CometExchange` | decoding and decompression time | 58.8 ms |
| `CometHashJoin` | build side | 8,249 rows in 8 batches, 0.5 MiB, 0.1 ms |
| `CometHashJoin` | probe side / output | 1,893,939 rows in 256 batches /
470 batches |
| `CometHashJoin` | Total time for joining | 130.6 ms |
| `CometProject` | time | 0.1 ms |
| `CometColumnarToRow` (JVM) | batches to rows | 235 to 1,893,939, **no time
metric** |
| `WholeStageCodegen (1)` | duration | 658 ms |
Native operator metrics cover about 240 ms of the 661 ms. The remaining ~420
ms is not covered by any metric. The one operator without a timer is
`CometColumnarToRow`, which turns the struct, array and map output back into
rows for the sink, so it is the likely place. Without a metric that is an
inference, not a measurement. With a flat payload, the same hot task is 1.6x
faster in Comet than in Spark (135 vs 219 ms).
#### Metrics: BHJ, flat payload, uniform keys, AQE off
One run, 40 tasks each (32 probe tasks and 8 tasks that scan `dim` for the
broadcast).
- **Spark**: 325 ms total, median task 8 ms, longest 13 ms.
- **Comet**: 436 ms total, median task 13 ms, longest 16 ms.
- Every Comet probe task consumes the whole broadcast build side: 262,144
rows in one batch, 8.4 MiB of build-side memory per task. That is 8,388,608
build rows and 268.6 MiB per query, as #6013 describes.
- But the native build is small: `Total time for collecting build-side of
join` is 14.4 ms per query (0.45 ms per task), and `Total time for joining` is
77.2 ms (2.4 ms per task). So caching the built hash table (#6013, level 2)
would recover little of the ~110 ms gap at this size.
- The rest of the gap, about 3 ms per task, is not covered by any metric.
Decoding the broadcast IPC batches in each task (`CometBatchRDD.compute`),
importing them into the native plan, and per-task plan setup have no timers, so
I can't say which of them it is.
### Steps to reproduce
Data, written with Comet disabled:
- `fact`: 2,097,152 rows in 32 Parquet files, each its own map task
(`spark.sql.files.openCostInBytes` = `spark.sql.files.maxPartitionBytes` = 128
MiB). Key columns put 0%, 50% or 90% of rows on key 0 and hash the rest over
262,143 keys. Payload is either flat (`BIGINT`, `DOUBLE`, `STRING`) or nested
(`STRUCT<id BIGINT, customer STRUCT<name STRING, tier INT>, items
ARRAY<STRUCT<sku BIGINT, qty INT>>>`, `ARRAY<STRING>`, `MAP<STRING, BIGINT>`).
- `dim`: 262,144 rows with unique `k` and two flat columns.
Queries, run with `.write.format("noop")`:
```sql
-- SHJ
SELECT /*+ SHUFFLE_HASH(d) */ f.<key>, f.<payload columns>, d.d_long, d.d_str
FROM fact f JOIN dim d ON f.<key> = d.k
-- BHJ
SELECT /*+ BROADCAST(d) */ f.<key>, f.<payload columns>, d.d_long, d.d_str
FROM fact f JOIN dim d ON f.<key> = d.k
```
Configuration: `local[8]`, `spark.sql.shuffle.partitions=32`,
`spark.memory.offHeap.enabled=true`, `spark.memory.offHeap.size=8g`,
`spark.sql.autoBroadcastJoinThreshold=-1` and
`spark.sql.adaptive.autoBroadcastJoinThreshold=-1` (the BHJ asks by hint).
AQE-on runs set
`spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=4MB` and
`spark.sql.adaptive.advisoryPartitionSizeInBytes=1MB`. Each run gets its own
job group, and task metrics are read back from the event log.
The benchmark (`CometSkewedJoinBenchmark`) is not in a PR yet.
Environment: Comet `main` at d5de7b9927, Spark 4.1.3, Scala 2.13, JDK
17.0.19, macOS 15.7.4 on an Apple M3 Max, native release build.
### Expected behavior
Comet's hash joins should not need more task time than Spark's for the same
query, and the shuffled hash join's hot task should keep the speedup Comet
shows on its other tasks. At the least, the time should be attributable from
the operator metrics.
### Additional context
- **AQE skew splits disable Comet for the join.** With AQE on and a hot key,
the skew-split join runs as Spark's `ShuffledHashJoin(skew=true)` behind
`ColumnarToRow`, with no fallback reason recorded. `CometExecRule.convertNode`
turns `ShuffleQueryStageExec(CometShuffleExchangeExec)` into a
`CometExchangeSink`, but `OptimizeSkewedJoin` wraps that stage in an
`AQEShuffleReadExec` before Comet's rule runs, and there is no case for it. So
the join above never gets a native child and is never attempted. That is a
separate problem from this issue (and not #6442, which is about a declined join
losing its reason).
- Possible directions, not verified: a time metric on `CometColumnarToRow`
and `CometNativeColumnarToRow`, and a metric for the per-task broadcast decode
and import. Both would show where the unattributed time goes.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]