Skip to content

perf: generate the columnar-to-row projection once per executor instead of per partition - #6858

Open
andygrove wants to merge 2 commits into
apache:mainfrom
andygrove:andygrove/datafusion-comet-6855
Open

andygrove wants to merge 2 commits into
apache:mainfrom
andygrove:andygrove/datafusion-comet-6855

Conversation

@andygrove

Copy link
Copy Markdown
Member

Which issue does this PR close?

Part of #6855.

Rationale for this change

#6855 reports that JVM columnar shuffle with 201 partitions is up to 20x slower than Spark on the fuzz-generated nested schema in CometShuffleBenchmark (100 columns, 1,000 rows), and suggests a per-partition cost in the native writer. A profile of the max-depth-6 case shows that the writer is not where the time goes: the native library takes about 10% of the samples. The cost is on the reduce side, and every one of the 201 tasks pays it:

  • The largest share went to CometColumnarToRowExec.doExecute, which calls UnsafeProjection.create in every partition. That call generates the projection's Java source each time; Spark caches only the compiled class. The plan runs outside whole-stage codegen because the schema has more nested fields than spark.sql.codegen.maxFields. Spark's own scans return rows for such schemas, so Spark rarely runs ColumnarToRowExec on one, but Comet's operators always produce batches. Generating the source once per executor saves about 19 ms per partition, half the time of the case.
  • Most of the rest is importing the wide schema's Arrow arrays for each batch, and wrapping and closing the vectors. This PR does not change that.

The same per-partition cost applies to any query whose CometColumnarToRowExec runs outside whole-stage codegen, which a wide or deeply nested output always does.

In this benchmark, the "Comet (native Shuffle)" arm runs Spark's shuffle. repartition(n) is round-robin, which native shuffle declines by default (spark.comet.shuffle.native.partitioning.roundrobin.enabled), and the fuzz schema is read by Spark's scan. So that arm's parity with Spark says nothing about native shuffle.

What changes are included in this PR?

  • CometUnsafeProjection creates the projection that CometColumnarToRowExec.doExecute copies rows through. It generates and compiles the class once per executor for each column layout, meaning the column types and nullability plus spark.sql.codegen.methodSplitThreshold, the one setting the generated source depends on. It keeps up to 100 classes and returns a new instance on every call. Spark's GenerateUnsafeProjection.create returns only an instance, so the class is built from the same template, which is identical in Spark 3.4 through 4.1. It honors spark.sql.codegen.factoryMode as UnsafeProjection.create does.
  • CometColumnarToRowExec.doExecute uses it. The whole-stage codegen path and the broadcast path, which builds one projection on the driver, are unchanged.

The in-memory cache reader (CachedBatchRowIterator) also generates its code in every partition and could cache its class the same way. That is left for a follow-up.

How are these changes tested?

CometUnsafeProjectionSuite checks that:

  • rows match UnsafeProjection.create for structs, arrays and maps with nulls, including with a method split threshold small enough to split the field writes
  • a layout's class is generated once, and each call returns a separate projection
  • a different split threshold generates its own class
  • the least recently used class is evicted after 100
  • NO_CODEGEN returns the interpreted projection
  • a query whose CometColumnarToRowExec runs outside whole-stage codegen matches Spark, and a second run generates no class

CometShuffleSuite and CometExecSuite pass, and the change compiles against Spark 3.4, 3.5 (strict warnings), 4.0 and 4.1.

CometShuffleBenchmark's fuzz-generated nested schema cases, run on a Ryzen 9 7950X3D with JDK 17, Spark 4.1 and local[1]. The times are the best of five for the JVM shuffle arm, with Spark's time from the same run:

Max depth Partitions Before After Spark (before / after runs)
2 5 114 ms 114 ms 117 / 127 ms
2 201 827 ms 404 ms 93 / 117 ms
6 5 681 ms 530 ms 525 / 454 ms
6 201 7,798 ms 3,992 ms 369 / 376 ms

The depth-6, 5-partition rows have a standard deviation near 200 ms, so only the 201-partition rows show a real change. The remaining gap at 201 partitions is the per-batch Arrow import described above.

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.
@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 entire three-file diff from dc880d6d265f8391b6202efa43574cfeea71e75a to ed8ba0c95b6e7fc628d77ecf88118f15fb2cd947. The PR is not a draft. The snapshot and live discussion contained no reviews, conversation comments, inline comments, or review threads.

