comphead opened a new issue, #6806: URL: https://github.com/apache/datafusion-comet/issues/6806
### What is the problem the feature request solves? With a 1,048,576-row build side, a Comet broadcast hash join probe task takes 31.2 ms, against 8.5 ms in Spark (mean over 32 probe tasks of 65,536 rows each, same data and query). About three quarters of the Comet task decodes the broadcast in the JVM, and every probe task decodes all of it again. Profile of one Comet probe task, run alone about 800 times, from a 1 ms stack sampler on the task thread plus macOS `sample` for native frames: | Where the Comet probe task's time goes | Share | About | |---|---:|---:| | lz4-java decompression of the broadcast | 58% | 18 ms | | ↳ of which xxhash32 block checksums | 14% | 4.4 ms | | `Channels.newChannel` copy loop in `Utils.decodeBatches` | 8% | 2.6 ms | | Other JVM decode and hand-off (`DataInputStream`, Arrow IPC read, export to native, releases) | 8% | 2.4 ms | | Native Parquet scan of the probe file | 7% | 2.2 ms | | Native hash table build (`try_create_array_map`, min and max of the build keys) | 6% | 1.9 ms | | Join probe and output | 5% | 1.5 ms | | UTF-8 check of the imported build batch | 2% | 0.7 ms | | Other | 6% | 1.9 ms | The native build matches the `Total time for collecting build-side of join` metric, 2.84 ms per task. Spark's probe task is about half Parquet scan and 40% output row writes. Spark builds its hash table once, on the driver. The decode path, at 9a6241f140: - `CometBroadcastExchangeExec` broadcasts Arrow IPC batches compressed with Spark's `spark.io.compression.codec`, lz4 by default ([`Utils.serializeBatches`](https://github.com/apache/datafusion-comet/blob/9a6241f140e9fd576ba9467ee26c2e21cf46ec02/spark/src/main/scala/org/apache/spark/sql/comet/util/Utils.scala#L254), [`Utils.coalesceBroadcastBatches`](https://github.com/apache/datafusion-comet/blob/9a6241f140e9fd576ba9467ee26c2e21cf46ec02/spark/src/main/scala/org/apache/spark/sql/comet/util/Utils.scala#L339), [`broadcastInternal`](https://github.com/apache/datafusion-comet/blob/9a6241f140e9fd576ba9467ee26c2e21cf46ec02/spark/src/main/scala/org/apache/spark/sql/comet/CometBroadcastExchangeExec.scala#L170)). - Spark's `TorrentBroadcast` compresses the broadcast blocks again with the same codec when `spark.broadcast.compress` is on, which is the default. So the payload is compressed twice, and the inner layer is the one every task undoes. - [`CometBatchRDD.compute`](https://github.com/apache/datafusion-comet/blob/9a6241f140e9fd576ba9467ee26c2e21cf46ec02/spark/src/main/scala/org/apache/spark/sql/comet/CometBroadcastExchangeExec.scala#L300) decodes the whole broadcast in every task through [`Utils.decodeBatches`](https://github.com/apache/datafusion-comet/blob/9a6241f140e9fd576ba9467ee26c2e21cf46ec02/spark/src/main/scala/org/apache/spark/sql/comet/util/Utils.scala#L309): `codec.compressedInputStream`, then a `DataInputStream`, then [`Channels.newChannel`](https://github.com/apache/datafusion-comet/blob/9a6241f140e9fd576ba9467ee26c2e21cf46ec02/spark/src/main/scala/org/apache/spark/sql/comet/util/Utils.scala#L319) into `ArrowReaderIterator`. - Here, lz4-java ran its pure-Java decompressor (`LZ4JavaUnsafeFastDecompressor`) and verified an xxhash32 checksum on every block. The cost per task grows with the build side. In an earlier sweep of the same benchmark, without converting Comet's output to rows: | Build rows | Comet probe task minus Spark's | One decode of the serialized broadcast, timed on the driver | |---:|---:|---:| | 65,536 | −0.4 to −0.8 ms (Comet faster) | 1.0 ms | | 262,144 | 2.9–3.0 ms | 4.4–4.5 ms | | 1,048,576 | 22.3 ms | 18.7–19.5 ms | Converting the output to rows, as a query ending in a row-based operator does, adds about 2 ms to every Comet probe task. ### Describe the potential solution - Decode the broadcast once per executor and share it across tasks. That is level 1 of #6013. The draft #6037 goes further and shares the prepared build for the joins it covers. At this size, either would remove most of the gap, because the hash table build itself is about 2 ms. - Make each decode cheaper, whether or not it is cached: - Skip Comet's own codec layer for broadcast payloads, which `TorrentBroadcast` already compresses. - Or decode natively. In the same benchmark's shuffle profile, the native lz4_flex decoder was about 5x faster per byte than lz4-java. - Read without `Channels.newChannel`, as #6805 does for shuffle reads. Its copy loop alone is 8% of the task. ### Additional context - Related to #6013 and #6037. At this size, this profile shows the decode, not the hash table build, dominates the per-task cost. - Found while profiling #6528. - To reproduce: - `fact`: 2,097,152 rows in 32 Parquet files, one probe task each, with columns `key` (`PMOD(HASH(id), 65536)`), `f_long`, `f_double` and `f_str`. - `dim`: 1,048,576 rows with `k`, `d_long` and `d_str`. - Query: `SELECT /*+ BROADCAST(d) */ f.key, f.f_long, f.f_double, f.f_str, d.d_long, d.d_str FROM fact f JOIN dim d ON f.key = d.k`. - Settings: AQE off, `local[8]`, `spark.memory.offHeap.size=8g`, `spark.sql.autoBroadcastJoinThreshold=-1`. Comet's output is counted from the native plan's batches, without converting it to rows. - Measured with a local benchmark that is not in a PR (`CometHashJoinTaskTimeBenchmark -- bhj profile`). - Environment: Comet `main` at 9a6241f140 with the unmerged #6805, which changes only the shuffle readers. Spark 4.1.3, Scala 2.13, JDK 17.0.19, macOS 15.7.4 on an Apple M3 Max, release build. On Linux, lz4-java may load its JNI decompressor, which would shrink the decompression share (not verified). ### Willingness to contribute I can contribute a fix for this bug independently -- 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]
