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


##########
spark/src/main/scala/org/apache/comet/udf/codegen/CometScalaUDFCodegen.scala:
##########
@@ -144,9 +169,14 @@ 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
+    // key serves every plan in the task. Only a kernel with a 
nondeterministic node is stored per
+    // plan.
+    val sharedKey = key.copy(planId = CometScalaUDFCodegen.NoPlan)
+    kernelCache.get(sharedKey).orElse(kernelCache.get(key)) match {

Review Comment:
   [P2] Avoid hashing large closures twice on every nondeterministic cache hit. 
For these kernels, `get(sharedKey)` always misses before `get(key)` succeeds, 
and each case-class hash traverses the entire serialized closure through 
`ByteBuffer.hashCode()`. This adds another full closure scan on every batch, 
even within one native plan where the base already returned correct results. A 
bounded native query with a captured 1 MiB lookup array slowed from a 271 ms 
median with the base dispatcher to 407 ms at this head, approximately 50%. 
Reusing the bytes hash in a disposable variant reduced this to 298 ms. Could 
the two lookups reuse a precomputed closure/schema hash while preserving plan 
isolation and deterministic sharing?
   
   Evidence: Reproduced with Spark 4.1.3/JDK 17 and the verified 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`. Set `val f = udf(new 
Lookup(Array.fill[Byte](1048576)(7)))`, then collect `spark.range(0, 1048576, 
1, 1).select(f(spark_partition_id()).as("x")).agg(sum("x"))`. After six 
warmups, nine samples gave medians of 271.107 ms with the dispatcher freshly 
compiled from base 8c783aa8, 407.237 ms with head, and 297.807 ms with the 
disposable hash-reuse variant. All returned 7,340,032. Only the dispatcher 
varied between otherwise-identical head builds. A direct 4096-row dispatcher 
benchmark also measured approximately 1.25 versus 2.28 ms per batch. Sources 
and logs are under `/tmp/pr6700-bf1f69ab-review/`, including 
`PR6700EndToEndCostSuite.scala`, `PR6700CacheCost.scala`, and 
`native-cost-*.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