dwsmith1983 opened a new pull request, #6712: URL: https://github.com/apache/datafusion-comet/pull/6712
## Which issue does this PR close? Closes #6711. ## Rationale for this change The JVM codegen dispatcher keeps one compiled kernel per task for each distinct pair of serialized expression bytes and column specs. The kernel holds the deserialized expression, and a Catalyst `Nondeterministic` node such as `monotonically_increasing_id()` keeps its counter inside it. Every dispatched subtree is bound on its own, so two occurrences of the same subtree in one projection serialize to identical bytes and share one kernel: the second occurrence continues the first one's counter. On a single-partition table with 8-row batches, `javaId(monotonically_increasing_id())` in two columns returned `a = 0..7, b = 8..15` per batch where Spark returns `a = b = id`. Spark gives every occurrence its own state. ## What changes are included in this PR? - `DispatchOccurrence`, a marker expression that carries an occurrence id. `CometScalaUDF.emitJvmCodegenDispatch` wraps the bound subtree in it before serializing when the tree contains a `Nondeterministic` node, taking the id from a JVM-wide counter on the driver. The id travels in the native plan bytes, so every batch and every task of a plan carries the same id for one occurrence, and two occurrences never share one. Deterministic subtrees are not wrapped, so identical ones keep sharing a compiled kernel. - `CometScalaUDFCodegen.lookupOrCompile` strips the wrapper right after deserializing, before compiling. The cache key is unchanged; the id changes the bytes. Two identical nondeterministic projections now serialize to different native plan bytes, so AQE may no longer treat them as equal for exchange reuse. That only loses a reuse Spark would not perform either, since Spark gives the two projections distinct state. Out of scope: user functions that keep their own state in one object Spark shares between calls (a `ScalaUDF` marked nondeterministic, `Invoke`, `StaticInvoke`) are not Catalyst `Nondeterministic` nodes and are not wrapped. #5526 covers them separately. ## How are these changes tested? Three tests in `CometCodegenSuite`: - `identical non-deterministic dispatched expressions keep their own state`: the query above with a Java UDF (no encoders, so both calls serialize identically) and a Scala UDF (whose per-occurrence encoders already kept the calls apart), compared with Spark across eight batches. Fails on `main` with `a = 0..7, b = 8..15`. - `non-deterministic dispatched subtrees serialize per occurrence`: two dispatches of a nondeterministic tree produce different payload bytes; a deterministic tree is not wrapped and its payload equals the plain serialization of the bound tree. - `dispatcher keeps one kernel per occurrence id and its state across batches`: two ids give two cache entries, a repeated id gives one hit, and the counter continues across batches for one id. `CometCodegenSuite`, `CometNativeUdfSuite`, `CometScalaUDFClassLoaderSuite` and `CometUdfBridgeSuite` pass on the Spark 4.1 and 3.5 profiles. -- 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]