Applied review-comet-pr and review-comet-ffi-pr, including their contributor documentation. Found one reproducible P2 CI integration issue. No additional runtime correctness or performance issue meeting the P1/P2 bar was identified.

Summary

  • Prior state and problem: Outside whole-stage codegen, CometColumnarToRowExec regenerated projection source for every partition. Spark caches compiled code, so wide nested schemas still incurred repeated source-generation work.
  • Design approach: CometUnsafeProjection caches up to 100 generated classes per executor while constructing a separate mutable projection for each caller.
  • Correctness: The helper delegates field serialization to Spark’s GenerateUnsafeProjection.createCode. Standalone checks matched Spark for nested nulls, decimals, binary data, scalar boundaries, signed zero, NaN, timestamps, intervals, and empty output. Cache hits retained independent row buffers. No row-copy correctness defect was found.
  • Compatibility analysis: Compared Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. Their GenerateUnsafeProjection implementations are identical, including the copied class template. The helper retains Spark’s codegen factory-mode handling. Runtime validation covered Spark 4.1.3 only.
  • Key design decisions: The key includes column types, nullability, and method-splitting threshold. Only classes needing no captured references enter the cache. Compilation happens outside the lock, so simultaneous cold misses can still generate independently.
  • Implementation sketch: One call in doExecute switches to the new helper. The helper provides class caching and interpreted fallback. The new suite covers serialization parity, instance isolation, configuration separation, eviction, and query integration.
  • Performance: Source inspection and generation-count checks confirm that warm cache hits avoid repeated source generation. The author reports the depth-six, 201-partition benchmark improving from 7,798 ms to 3,992 ms. I did not independently rerun that benchmark. No evidence-backed performance regression was found.
  • Design: The change is narrowly scoped to projection construction. Existing batch iteration, Arrow ownership, broadcast conversion, and whole-stage execution remain intact. No further design issue met the reporting bar.
  • Abstraction & complexity: The helper isolates the cache and the small duplicated Spark class template. Spark still supplies the serialization implementation, avoiding a second set of row-writing semantics. No additional abstraction is needed for this scope.
  • Behavioral changes worth calling out: Compared with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1, the intended change is reuse of projection classes between partitions and queries. No unintended result or fallback change was identified.
  • Suggested improvements: Register org.apache.spark.sql.comet.CometUnsafeProjectionSuite in both Linux and macOS workflow matrices. Its omission currently fails preflight and blocks downstream validation.

Exact-head CI: Preflight failed at Check missing suites, and Required Checks failed. Linux/macOS builds and Spark SQL jobs were skipped. CodeQL, Analyze Actions, labeling, and title checks passed. The failed preflight job reports the same missing suite reproduced locally.

Validation and limits: Compiled the unchanged production helper against Spark 4.1.3, Scala 2.13.17, and JDK 21. A disposable standalone harness passed the five non-integration tests adapted from the new suite plus scalar-boundary/cache-hit and empty-output checks. python3 dev/ci/check-suites.py reproducibly exited 255. The full Comet/native build, Parquet integration test, cross-version runtime suites, and benchmarks were not run. Project files remain unchanged.

import org.apache.spark.sql.internal.SQLConf
import org.apache.spark.sql.types._

class CometUnsafeProjectionSuite extends CometTestBase with AdaptiveSparkPlanHelper {

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.

[P2] Register the new suite in both CI workflow matrices. Adding CometUnsafeProjectionSuite without listing its fully qualified name in .github/workflows/pr_build_linux.yml and .github/workflows/pr_build_macos.yml makes the mandatory suite-inventory check fail. Expected behavior is for preflight to pass and the new regression tests to run. Instead, exact-head preflight fails and downstream builds/tests are skipped. Please add the suite to the appropriate group in both workflows.

Evidence: At ed8ba0c, python3 dev/ci/check-suites.py exits 255 with Suite not found in workflow .github/workflows/pr_build_linux.yml: org.apache.spark.sql.comet.CometUnsafeProjectionSuite. The suite name is absent from both workflow files. Exact-head CI reproduces this at https://github.com/apache/datafusion-comet/actions/runs/38063389023/job/114246025282, and Required Checks fails.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Registered CometUnsafeProjectionSuite in the exec group of both workflows in 6d34718; check-suites.py and Preflight pass.

@andygrove andygrove added the run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue label Oct 10, 2026
@andygrove
andygrove requested a review from comphead October 10, 2026 17:54
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 run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants