Repository navigation
[EPIC] Nested-type conversion performance: C2R, R2C and JVM shuffle slower than Spark #6855
Description
Activity
- addedenhancementNew feature or requestNew feature or request
on Oct 10, 2026 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.
libcomettakes about 10% of the samples. The time goes to the reduce side, in each of the 201 tasks:CometColumnarToRowExec.doExecutecallsUnsafeProjection.createin every partition, which regenerates the projection's Java source. The schema is overspark.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.Follow-up to the comment above:
- test: extend C2R and R2C microbenchmarks to cover nested types #6856 now measures native shuffle in the ten nested shapes (36acd72). On the same 7950X3D, native shuffle runs them 1.1X to 3.8X faster than Spark at both 5 and 201 partitions, while JVM shuffle is at about parity (0.9X to 1.1X) at 5 partitions and 0.7X to 0.9X at 201. The fuzz-schema cases skip the native arm, since Spark's scan reads that schema.
- perf: generate the in-memory cache reader's class once per executor #6860, stacked on perf: generate the columnar-to-row projection once per executor instead of per partition #6858, does the same for Comet's in-memory cache reader (
CachedBatchRowIterator), which also generated its source in every partition when Spark reads a wide cached relation as rows. Reading the depth-6 fuzz schema from a 201-partition cache went from 16,971 ms to 1,390 ms, and a 1,500-column table at 200 partitions from 9,938 ms to 5,802 ms.
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
CometSparkToColumnarExec) loses to Spark on nested shapes. Plan is scan, native filter, columnar-to-row. Relative to Spark: nullablearray<string>0.7X,map<string,string>0.8X, struct holdingarray<string>0.8X, 20-field struct 0.8X, depth-3 struct 0.9X, nullable struct 0.9X, flat struct 1.0X. The same plan on the native scan is 1.3-2.4X. Related: [EPIC] Faster Spark-to-Arrow conversion in CometSparkToColumnarExec #6565, Reuse Arrow buffers across batches in Spark-to-Arrow conversion instead of reallocating and zeroing them #6720, Write strings and nested values in bulk when converting Spark rows to Arrow #6721, Admit all array and map types in CometSparkToColumnarExec #6722 (arrays other thanarray<string>and maps other thanmap<string,string>are not converted at all today).array<string>,array<struct>,array<array<int>>,map<string,int>,map<int,array<string>>, structs with strings, nullable structs and arrays). On the fuzz-generated nested schema (1000 rows, 201 partitions) it is far worse: 822 ms vs 94 ms for Spark at max depth 2 and 7,996 ms vs 380 ms at max depth 6. At 5 partitions it is at or near parity. That shape suggests a per-partition fixed cost (builder construction per partition per column) rather than per-row work. Related: perf: shuffle reader and writer performance review (native, JVM columnar, reader, Celeborn) #5905.rowIterator+UnsafeProjection(a proxy for Spark's path): flat schema 0.4X / 0.3X / 0.1X at batch sizes 8192 / 512 / 32; nested schema 0.9X / 0.7X / 0.2X, which is 1.55 us per row at batch size 32 (about 42 us per batch above the flat case). Small batches come from filters, joins and shuffle reads. Related: Native columnar-to-row conversion is much slower than the JVM implementation for small batches #5112, which has the flat numbers.CometColumnarToRowBenchmarkgroup (1.1-2.2X), so this is a tuning question rather than a regression against Spark.array<struct<int,string>>key is 0.9X andstruct<map<string,int>,int>keys are 0.9X, with the Spark-shuffle and JVM-shuffle arms at 0.5-0.9X.Nested conversion cost (Comet only, no Spark counterpart)
ArrowWriternested columns cost 10-30x a flat column, and the row path is 3-6x slower than the columnar path. ns per row, columnar / rows: int 0.8 / 6.3,array<string>17 / 60,map<string,string>25 / 100,array<struct<int,string>>35 / 96,array<array<int>>48 / 82,map<string,array<int>>60 / 137. Related: [EPIC] Faster Spark-to-Arrow conversion in CometSparkToColumnarExec #6565, Write strings and nested values in bulk when converting Spark rows to Arrow #6721.Benchmark follow-ups
row_columnarcaselist_conversion/elements_100takes 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.CometC2RIsolatedBenchwrites no results file (it is excluded frombenchmarks/micro/run.py), so its numbers are not tracked. Make it aBenchmarkBasesuite.Method
Suites run one at a time on the same machine:
CometColumnarToRowBenchmark,CometSparkToColumnarBenchmark,CometShuffleBenchmark --nested-only,CometArrowWriterBenchmark,CometC2RIsolatedBench, and the criterionrow_columnarbench. Spark'sBenchmarkreports 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
CometSparkToColumnarExec