Skip to content

perf: generate the in-memory cache reader's class once per executor - #6860

Open
andygrove wants to merge 3 commits into
apache:mainfrom
andygrove:perf/cache-reader-codegen-once
Open

andygrove wants to merge 3 commits into
apache:mainfrom
andygrove:perf/cache-reader-codegen-once

Conversation

@andygrove

Copy link
Copy Markdown
Member

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 CometColumnarToRowExec regenerated its projection's source in every partition when it ran outside whole-stage codegen, which a schema with more fields than spark.sql.codegen.maxFields always does. Comet's in-memory cache reader, CachedBatchRowIterator from #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?

  • GeneratedClassCache holds the executor-wide cache that perf: generate the columnar-to-row projection once per executor instead of per partition #6858 added inside CometUnsafeProjection: the most recently used 100 classes, keyed by a ColumnLayout of column types, nullability and spark.sql.codegen.methodSplitThreshold. CometUnsafeProjection now uses it.
  • CachedBatchRowIterator generates 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. The spark.sql.codegen.hugeMethodLimit check runs on every call against the size recorded with the class, so a class generated under one limit is still checked under another.
  • The fallback for a reader with a method over that limit builds its projection through CometUnsafeProjection, so it is not generated in every partition either.

How are these changes tested?

New tests in CachedBatchRowIteratorSuite check 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. CometUnsafeProjectionSuite covers the refactored projection cache. CachedBatchRowIteratorSuite, CometUnsafeProjectionSuite, CometInMemoryCacheSuite, CometInMemoryCachePruningSuite and CometExecSuite pass on Spark 4.1, CachedBatchRowIteratorSuite also on Spark 3.5, and the change compiles with strict warnings.

A Spark operator (a noop write with Comet off) reading a cached table in Comet's format, with local[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 from CometShuffleBenchmark (100 columns, max depth 6). Before is #6858; the benchmark harness is not part of this PR.

Cached table Partitions Before After
200 columns 4 483 ms 454 ms
200 columns 200 1,343 ms 796 ms
1,500 columns 4 4,992 ms 4,716 ms
1,500 columns 200 9,938 ms 5,802 ms
fuzz nested, max depth 6 4 626 ms 306 ms
fuzz nested, max depth 6 201 16,971 ms 1,390 ms

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.

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.
@github-actions github-actions Bot added enhancement New feature or request performance area:ffi Arrow FFI / JNI boundary labels Oct 10, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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: GeneratedClassCache provides 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. GenerateUnsafeProjection and 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.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1, 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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:ffi Arrow FFI / JNI boundary enhancement New feature or request performance

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants