andygrove opened a new pull request, #6709:
URL: https://github.com/apache/datafusion-comet/pull/6709

   ## Which issue does this PR close?
   
   Closes #6705.
   
   ## Rationale for this change
   
   The codegen dispatcher paid a cost on every batch that grew with the size of 
the serialized expression. For each batch, `CometScalaUDFCodegen.evaluate` 
copied the closure-serialized expression out of its first argument and hashed 
the copy byte by byte to find the kernel. That is 7,349 bytes for `(x: Long) => 
x + 1`. The native side copied the same bytes twice more on the way to the JVM: 
`Literal::evaluate` clones its value, and `to_array_of_size(1)` copies the 
clone into the array sent over FFI.
   
   @mbutrovich measured a fixed cost of about 8 µs per batch in the dispatcher 
on #6697 
([comment](https://github.com/apache/datafusion-comet/pull/6697#issuecomment-6006552028)),
 about 6 µs of it between calling `evaluate` and calling its kernel directly.
   
   ## What changes are included in this PR?
   
   - The serde computes a SHA-256 digest of the serialized expression once on 
the driver and sends it as the dispatcher's first argument, ahead of the bytes. 
The dispatcher keys its kernel cache on the 32-byte digest and reads the 
serialized expression only to compile a kernel on a cache miss. This is the 
first option the `perf-cache-key` TODO listed, and the TODO is removed. The 
digest is SHA-256 rather than a cheaper hash because a hit is trusted without 
comparing the bytes. The key still identifies exactly the expressions the bytes 
did, so which calls share a kernel does not change.
   - `JvmScalarUdfExpr` builds the one-row array for each literal argument when 
it is created and sends the same array with every batch. This applies to every 
JVM UDF call, not only the dispatcher. A `CometUDF` that wrote into a literal 
argument's buffers would now change what later batches see. The [Arrow C data 
interface](https://arrow.apache.org/docs/format/CDataInterface.html) already 
says a consumer should treat imported data as immutable, and the dispatcher, 
the only `CometUDF` in Comet, only reads its inputs.
   - The two places that call the dispatcher directly, a test in 
`CometCodegenSuite` and `CometTimeExtractBenchmark`, pass the digest.
   - A new benchmark, `CometScalaUDFDispatchBenchmark`, adapted from the one 
@mbutrovich wrote for #6697 without the vectorized UDF cases, which need #6697. 
It runs `SELECT max(f(c))` end to end at three batch sizes, and calls the 
dispatcher directly next to the kernel it runs.
   
   After this change, no per-batch work depends on the size of the serialized 
expression.
   
   With `CometScalaUDFDispatchBenchmark` on an Apple M3 Ultra (JDK 17, Spark 
4.1, `make release`), calling the dispatcher directly on 1,024-row batches, 
4,096 batches per iteration (best time in ms):
   
   | | main | this PR |
   | --- | --- | --- |
   | generated kernel alone, `(x: Long) => x + 1` | 12 | 12 |
   | `CometScalaUDFCodegen.evaluate`, same UDF | 37 | 12 |
   
   On main, `evaluate` cost about 6.1 µs per batch more than the kernel it 
runs. With this PR the difference is too small for the harness to show; it 
reports whole milliseconds, which is 0.24 µs per batch here.
   
   End to end, `SELECT max(f(c))` over 4M rows, in ms above `max(c)` with no 
function (best of at least 5 iterations):
   
   | `f` | batch size | main | this PR |
   | --- | --- | --- | --- |
   | `(x: Long) => x + 1` | 1024 | 97 | 70 |
   | `(x: Long) => x + 1` | 8192 | 39 | 35 |
   | `(x: Long) => x + 1` | 65536 | 30 | 28 |
   | `(x: java.lang.Long) => ...` | 1024 | 173 | 132 |
   
   The saving is a fixed amount per batch, about 6.6 µs at batch size 1024, so 
it matters most when batches are small. Of the primitive case's 27 ms there, 
the digest accounts for about 23 ms and building literal arguments once for 
about 4 ms. I ran the digest-only and complete versions twice each and saw that 
4 ms in both runs, but it is close to the noise. The boxed case, whose 
serialized expression is 8,853 bytes, is noisier, so only its batch size 1024 
row is shown.
   
   The other two costs @mbutrovich found, Spark's null guard around a primitive 
parameter and the encoder behind a boxed one, are #6704 and #6706.
   
   ## How are these changes tested?
   
   - A new test in `CometCodegenSuite` calls the dispatcher twice with the same 
digest. The second call passes a null in place of the serialized expression, so 
it passes only if the digest alone finds the kernel the first call compiled. 
The test also checks that a call in the old layout, with the serialized 
expression first, is refused instead of being read as a digest.
   - A new Rust unit test checks that a literal argument's array is built when 
the expression is created and that a column argument's is not.
   - Existing suites that run the dispatcher pass locally: `CometCodegenSuite`, 
`CometCodegenSourceSuite`, `CometCodegenHOFSuite`, 
`CometScalaUDFClassLoaderSuite`, `CometUdfBridgeSuite`, 
`CometSpecializedGettersDispatchSuite`, `CometRegExpJvmSuite`, 
`CometJsonJvmSuite` and `CometDecimalPromotionSuite`.
   - Labeled `run-spark-4.1-tests`, because Spark's own SQL tests send many 
expressions through the dispatcher, and `run-benchmark-check` for the new 
benchmark.
   


-- 
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