andygrove opened a new pull request, #6553: URL: https://github.com/apache/datafusion-comet/pull/6553
## Which issue does this PR close? None yet. This is a draft prototype, opened to discuss the approach. ## Rationale for this change `RangeExec` (`spark.range`, SQL `range()`) stays on Spark today, so the operators above it in the same stage stay on Spark too. Comet only takes over after the first shuffle. The one existing route is `spark.comet.sparkToColumnar.enabled`, which is off by default and whose operator list includes `Range`. It converts Spark's rows to Arrow on the JVM, at about 15 ns per row in the benchmark below. ## What changes are included in this PR? - `CometRangeExec` follows the `CometLocalTableScanExec` pattern (`CometSink` + `CometNativeArrowSource`). It writes each partition's values straight into `BIGINT` Arrow batches on the JVM (`RangeArrowReader`) and hands them to the native plan above it. There are no protobuf or native changes. - Exact results: Spark's generated code and its interpreted `RangeExec.doExecute` disagree where the arithmetic overflows. For example, `range(Long.MinValue, Long.MaxValue, 1L << 62, 1)` returns 0 rows from the generated code and 4 from the interpreted path. `RangeArrowReader` ports the generated code's loop, including its 1000-value batches and wrapping `long` arithmetic, so it matches what Spark runs by default. - Plan identity: `equals`, `hashCode` and `doCanonicalize` cover the range parameters, so exchanges over different ranges are not reused, and exchanges over equal ranges still are. - New config `spark.comet.exec.range.enabled`, default `false`. When the Spark-to-Arrow conversion is also enabled, Range uses `CometRangeExec`. - The operators page moves `RangeExec` from "Not currently planned" to the Scans table as experimental. ## How are these changes tested? `CometRangeExecSuite` (11 tests) is registered in the Linux and macOS PR workflows. It covers: - edge-case bounds, steps and splits, including the cases from Spark's `DataFrameRangeSuite` - empty ranges - `DataFrameRangeSuite`'s randomized-parameters test - ordering across 7-row batches - SQL `range()` - the output-rows metric - off by default - precedence over Spark-to-Arrow - exchange reuse in both directions It passes on the default profile (Spark 4.1) and on Spark 3.5 with Scala 2.12. Replacing the generated-code loop with `start + k * step` fails the overflow case, so that test is not vacuous. I have not run Spark's SQL test suites with the operator enabled. `CometRangeBenchmark`: 64M rows, `local[1]`, M3 Max under load. Times are in ms, from the second of two runs that agreed. | Query | Spark | SparkToColumnar | CometRange | | ------------------------------------ | ----- | --------------- | ---------- | | range → sum | 34 | 963 | 104 | | range → filter → sum | 51 | 1030 | 193 | | range → xxhash64 → sum | 90 | 1021 | 198 | | range → cast to string → sum(length) | 1587 | 1582 | 754 | | range → group by 100 keys | 392 | 1217 | 365 | | range → group by 1M keys | 6008 | 2015 | 1258 | `CometRange` is faster than the SparkToColumnar route on every query. Against Spark, it wins when the operators above the range do substantial work. It loses 2-3x on queries that whole-stage codegen fuses with the range into one loop. A diagnostic split showed where that loss comes from. The JVM generation loop costs about 0.9 ns per row, and the per-batch hand-off to native costs about 25 ms at 8192-row batches. A native Rust generator would remove both, and I'm prototyping one next. -- 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]
