Skip to content

[EPIC] Nested-type conversion performance: C2R, R2C and JVM shuffle slower than Spark #6855

Description

@andygrove

What is the problem the feature request solves?

The microbenchmarks for columnar-to-row (C2R) and row-to-columnar (R2C) conversion covered nested types thinly: no nulls in nested data, no collections of structs or arrays, no string-element collections in the native R2C bench, and no nested end-to-end case for CometSparkToColumnarExec. #6856 closes those gaps (see the list there). Running the extended suites on a 32-thread Linux box (Ryzen 9 7950X3D, JDK 17, local[1], 1M rows, release build) turned up the cases below where Comet is slower than Spark, or where the nested cost is out of proportion. This epic tracks them.

Where Comet is slower than Spark

Nested conversion cost (Comet only, no Spark counterpart)

Benchmark follow-ups

  • The pre-existing row_columnar case list_conversion/elements_100 takes about 1.6 ms at both 1,000 and 10,000 rows, and bytes written are not monotonic in the row count (141 KB at 1k rows, 31 KB at 10k; the 10%-null variant writes 6 KB at 1k and 1.4 MB at 10k). Either the benchmark data or a data-dependent path is wrong. Do not read the 100-element list numbers until this is understood.
  • CometC2RIsolatedBench writes no results file (it is excluded from benchmarks/micro/run.py), so its numbers are not tracked. Make it a BenchmarkBase suite.
  • No Spark-comparison benchmark for R2C from a row source (RDD, local relation) with nested types.
  • Add a small-batch variant of the nested C2R and Spark-to-Arrow cases to the SQL benchmarks, since the SQL path always uses full batches.

Method

Suites run one at a time on the same machine: CometColumnarToRowBenchmark, CometSparkToColumnarBenchmark, CometShuffleBenchmark --nested-only, CometArrowWriterBenchmark, CometC2RIsolatedBench, and the criterion row_columnar bench. Spark's Benchmark reports best time over at least 2 s of iterations after a 2 s warmup; criterion runs were not CPU-pinned except one rerun. Numbers are from a single run per case.

Related

Activity

  1. self-assigned this
    on Oct 10, 2026
  2. andygrove commented on Oct 10, 2026

    @andygrove
    MemberAuthor

    On the "JVM shuffle with 201 partitions" item: a profile of the fuzz schema case (max depth 6) shows the per-partition cost is not in the native writer. libcomet takes about 10% of the samples. The time goes to the reduce side, in each of the 201 tasks:

    • CometColumnarToRowExec.doExecute calls UnsafeProjection.create in every partition, which regenerates the projection's Java source. The schema is over spark.sql.codegen.maxFields, so the plan is never under whole-stage codegen. That was about 19 ms per partition at depth 6 and 2 ms at depth 2. perf: generate the columnar-to-row projection once per executor instead of per partition #6858 generates it once per executor, which halves those cases (827 to 404 ms, and 7,798 to 3,992 ms).
    • Most of what remains is importing the wide schema's Arrow arrays for each batch (ArrowImporter.importVector, about 15% of the executor thread), and wrapping and closing the CometVectors (about 15%).

    Also, in these cases and in the ten nested shapes #6856 adds, 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. Its parity with Spark says nothing about native shuffle. For the ten shapes, enabling round-robin in that arm, or shuffling by hash on a key, would measure native shuffle.

  3. andygrove commented on Oct 10, 2026

    @andygrove
    MemberAuthor

    Follow-up to the comment above:

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

Metadata

Metadata

Assignees

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions