Skip to content

Use native shuffle for row input by converting it with CometSparkToColumnarExec #6598

Description

@andygrove

What is the problem the feature request solves?

When a shuffle's child is a Spark row-based operator, Comet uses its JVM columnar shuffle. It serializes each UnsafeRow into off-heap pages, sorts the rows by partition, and converts them to Arrow natively, one partition at a time. Converting those rows with CometSparkToColumnarExec instead and using native shuffle is faster.

I measured both on the same query: an RDD scan of 4M rows, repartition(200, k), then an aggregate, at local[1] with AQE off. The two arms alternated over 7 runs on a 1.1.0-SNAPSHOT release build.

Columns JVM columnar shuffle CometSparkToColumnarExec + native shuffle
int, long, double 1602 ms 1354 ms (-15%)
int, long, string, double, decimal(18,2) 2160-2175 ms 1687-1759 ms (-19% to -22%)

The second arm set spark.comet.sparkToColumnar.enabled, whose default operator list includes RDDScan. That only converts leaf operators. The build predates #6566, whose faster decimal path should widen the second gap.

CometSparkToColumnarExec converts rows more slowly than the native converter does (#6596 has the per-type numbers). So the columnar shuffle's extra 60-120 ns per row is spent elsewhere in its write path. I haven't profiled where.

Describe the potential solution

When a shuffle would otherwise use the JVM columnar shuffle, insert CometSparkToColumnarExec over its child and use native shuffle. Put it behind a config that is off by default, and keep the JVM columnar shuffle as the fallback for whatever native shuffle or CometSparkToColumnarExec can't take:

Additional context

This reuses the ArrowWriter that CometSparkToColumnarExec shares with the cache serializer and CometLocalTableScanExec, so #6566 applies to these shuffles too.

One difference to weigh before turning it on by default: CometArrowAllocator allocates the converted batches outside any memory pool, while the JVM columnar shuffle accounts for its pages.

#6593 proposes a spark.comet.convert.* config for each Spark-to-Arrow conversion. This would be another one.

Metadata

Metadata

Assignees

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions