Repository navigation
Conversation
Add nested coverage to the columnar-to-row, row-to-columnar and shuffle microbenchmarks: nullable nested columns, collections of structs and arrays, large collections, an isolated native C2R case with a nested schema, a Spark-to-Arrow benchmark for nested shapes, and string and nested cases for the native shuffle row-to-columnar bench.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Nested conversion benchmarks underrepresented nullable values, variable-width collections, deeper nesting, and large collections. This PR expands those fixtures across six benchmark files.
- Design approach: Extend existing harnesses and add a three-arm Spark-to-Arrow benchmark comparing Spark, native scanning, and conversion from Spark scanning.
- Correctness: Found one P2 measurement issue: the newly added shuffle cases label Spark fallback as native shuffle. Separately, 33 sampled Rust-encoded rows matched Spark
UnsafeProjectionbyte-for-byte. Nested vector and SQL fixture checks passed. - Compatibility analysis: Compared relevant row-layout methods and
Repartition.partitioningagainst Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0 sources. No additional compatibility issue met the reporting bar. Compared touched benchmark files withbranch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1. This PR changes benchmark behavior without changing production execution. - Key design decisions: Fixture construction and scan-plan verification happen outside timed cases. Isolated C2R reuses independently owned batches at three batch sizes. Larger collection fixtures use fewer rows.
- Implementation sketch: A recursive Rust
Valueencoder constructs nested UnsafeRows. Scala adds nested schemas and SQL fixtures, reserves flattened struct-child capacity, and introducesCometSparkToColumnarBenchmarkand--nested-only. - Performance: No production overhead is added. Setup remains outside measurement, but the new shuffle matrix cannot currently support native-versus-JVM conclusions because its native arm falls back. Full timing runs were not repeated.
- Design: Reusing existing harnesses keeps the extension understandable. The new scan benchmark's explicit plan checks are useful. The added shuffle cases need equivalent protection against measuring fallback.
- Abstraction & complexity: The local encoder centralizes null bits, offsets, alignment, and nested payloads. The abstractions remain confined to benchmarks. No separate P1/P2 complexity concern was identified.
- Behavioral changes worth calling out: Default suites run more cases,
--nested-onlylimits shuffle groups, and the isolated flat C2R consumer now reads all four fields, affecting comparison with older timings. - Suggested improvements: Address the finding by explicitly selecting supported native partitioning for the added shuffle cases and verifying a
CometNativeShuffleexchange outside timing.
Reviewed the full six-file PR diff at 25b3384393d0c01950c20b72034094d31e0491a1 against supplied base dc880d6d265f8391b6202efa43574cfeea71e75a, using merge base a32529c4d84f5157be7a23aa0742ecdc15090eb7. GitHub's changed-file list agrees. The PR remained non-draft. The snapshot and refreshed discussion endpoints contained no existing reviews or comments.
Routed skills: review-comet-pr, review-comet-shuffle-pr, and review-comet-ffi-pr. Also consulted memory guidance for fixture capacity and ownership.
Exact-head CI at the final check: 9 successful, 10 in progress, and 14 skipped, with no failures reported. Successful checks included lint, Scala syntactic lint, and CodeQL. Native build, Rust tests, and several compatibility checks remained running. Spark SQL, Iceberg, macOS, and benchmark checks were skipped.
Validation: check-benchmark-runner.py passed for 58 suites. cargo check --offline --locked -p datafusion-comet-shuffle --bench row_columnar passed. All five changed Scala benchmark sources compiled against cached dependencies. Spark 4.1.3 checks passed for the Rust encoder samples, 19 column-vector fixture types, 256 isolated nested rows, eight Spark-to-Arrow SQL shapes, and all three new C2R SQL groups. A probe compiling the exact-head shuffle selector and configuration reproduced the finding and verified that enabling native round-robin resolves admission.
Limits: No full reactor build, complete benchmark execution, or cross-version runtime matrix was run. Scala compilation and planner probes used cached supporting dependencies. Project source remained unchanged. One introduced P2 issue is reported, with no additional introduced P1/P2 issues found.
| .parquet(filename) | ||
| } | ||
| Seq(5, manyPartitions).foreach { partitionNum => | ||
| shuffleDeeplyNestedBenchmark( |
There was a problem hiding this comment.
[P2] Make the added native-shuffle cases actually select native shuffle. Each new shape calls shuffleDeeplyNestedBenchmark, whose Comet (native Shuffle) arm uses .repartition(partitionNum) and sets spark.comet.shuffle.mode=native. At both 5 and 201 partitions this requests RoundRobinPartitioning, but spark.comet.shuffle.native.partitioning.roundrobin.enabled remains false. Consequently, the arm falls back to Spark shuffle instead of measuring native shuffle, duplicating the Spark-shuffle comparison and leaving all 20 added native measurements mislabeled. Please explicitly enable native round-robin, or repartition these fixtures by their scalar key, and verify a CometNativeShuffle exchange outside timing.
Evidence: Compiled the exact-head CometConf.scala and CometShuffleExchangeExec.scala with a disposable Spark 4.1.3 planner probe. For five representative added nested types at both partition counts, mode native returned None from shuffleSupported, with reasons spark.comet.shuffle.native.partitioning.roundrobin.enabled is disabled and Comet columnar shuffle not enabled. Mode jvm returned Some(CometColumnarShuffle). Enabling the round-robin flag returned Some(CometNativeShuffle). Spark sources for every supported profile map unkeyed repartitioning with more than one partition to RoundRobinPartitioning. The benchmark and standard launcher do not enable that flag.
There was a problem hiding this comment.
Fixed in 36acd72: the native arm turns on round-robin, and each Comet arm's planned shuffle is checked before timing, skipping an arm that falls back. All ten shapes now run native shuffle; the fuzz cases skip it because Spark's scan reads that schema.
The native arm of shuffleDeeplyNestedBenchmark planned Spark's shuffle. `repartition(n)` plans a round-robin, which native shuffle declines unless spark.comet.shuffle.native.partitioning.roundrobin.enabled is on. Turn it on for that arm, check before timing which shuffle each Comet arm plans, and skip an arm that does not plan the shuffle it is named after. The fuzz-generated schema is read by Spark's scan, which native shuffle cannot take, so its native arm is skipped with a note.
Which issue does this PR close?
Part of #6855.
Rationale for this change
The C2R and R2C microbenchmarks covered nested types thinly: nested test data never contained nulls, collections held 2-5 simple elements, the isolated C2R bench and the native
row_columnarbench had no strings or collections of structs, and nothing measuredCometSparkToColumnarExecon nested input. Extending them found several cases where Comet is slower than Spark (filed under #6855).What changes are included in this PR?
Benchmarks only, no production code.
CometColumnarToRowBenchmark: three new groups. Nullable nested types (null structs, arrays and maps, empty arrays, null elements), nested collections (array of struct, array of array, struct of array of struct, map of struct, decimals, timestamps and binary inside nested types), and large collections (100-element arrays, 20-entry maps).CometC2RIsolatedBench: a nested schema (struct, array, map, array of struct, with nulls) alongside the flat one, at batch sizes 8192, 512 and 32.native/shuffle/benches/row_columnar.rs: a small UnsafeRow encoder and 12 new cases coveringlist<utf8>,list<struct>,list<list>, lists with null elements, structs with strings, nested structs with strings, a struct holding a list, andmap<utf8,*>variants.CometArrowWriterBenchmark: nine more nested types (nested structs,array<struct>,array<array>,struct<array>, maps with array values, and others).CometSparkToColumnarBenchmark(new): Spark Parquet scan converted to Arrow versus Spark and the native scan, for eight nested shapes, with plan checks on each arm.CometShuffleBenchmark: ten nested shapes across Spark shuffle, Comet Spark shuffle, and JVM and native shuffle, plus a--nested-onlyflag.How are these changes tested?
The new benchmarks were built in release mode and run on a Linux machine; results and findings are in #6855. The native bench encoder was checked to produce byte-identical rows to the existing hand-built
listandmaprows.cargo fmt,cargo clippy --benches -D warnings,spotlessanddev/ci/check-benchmark-runner.pypass.