andygrove commented on code in PR #6700:
URL: https://github.com/apache/datafusion-comet/pull/6700#discussion_r4194899560
##########
spark/src/main/scala/org/apache/comet/udf/codegen/CometScalaUDFCodegen.scala:
##########
@@ -144,9 +154,13 @@ class CometScalaUDFCodegen extends CometUDF with Logging {
private def lookupOrCompile(
key: CometScalaUDFCodegen.CacheKey,
bytes: Array[Byte],
- specs: IndexedSeq[ArrowColumnSpec]): CometScalaUDFCodegen.CacheEntry = {
+ specs: IndexedSeq[ArrowColumnSpec],
+ partitionIndex: Int): CometScalaUDFCodegen.CacheEntry = {
assert(Thread.holdsLock(this), "lookupOrCompile must run under
this.synchronized")
- kernelCache.get(key) match {
+ // A deterministic kernel never reads the partition index, so one instance
under the planless
Review Comment:
The shared-key decision reads `boundExpr.deterministic`, and that is only
the root's flag. Spark doesn't make it transitive everywhere.
`Invoke.deterministic` is `isDeterministic &&
arguments.forall(_.deterministic)` and skips `targetObject`. Spark 4 lowers
`make_valid_utf8(x)` and `is_valid_utf8(x)` to `Invoke(x, "makeValid")` and
`Invoke(x, "isValid")`, and `CometInvoke` dispatches the whole subtree. So a
partition-seeded child under one of them still goes into the shared kernel. I
ran this on the branch:
```scala
spark.range(0, 8, 1, 2)
.selectExpr(
"id",
"make_valid_utf8(cast(spark_partition_id() AS STRING)) AS p",
"make_valid_utf8(cast(round(rand(42), 6) AS STRING)) AS r")
.coalesce(1)
```
The four rows from the second partition come back with `p = 0`, and with the
first partition's `rand` sequence continued. Spark returns `p = 1` for them.
The same query under a union matches. Could the check walk the tree instead,
for example `boundExpr.exists(!_.deterministic)`? With that change the query
matches Spark locally, and a deterministic expression still compiles once
across 16 coalesced plans. Adding this shape to the new test, guarded with
`isSpark40Plus`, would pin it.
##########
spark/src/main/scala/org/apache/comet/udf/codegen/CometScalaUDFCodegen.scala:
##########
@@ -168,14 +182,14 @@ class CometScalaUDFCodegen extends CometUDF with Logging {
}
val compiled = CometBatchKernelCodegen.compile(boundExpr, specs)
val kernel = compiled.newInstance()
- kernel.init(CometScalaUDFCodegen.currentPartitionIndex())
+ kernel.init(partitionIndex)
val outputField = CometBatchKernelCodegen.toFfiArrowField(
"codegen_result",
boundExpr.dataType,
boundExpr.nullable)
val entry =
CometScalaUDFCodegen.CacheEntry(compiled, kernel,
boundExpr.dataType, outputField)
- kernelCache.put(key, entry)
+ kernelCache.put(if (boundExpr.deterministic) sharedKey else key, entry)
Review Comment:
I agree with sunchao that this needs fixing here. On this branch,
`round(rand(42), 6)` under `coalesce(1)` over 16 partitions compiled 16
kernels, and nothing removes them before the task ends.
`CometExecIterator.close()` already has the plan id, because it is the
iterator's `id`. Could it drop that plan's entries right after
`nativeLib.releasePlan(plan)`? For example, a
`CometUdfBridge.releasePlan(taskAttemptId, id)` that calls a hook only the
dispatcher implements. A test that runs a coalesce over several partitions and
then checks what the cache holds would cover the cleanup. It would also cover
the claim that deterministic kernels are shared across plans, which no test
checks today. The caching diagram at the top of the class needs updating as
well. It still says the kernel cache is keyed on the bound expression and input
shapes, for the whole task.
--
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]