Repository navigation
Conversation
CometColumnarToRowExec.doExecute built its UnsafeProjection with UnsafeProjection.create in every partition, which regenerates the projection's Java source each time. The plan runs that path whenever it is outside whole-stage codegen, which a schema with more nested fields than spark.sql.codegen.maxFields always is. For 100 columns nested three levels deep, generating the source took about 19 ms per partition, half the time of a 201-partition JVM shuffle of 1,000 rows. Cache the compiled class per executor by column layout and create a new instance for each partition.
Spark reads Comet's in-memory cache through CachedBatchRowIterator, once per partition, whenever it does not read the cache as batches. That is always the case for a schema with more fields than spark.sql.codegen.maxFields. Every call generated the reader's source again; Spark's compiled class cache only saved the compile. For a 1,500-column table that was about 18 ms per partition, and for 100 columns nested three levels deep about 78 ms. Keep the class per executor by column layout, as CometUnsafeProjection does for CometColumnarToRowExec, with the cache both now use moved into GeneratedClassCache. Each reader gets its own batches through its own copy of the references. The fallback for a method over the huge-method limit builds its projection through CometUnsafeProjection.
sunchao
left a comment
There was a problem hiding this comment.
Reviewed the full nine-file comparison from b94d52f43c9edaf341f5c2cc00fa2e351ba9222b to 31106da3d22cdcad9b2b6dd2f3404cdb73ec5c2d, including both prerequisite commits. The PR is not a draft. Read the supplied snapshot and live reviews, conversation comments, inline comments, and threads. PR #6860 has no existing discussion. The prerequisite’s earlier suite-registration concern is fixed.
Applied review-comet-pr, review-comet-ffi-pr, review-comet-memory-pr, and review-comet-shuffle-pr, including their contributor documentation. The JNI difference reverses the base-only #6695 refactor. I reviewed that difference and confirmed the PR commits leave jni_api.rs unchanged relative to their common ancestor.
No introduced P1/P2 issues found within this review.
Summary
- Prior state and problem: Both non-whole-stage columnar-to-row conversion and cached-batch row reading regenerated Java source per partition. Spark’s compiled-class cache avoided compilation but retained source-generation costs for wide and nested schemas.
- Design approach:
GeneratedClassCacheprovides bounded class reuse for both paths. Each caller still receives a separate projection or reader instance. - Correctness: Row serialization continues through Spark’s
GenerateUnsafeProjection.createCode. Checks passed for independent readers, projection-buffer isolation, nulls, nested values, decimals, scalar boundaries, empty input, and wide schemas. Cached readers clone their references array before installing the caller’s batches. Source inspection found no new batch-retention or release issue. - Compatibility analysis: Compared relevant Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0.
GenerateUnsafeProjectionand vector-access generation agree across these versions. Spark’s factory-mode handling remains authoritative. Also checked Spark’s columnar-to-row evaluator and the cache-scan width gate. Local runtime validation covered 4.1.3 only. - Key design decisions: Each cache retains at most 100 entries. Keys include column types, nullability, and method-splitting threshold. Projections are shared only without captured references, and readers only with the replaceable batches reference. The huge-method limit is checked on every request, including cache hits.
- Implementation sketch: The prerequisite redirects projection construction to
CometUnsafeProjection. The final commit extracts the common cache, caches reader classes and method-size statistics, and uses the cached projection helper for oversized-reader fallback. Both workflow matrices register the new projection suite. - Performance: Generation-count checks confirm warm calls avoid regenerating source. Compilation occurs outside the cache lock, allowing concurrent cold misses to generate independently. The author reports the nested 201-partition cache benchmark improving from 16,971 ms to 1,390 ms. I did not rerun that benchmark. No evidence-backed performance regression was identified.
- Design: Sharing compiled classes while keeping mutable execution state local is appropriate for reuse across partitions. The reviewed JNI comparison preserves the execution, metrics, callback-registration, and cleanup behavior of the base refactor.
- Abstraction & complexity: A small generic cache serves two concrete callers. Spark still owns serialization logic. The duplicated projection class template is limited to obtaining a generated class rather than an instance. No complexity concern met the P1/P2 bar.
- Behavioral changes worth calling out: Compared with
branch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1, class reuse is the intended change. The cache reader’s reusable-row behavior already exists at the supplied base and is not introduced here. No unintended result, error, or fallback change was identified. - Suggested improvements: No additional improvement meeting the P1/P2 reporting bar was identified. The prerequisite’s requested Linux/macOS suite registration is present and passes the inventory check.
Exact-head CI at 2026-10-10 16:54 UTC: Preflight, lint, Scala syntactic lint, CodeQL, Analyze Actions, title checking, and labeling passed. Native build/tests and compatibility checks were running or queued. No failures were reported. Spark SQL, Iceberg, macOS, and benchmark jobs were skipped. The CI run had no final verdict.
Validation and limits: Compiled the three exact-head Scala helpers against Spark 4.1.3, Scala 2.13.17, and JDK 21. Disposable harnesses passed seven projection checks and twelve cached-reader checks, including 1,500-column input. dev/ci/check-suites.py passed. The reader harness used on-heap vectors and omitted the Arrow-backed lifecycle tests. Full Comet/native builds, query integration, cross-version runtime suites, and benchmarks were not run locally. Project files remain unchanged.
Which issue does this PR close?
Part of #6855. Stacked on #6858; only the last commit is new.
Rationale for this change
#6858 found that
CometColumnarToRowExecregenerated its projection's source in every partition when it ran outside whole-stage codegen, which a schema with more fields thanspark.sql.codegen.maxFieldsalways does. Comet's in-memory cache reader,CachedBatchRowIteratorfrom #5859, has the same cost. Spark reads a cache through it once per partition whenever it does not read the cache as batches, which is always the case for such a schema, and every call generated the reader's source again; Spark's compiled class cache only saved the compile. With the cache on by default since #5634, any Spark operator reading a wide or deeply nested cached relation pays it in every partition.What changes are included in this PR?
GeneratedClassCacheholds the executor-wide cache that perf: generate the columnar-to-row projection once per executor instead of per partition #6858 added insideCometUnsafeProjection: the most recently used 100 classes, keyed by aColumnLayoutof column types, nullability andspark.sql.codegen.methodSplitThreshold.CometUnsafeProjectionnow uses it.CachedBatchRowIteratorgenerates its reader class once per layout. The class takes its batches as a reference, so each reader gets its own batches in its own copy of the references. A class is shared only when the batches are its one reference, which is always the case today. Thespark.sql.codegen.hugeMethodLimitcheck runs on every call against the size recorded with the class, so a class generated under one limit is still checked under another.CometUnsafeProjection, so it is not generated in every partition either.How are these changes tested?
New tests in
CachedBatchRowIteratorSuitecheck that a second reader for a layout generates no class and still reads only its own batches, and that a class generated under a 1-byte huge-method limit falls back there but is used without being generated again once the limit allows it.CometUnsafeProjectionSuitecovers the refactored projection cache.CachedBatchRowIteratorSuite,CometUnsafeProjectionSuite,CometInMemoryCacheSuite,CometInMemoryCachePruningSuiteandCometExecSuitepass on Spark 4.1,CachedBatchRowIteratorSuitealso on Spark 3.5, and the change compiles with strict warnings.A Spark operator (a
noopwrite with Comet off) reading a cached table in Comet's format, withlocal[1]on a Ryzen 9 7950X3D, JDK 17 and Spark 4.1. Times are the best of at least five iterations. The flat tables hold 200,000 rows of nullable bigint columns; the nested one is the 1,000-row fuzz-generated schema fromCometShuffleBenchmark(100 columns, max depth 6). Before is #6858; the benchmark harness is not part of this PR.The savings match the time each call spent generating the reader's source: called directly, after the first call compiled the class, the old reader took about 5 ms per call for 200 columns and 18 ms for 1,500 (Spark's compiled class cache hit every time), against under 0.5 ms now. For the nested schema it was about 78 ms per partition.