sunchao commented on code in PR #6700:
URL: https://github.com/apache/datafusion-comet/pull/6700#discussion_r4192121583


##########
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:
   [P2] Release nondeterministic cache entries when their native plan finishes. 
Every parent partition consumed by `coalesce` gets a fresh `planId`, so this 
insertion retains another serialized expression, deserialized closure, and 
kernel. `CometUdfBridge` removes the dispatcher only at task completion, while 
`CometExecIterator.close()` does not remove these entries. A bounded 
32-partition query retained approximately 32 MiB in cache-key payloads alone 
after consuming all parents. The base dispatcher retained one approximately 1 
MiB payload for the same expression. Consequently, coalescing many partitions 
now accumulates JVM heap proportional to every completed plan, with additional 
closure copies beyond those measured bytes. Tie entry cleanup to native-plan 
lifetime while preserving state for live or interleaved plans.
   
   Evidence: Reproduced on Spark 4.1.3/JDK 21 with the exact-head native 
artifact. Define a serializable `Lookup(val data: Array[Byte]) extends (Int => 
Int)` whose `apply(i)` returns `data(Math.floorMod(i, data.length)).toInt`. 
Evaluate `val f = udf(new Lookup(Array.fill[Byte](1024 * 1024)(7)))` and 
`spark.range(0, 32, 1, 
32).select(f(spark_partition_id()).as("x")).coalesce(1)`. Inside 
`rdd.mapPartitions`, exhaust the rows, then inspect the current task's 
dispatcher cache before task completion. The plan used `CometCoalesce` over 
`CometProject` and retained 32 entries containing 33,794,624 serialized-key 
bytes. A separate identical-expression comparison against the dispatcher 
freshly compiled from base d16f7bbb measured head: 32 entries/33,651,968 bytes; 
base: 1 entry/1,051,624 bytes. Reproduction source and logs: 
`/tmp/6700-pinned-review-4qldjlpo/PR6700ReviewSuite.scala` and `tests.log`.



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