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]