comphead opened a new issue, #6485: URL: https://github.com/apache/datafusion-comet/issues/6485
### What is the problem the feature request solves? Comet evaluates native expressions through DataFusion's `PhysicalExpr` tree. Each node processes the whole batch and materializes its result, so `a + b * c` writes a temporary array for `b * c`, and a filter copies the columns it keeps before the projection reads them. Spark avoids this kind of overhead with whole-stage codegen. Comet has JVM-side codegen (the Janino dispatcher from #4267), but nothing that fuses native expression evaluation. This issue is to assess whether generating native kernels for expression trees is worth pursuing, and if so, how. ### Describe the potential solution Generate one kernel per fusable expression subtree (or filter plus projection) at native plan time, compile it, and call it from a `PhysicalExpr`. Options to evaluate: - Generate Rust source and compile it with `rustc` on the executor, loading the result through a C ABI similar to the one in #4459. - Compile once on the driver and ship the binary to executors. - JIT in process (Cranelift or LLVM). DataFusion shipped a Cranelift JIT in 2022 and removed it in apache/datafusion#6164, because Cranelift could not inline Rust functions and the generated code was slower. - Kernels pre-generated at build time for common shapes, with no runtime compilation. Supporting pieces: run the interpreted path until a kernel is ready, cache compiled kernels (per process, per node, per cluster), pass literals as parameters so one kernel serves many queries, and let DataFusion evaluate subexpressions the generator does not support. ### Additional context **Proof of concept.** Hand-written kernels in the shape a generator would emit: an `extern "C"` function over raw Arrow buffers that writes into buffers the host allocates, called from a `PhysicalExpr` through a function pointer. Baselines are built the way Comet plans these expressions today (`BinaryExpr`, `CaseExpr`, `filter_record_batch`, arrow's checked kernels for ANSI). 8192-row batches with 10% nulls, DataFusion 55.1, arrow 59.3, Apple Silicon, single thread. Every fused result matched DataFusion's. | Expression | DataFusion | Fused | Speedup | |---|---|---|---| | `a * (1 - b) * (1 + c) + d` (double) | 10.1 µs | 3.4 µs | 2.9-3.0x | | `(x + y) * (z - w) + x * 3` (long, legacy) | 13.3 µs | 5.7 µs | 1.9-2.3x | | same, ANSI | 33.8 µs | 14.0 µs | 2.2-2.4x | | `CASE WHEN ... THEN ... ELSE ... END`, 1 branch | 94.3 µs | 8.3 µs | 11.4x | | `CASE`, 3 branches | 150.7 µs | 12.4 µs | 12.2x | | filter, then project 2 columns | 21.2 µs | 11.3 µs | 1.9-2.1x | Times are per batch. Speedups are the range over three runs. Most of the `CASE` gap does not need codegen. A plain vectorized select over 64 rows at a time brings the two `CASE` cases to 13.3 µs and 30.2 µs, 5 to 7x faster than `CaseExpr`, but that is only valid when no branch can fail. Fusion adds another 1.6 to 2.4x on top. Arrow's `zip` is slow on this input (63 µs and 122 µs). **Compile cost.** `rustc -O` compiled a dependency-free generated kernel in 0.07 to 0.17 s (up to 200 fused expressions) and one using `arrow` in 0.14 to 0.46 s. At about 0.1 s per compile, an arithmetic kernel pays for itself only after roughly 40 to 120 million rows, so compiled kernels have to be cached and reused. A Rust toolchain is about 1.4 GB, which matters for executor images. **To assess:** - [ ] Measure how much time real queries (TPC-H, TPC-DS) spend in expression evaluation, to bound the end-to-end gain. - [ ] Compare compile strategies (`rustc` on executors, `rustc` on the driver, JIT, pre-generated kernels) on latency, deployment, and code quality. - [ ] Decide how to handle branches that can fail (ANSI arithmetic, casts), where a vectorized plan has to filter per branch but a fused kernel can decide row by row. - [ ] Define how generated code is tested for Spark parity, for example differential testing against the interpreted path. - [ ] Consider improvements that need no codegen first: a faster primitive select in arrow's `zip`, and a `CaseExpr` fast path for branches that cannot fail. -- 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]
