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]

Reply via email to