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]
