Skip to content

Revisit ExternalSorter sort strategy and sort_in_place_threshold_bytes #21543

Description

@mbutrovich

Context

ExternalSorter branches on sort_in_place_threshold_bytes (default 1MB) in in_mem_sort_stream():

  • Below 1MB: concatenate all buffered batches into one RecordBatch, sort in place
  • Above 1MB: sort each batch individually, then streaming-merge them

This threshold was introduced in May 2023 by @tustvold in #6163 ("Adaptive in-memory sort") with the comment: "This is a very rough heuristic and likely could be refined further." It was later extracted to a config option by @alamb in #7130 with the same 1MB default. The default hasn't changed since, though the surrounding sort architecture has evolved significantly: multi-level merge (#15700), chunked sort output (#19494), IncrementalSortIterator (#20314), and PartialSortExec (#9125).

Problem

The sort-each-batch-then-merge path dominates real workloads because typical in-memory buffer sizes exceed 1MB. In this path, each batch (often 1024–8192 rows) is sorted individually via lexsort_to_indices and then merged via StreamingMergeBuilder. This means:

  1. Per-batch sort kernels can't amortize overhead. Row-format sorting (e.g., MSD radix sort on RowConverter output, feat(arrow-row): add MSD radix sort kernel for row-encoded keys arrow-rs#9683) is 2–3x faster than lexsort_to_indices at 32K+ rows, but at 1K–8K rows the RowConverter encoding cost dominates. The sort-then-merge path never gives these kernels enough rows to benefit. perf: Bring over apache/arrow-rs/9683 radix sort, integrate into ExternalSorter #21525 attempted to integrate the radix sort kernel into ExternalSorter and saw no improvement for this reason.

  2. The concat path is gated on memory, not row count. The 1MB threshold is a memory proxy, but the actual concern is the temporary 2x memory spike from concat_batches (plus RowConverter allocation on top). A row-count or batch-count heuristic might be a better fit.

  3. The sort benchmark doesn't exercise the merge path. The benchmark produces 8 partitions of ~12 batches at 1024 rows each. In the sort partitioned variant, each partition's ~12K rows (~100KB for integers) falls well below the 1MB threshold, so it always takes the concat-and-sort-in-place path. This means benchmark results don't reflect the sort-then-merge path that dominates at larger data sizes.

Prior art: DuckDB

For comparison, DuckDB's sort redesign encodes into its normalized key format as data arrives during the sink phase, accumulating into large thread-local sorted runs. The encoding cost is amortized across the entire input stream, and sorting happens once per run on already-encoded data. This avoids the small-batch problem entirely — by the time sorting begins, each thread has a single large run to sort.

Possible directions

These aren't mutually exclusive:

  • Raise or rethink the threshold. The 1MB limit was chosen conservatively. With IncrementalSortIterator (added in fix: Unaccounted spill sort in row_hash #20314) now yielding sorted output in chunks, the peak memory of the concat-and-sort path may be more manageable than it was in 2023. Could we raise it, or gate on row count instead?

  • Coalesce batches before sorting in the merge path. When the merge path is taken, we could concatenate small batches into larger ones (e.g., 32K–64K rows) before sorting, giving row-format kernels enough rows to amortize encoding. This also reduces the merge fan-in, which is related to Improve performance of large sorts with Cascaded merge / tree #7181 (cascaded merge for large fan-in). This doesn't require concatenating the entire buffer — just local coalescing.

  • Incremental Rows encoding. RowConverter::append already supports incrementally extending a Rows buffer across batches. ExternalSorter could maintain a Rows alongside its in_mem_batches, calling append as each batch arrives (similar to DuckDB's approach). At sort time, the encoding is already done — you just sort the accumulated Rows and use the indices to reorder the original batches. The tradeoff is higher memory during accumulation (raw batches + encoded rows), but encoding cost is fully amortized and radix sort gets a large contiguous run to work with.

  • Fix the benchmark. Increase BATCH_SIZE and/or INPUT_SIZE in benches/sort.rs so that sort partitioned exercises the sort-then-merge path.

Related issues

References

Activity

  1. mbutrovich commented on Apr 10, 2026

    @mbutrovich
    ContributorAuthor

    Leaving some notes to myself for next week, and anyone else to give feedback on. One idea I came up with:

    Chunked sort with incremental accumulation

    1. Buffer RecordBatches as they arrive (same as today).
    2. When ~32K rows accumulate, concat the chunk, sort it, materialize via take(), stash as a sorted run.
    3. On memory pressure, spill sorted runs directly — they're already sorted, no re-sorting needed.
    4. On input exhaustion, k-way merge all runs using the existing loser tree / multi-level merge.

    For radix-eligible types: encode the chunk to Rows, radix sort, drop the Rows. Encoding is transient. For dictionaries/nested types: lexsort_to_indices on the concatenated chunk, same as the current sub-1MB path.

    This design is similar to the new immediate-mode shuffle writer in DataFusion Comet: apache/datafusion-comet#3845

    Why this helps

    • Radix sort: 2–3x faster than lexsort_to_indices at 32K rows (feat(arrow-row): add MSD radix sort kernel for row-encoded keys arrow-rs#9683), but needs large batches to amortize RowConverter encoding. This gives it that.
    • Any sort kernel: Fewer sorted runs means lower merge fan-in. 100K rows goes from ~100 runs of 1K to 3 runs of 32K.
    • Memory: Peak during a chunk sort is raw_chunk + Rows + sorted_output for one 32K-row chunk. Bounded and predictable. The 1MB threshold question mostly goes away — sort always operates on fixed-size chunks, memory pressure is handled by spilling sorted runs.

    For context, DuckDB's sort redesign takes this further — encoding into normalized keys as data arrives, accumulating large thread-local runs. This proposal is a step in that direction while staying within DataFusion's existing architecture.

    Degenerate cases

    I don't think there are new ones beyond what already exists. Fan-in is strictly better (fewer, larger runs). Skew doesn't affect radix sort (O(n × key_width)). Small datasets (< 32K rows) reduce to a single chunk — same as the current sub-1MB path. The existing multi-level merge handles large fan-in when many chunks are needed.

  2. Dandandan commented on Apr 11, 2026

    @Dandandan
    Contributor

    Using https://docs.rs/arrow/latest/arrow/compute/struct.BatchCoalescer.html
    might also help a bit (instead of using concat_batches).

    (It still uses take but avoids the double memory usage caused by concat_batches.
    Additionally it allows doing the copying earlier, which also reduces the "final" allocation/CPU spike and might also be a bit more efficient as batch is likely be in CPU cache when processed (and not when doing concat_batches on a large number of batches), spreading out deallocations, etc.

    (There is potential to make it even faster by fusing take + push_batch (in push_batch_with_indices, but that is yet to be implemented).

  3. Dandandan commented on Apr 11, 2026

    @Dandandan
    Contributor

    Another option would perhaps to increase the target batch_size of a downstream operator (which already use the coalesce kernel) whenever SortExec is present, this would avoid the double copying in both the operator and SortExec

  4. alamb commented on Apr 13, 2026

    @alamb
    Contributor

    Incremental Rows encoding. RowConverter::append already supports incrementally extending a Rows buffer across batches. ExternalSorter could maintain a Rows alongside its in_mem_batches, calling append as each batch arrives (similar to DuckDB's approach). At sort time, the encoding is already done — you just sort the accumulated Rows and use the indices to reorder the original batches. The tradeoff is higher memory during accumulation (raw batches + encoded rows), but encoding cost is fully amortized and radix sort gets a large contiguous run to work with.

    I think this is a good idea to pursue as well -- I also wonder if we have already created data in the row format, we could avoid the second copy entirely perhaps by keeping a list of sorted indices and then merging using those rather than copying the data again 🤔

  5. mbutrovich commented on Apr 13, 2026

    @mbutrovich
    ContributorAuthor

    I think this is a good idea to pursue as well -- I also wonder if we have already created data in the row format, we could avoid the second copy entirely perhaps by keeping a list of sorted indices and then merging using those rather than copying the data again 🤔

    Yeah I was looking at this right now, in fact: could a radix sort in row format benefit the merge phase? It's more plumbing and might be one of the later PRs in a sequence of changes, but on paper it seems like a good idea.

  6. mbutrovich commented on Apr 16, 2026

    @mbutrovich
    ContributorAuthor

    So the TPC-H and TPC-DS results in #21629 are very exciting but I am concerned about regressions with string types (not string view) and generally wide schemas. I wonder:

    • should we only coalesce sort expression columns, convert those to rows, then map those back to their original batches for the take? I think joins do something similar when join keys span input batches
    • should we do something data-type specific? It's hard to find a general solution here :(

    Sorts are fun.

  7. Dandandan commented on Apr 16, 2026

    @Dandandan
    Contributor

    We should be able to run the benchmarks with string types (with the right env variable)

  8. mbutrovich commented on Apr 17, 2026

    @mbutrovich
    ContributorAuthor

    I worked on several implementations tonight, summarizing findings (thanks Claude):

    ExternalSorter optimization — full summary of learnings

    The problem

    ExternalSorter sorts each incoming batch individually, producing one sorted run per batch. At scale (TPC-H SF10, 60M rows), this creates ~7300 sorted runs with high merge fan-in. Coalescing batches before sorting reduces fan-in and gave 1.2-1.5x TPC-H speedups — but introduced regressions for multi-column string sorts.

    Iterations

    externalsorter (PR #21629): BatchCoalescer coalesces all columns to 32K → lexsort → take.

    • Fixed-width: 1.3-1.6x faster. StringView: 1.1-1.2x faster. TPC-H: 1.08-1.51x faster.
    • Multi-column StringArray: 2.5x slower. Dictionary: 2.6x slower.
    • Root cause: take scatter-gathers across ~1.9MB of string heap at 32K rows, exceeding L2 cache.

    externalsorter2: Key-only extraction + interleave. Extract sort-key columns (promote to StringView), concat keys only, sort, translate indices to (batch_idx, row_idx) pairs, interleave_record_batch on original batches.

    • Fixed single-column regressions vs es1.
    • Multi-column still slow: interleave_record_batch is expensive for StringArray (same scatter-gather) and DictionaryArray (dictionary merging + key remapping per chunk).
    • Learning: existing benchmarks sort on ALL columns, so key-only extraction provided zero benefit for them. Need wide-schema benchmarks to measure the real benefit.

    externalsorter3: Internal StringView representation. Convert ALL columns StringArray→StringView on input, concat, sort, take on views (16-byte copies), convert back at output boundary.

    • Uniformly better than es2 — removed interleave overhead, direct take instead.
    • Big win: sort utf8 low cardinality 1.78x faster than main (StringView prefix comparisons for short inlined strings).
    • Multi-column tuples still 1.5-2.6x slower than main.

    externalsorter4: Key-only extraction + RowConverter sort + hybrid reconstruction. Extract key columns, concat, sort via RowConverter (for multi-column varlen) or lexsort (for single-col/fixed-width), take keys from concat batch, interleave values from originals.

    • RowConverter cut sort utf8 tuple from 180→144ms (22% vs es3). sort utf8 view tuple from 156→123ms (21%).
    • Fixed-width unchanged at 1.31x faster than main.
    • Multi-column strings now 1.22-1.30x slower than main (down from 2.5x at start).

    Key technical findings

    1. take on StringArray at large batch sizes is cache-hostile

    take with a random permutation scatter-gathers across the string offset/data buffers. At 32K rows × 3 StringArray columns ≈ 1.9MB — exceeds L2 cache. At 1024 rows × 3 columns ≈ 30KB — fits in L1. This is the fundamental source of the multi-column string regression.

    2. Multi-column sort regression scales superlinearly with column count

    Benchmark columns main (ms) es3 (ms)
    sort utf8 view high card 1 62.3 58.8 (faster)
    sort utf8 view tuple 3 100.6 156.1 (slower)

    1 column wins, 3 columns loses badly. The cause: lexsort_to_indices with multiple columns does cascading comparisons. With low-cardinality first columns, most comparisons cascade to column 2 and 3. At 32K rows the sort algorithm's random access pattern crosses cache boundaries on every comparison across 3 separate column arrays.

    3. RowConverter encoding solves the multi-column cache problem

    RowConverter encodes all sort key columns into one contiguous buffer of binary-comparable rows. The sort operates on this single buffer instead of jumping between separate column arrays. Cache prefetcher works because access stays within one memory region. arrow-rs#9683 benchmarks confirm: lexsort_rows is 1.3-2.5x faster than lexsort_to_indices for multi-column string/dictionary schemas at 32K rows.

    4. But RowConverter is slower for single-column and fixed-width sorts

    Encoding overhead exceeds the benefit when lexsort_to_indices can use SIMD comparisons on native column data. Schema-based kernel selection is necessary: RowConverter for multi-column varlen, lexsort for everything else.

    5. GC on StringView concat batches doesn't help

    We hypothesized that compacting 32 backing buffers to 1 would improve sort locality. Benchmarks showed no improvement — the cache problem is the total working set size (960KB for 32K × 3 columns), not the number of buffers.

    6. StringView promotion helps for short strings, not long strings

    StringView inlines strings ≤12 bytes in the 16-byte view struct — no pointer chase needed. For strings >12 bytes, comparison still chases a pointer, making StringView no better than StringArray. The benchmark's low-cardinality values ("value0"–"value99", 6-8 chars) benefit from inlining. High-cardinality values (20-char random strings) don't.

    7. Non-key columns should never be converted or concatenated

    Converting/concatenating value columns that aren't sort keys is pure overhead. The key-only approach (es2/es4) avoids this. For wide schemas (small key + large value), the savings are significant.

    8. interleave_record_batch is expensive for DictionaryArray

    interleave_dictionaries merges dictionary values across all input batches and recomputes key mappings per chunk. For 32 input batches × 3 dictionary columns × 4 chunks per run, this is a lot of dictionary merging work that take on a single batch doesn't need.

    9. Dictionary encoding is increasingly irrelevant

    Comet unpacks dictionaries to StringArray anyway. DataFusion native uses StringView from Parquet. Dictionary sort performance doesn't represent real workloads.

    10. The benchmark's BATCH_SIZE=1024 is pessimistic

    With 1024-row input batches, per-batch sort on main is very cheap (30KB working set fits in L1). Coalescing accumulates 32 batches per run, adding significant per-run overhead (32 key extractions, 32-batch concat). At DataFusion's default 8192-row batches, only 4 batches per run — much less overhead, and main's per-batch sort has a larger working set (less of a cache advantage).

    11. The merge uses RowConverter anyway

    The streaming merge (StreamingMergeBuilder) already encodes sort keys via RowConverter for merge comparisons. So using RowConverter during the sort step doesn't add new encoding work to the total pipeline — it moves the encoding earlier (sort time vs merge time).

    12. Arrow-rs improvements that would help

    • concat_batches for StringView producing 1 backing buffer directly (avoid N-buffer scatter)
    • Efficient Dictionary → StringView conversion (convert dictionary values only, not every row)
    • lexsort_rows / lexsort_radix from arrow-rs#9683 (purpose-built sort on row-encoded data)

    Final benchmark results (sort 1M, all iterations vs main)

    Benchmark main es1 es2 es3 es4
    sort i64 45.5 33.8 37.4 34.6 34.8
    sort f64 46.4 35.8 39.7 36.7 —
    sort utf8 low card 67.4 64.6 57.7 37.9 —
    sort utf8 high card 74.0 59.7 73.6 66.6 —
    sort utf8 view low card 34.9 25.6 29.0 26.3 —
    sort utf8 view high card 62.3 65.1 66.6 58.8 —
    sort utf8 dict 49.3 31.8 35.1 31.9 —
    sort mixed tuple w/ view 93.9 85.9 90.1 85.2 —
    sort mixed tuple 102.3 105.7 112.5 105.8 —
    sort utf8 dict tuple 78.5 92.4 96.5 92.0 91.3
    sort utf8 tuple 111.0 272.1 187.6 180.2 144.5
    sort utf8 view tuple 100.6 163.7 163.9 156.1 123.2
    sort mixed dict tuple 108.2 282.6 288.1 282.9 282.6

    Approach summary

    Iteration Key idea What improved What didn't
    es1 Coalesce all columns, lexsort, take Fixed-width, single string, TPC-H Multi-col string/dict (cache thrashing on take)
    es2 Key-only extract, interleave values Single-col sorts Multi-col (interleave expensive for string/dict)
    es3 Internal StringView, full-batch concat+take Low-card strings, uniformly better than es2 Multi-col tuple (working set exceeds L2)
    es4 Key-only + RowConverter sort + hybrid reconstruct Multi-col string tuple (22% improvement) Dict (not using RowConverter), residual overhead
  9. gratus00 commented on Apr 21, 2026

    @gratus00

    Hi, related to the "Fix the benchmark" direction, I noticed a small benchmark gap in datafusion/core/benches/sort.rs.

    The existing tuple/string/dictionary cases use make_sort_exprs(schema), which sorts by every column. Would a small benchmark-only PR adding a case that includes a table with a cheap sort-able key such as an i64 key and non-key utf8/dictionary payload columns but ONLY gets sorted by the i64 key be useful?

    The goal would be to make the sort key cheap and better measure the cost of take reordering wider payload column cases. I think this case is not covered at the moment.

    Would love this to be a way I start working on this project!

    @Dandandan @mbutrovich

  10. mbutrovich commented on Apr 22, 2026

    @mbutrovich
    ContributorAuthor

    The existing tuple/string/dictionary cases use make_sort_exprs(schema), which sorts by every column. Would a small benchmark-only PR adding a case that includes a table with a cheap sort-able key such as an i64 key and non-key utf8/dictionary payload columns but ONLY gets sorted by the i64 key be useful?

    Hi @gratus00! We could definitely use more sort benchmarks, and in fact some of what you're describing I have in #21688 if you want to look at those for inspiration and bring them to a separate PR.

    I had Claude summarize my notes from over the weekend of looking at this. I'm kinda stumped. I can make TPC-H faster, but at the expense of other schemas I care about (e.g., unconditional coalescing hurts wide/large schemas).

    ExternalSorter Investigation Summary

    Context

    PR #21688 (externalsorter4) rewrites ExternalSorter to coalesce batches before sorting, reducing merge fan-in. At TPC-H SF10 (~60M rows), this gives 11 queries faster (up to 1.51x), 0 regressions. Issue #21543 tracks the broader redesign.

    Four iterations were explored (es1–es4), each trying to address regressions introduced by the previous. PR #21629 (externalsorter, es1) is the predecessor.

    Key Findings

    1. BATCH_SIZE matters enormously

    All prior microbenchmark comparisons used BATCH_SIZE=1024. At the realistic default of 8192, the picture flips:

    Benchmark (1M rows) main es5 (=es4+8192) ratio
    i64 34.9 33.4 0.96x (same)
    f64 36.1 35.6 0.98x (same)
    utf8 tuple (3 col) 95.0 147.7 1.55x slower
    utf8 view tuple (3 col) 81.4 125.1 1.54x slower
    mixed tuple 83.9 127.2 1.52x slower
    mixed tuple w/ view 75.9 111.3 1.47x slower
    utf8 dict tuple 49.8 93.8 1.88x slower
    mixed dict tuple 76.6 275.1 3.59x slower
    utf8 dict 37.9 32.5 0.86x (faster)
    utf8 view low card 25.8 26.3 ~same
    utf8 high card 65.7 61.3 0.93x (same)

    At 1024, es4 showed wins on single-column sorts (1.3-1.6x faster). At 8192, those wins disappear because main's per-batch sort already has a meaningful working set — the cache advantage of coalescing shrinks when the expansion is only 4x (8192→32768) instead of 32x (1024→32768).

    Multi-column tuple regressions got worse at 8192, not better.

    2. The microbenchmark is in the wrong regime

    At 1M rows / 8192 batch size: ~122 batches → ~122 runs on main, ~30 on es4. Fan-in reduction is modest (4x) and the merge at fan-in 122 is already manageable.

    At TPC-H SF10 (~60M rows): ~7300 batches → ~7300 runs on main, ~1800 on es4. The fan-in reduction matters here — multi-level merge kicks in, cursor management for 7300 streams is expensive, etc.

    The crossover point where coalescing pays off is somewhere between 1M and 60M rows. The microbenchmark is below it; TPC-H is above it.

    3. Profiling reveals double RowConverter encoding

    Profiled sort utf8 tuple 1M on both branches (samply + Firefox Profiler, --profile profiling):

    main (inverted call tree, self time):

    Self % Function
    3.6% _platform_memmove
    3.5% arrow_row::variable::encode_one
    ~1.8% quicksort partition
    ~1.4% LengthTracker::push_variable

    es5 (inverted call tree, self time):

    Self % Function
    7.8% _platform_memmove
    7.6% arrow_row::variable::encode_one
    7.6% arrow_row::RowConverter::append
    ~3.7% arrow_select::take::take_impl
    ~3.5% sort internals

    es5 spends 2x the time in memmove and 2x in RowConverter encoding vs main. This is because es5 encodes to Rows twice:

    1. Sort phase: concat key columns (memmove) → encode to Rows → sort Rows → take
    2. Merge phase: RowCursorStream re-encodes sorted output batches to Rows for merge comparisons

    The Rows from the sort phase are discarded. The merge builds new ones from scratch. ~15% of total time is wasted on this double work.

    4. StreamingMerge only uses RowConverter for multi-column sorts

    StreamingMergeBuilder::build() dispatches:

    • Single-column primitive → FieldCursorStream<PrimitiveArray> — native comparisons, no RowConverter
    • Single-column Utf8/Utf8View/Binary → FieldCursorStream — byte comparisons, no RowConverter
    • Multi-column → RowCursorStream — RowConverter encoding per batch

    The double-encoding problem only affects multi-column sorts — which are exactly the benchmarks with the worst regressions.

    5. Key-only extraction doesn't help at realistic batch sizes

    New key/value benchmarks (sort on key columns only, carry non-key value columns):

    Benchmark (1M) main es5 ratio
    i64 key, 10x utf8 view value 46.1 53.7 1.16x slower
    utf8 view low card key, large value 28.0 32.1 1.15x slower
    utf8 view high card key, large value 51.0 66.5 1.30x slower
    (i64, utf8 view) key, 5x f64 value 66.5 98.6 1.48x slower

    Key-only extraction (es4's core optimization — don't touch value columns during sort) provides no benefit at 8192 batch size. The coalescing overhead exceeds the savings from avoiding value column work.

    6. Dictionary regressions don't matter for real workloads

    • DataFusion native Parquet readers produce StringView, not DictionaryArray
    • Comet unpacks dictionaries to StringArray before ExternalSorter
    • No major consumer sends DictionaryArray into the sort pipeline
    • The 3.59x mixed dict tuple regression is a benchmark artifact

    7. es3's StringView conversion wasn't a universal win

    es3 converted StringArray→StringView internally. Results were mixed:

    • Low-cardinality short strings: 1.70x faster (inlined in 16-byte views)
    • High-cardinality long strings (>12 bytes): 1.12x slower (views still chase pointers)
    • Dictionary columns: untouched (es3 never tried Dict→StringView)

    Proposed Architecture

    Two orthogonal improvements that reduce the coalescing overhead:

    A. Incremental RowConverter::append (eliminate concat)

    Current es4 (for varlen multi-column keys):

    batch₁ keys ─┐
    batch₂ keys ─┼── concat (memmove #1) ──► big key array ──► RowConverter::append (memmove #2) ──► Rows
    batch₃ keys ─┘
    

    Proposed:

    batch₁ keys ──► RowConverter::append ──┐
    batch₂ keys ──► RowConverter::append ──┼──► Rows  (one memmove)
    batch₃ keys ──► RowConverter::append ──┘
    

    RowConverter::append already supports incremental extension. This eliminates the key column concat entirely.

    B. Pass Rows from sort phase to merge phase (eliminate double encoding)

    Current: sort encodes to Rows → sort → discard Rows → merge re-encodes to Rows

    Proposed: sort encodes to Rows → sort → pass Rows to merge → merge uses them directly

    This requires plumbing changes in StreamingMergeBuilder to accept pre-encoded Rows via a new PartitionedStream implementation. Only applies when the merge uses RowCursorStream (multi-column sorts).

    Combined effect

    Eliminates 3 copies/encodings → 1:

    1. concat key columns (removed by A)
    2. encode for sort (kept — this is the one encoding)
    3. re-encode for merge (removed by B)

    From the profile, this could save ~15% of total time (7.8% memmove + 7.6% encoding), which would shift the crossover point lower and potentially make coalescing viable even at 1M rows.

    Scope and limitations

    • Only applies to multi-column sort keys (single-column sorts don't use RowConverter in merge)
    • Only applies to the RowConverter path (varlen keys) — lexsort path for fixed-width keys still needs concat
    • Reconstruction of output columns (take/interleave) is unchanged
    • Requires changes to PartitionedStream trait and StreamingMergeBuilder API

    Encode-Once Implementation and Profiling (branch: sort_rows)

    What we built

    Implemented the encode-once architecture on branch sort_rows:

    • uses_row_cursor(): shared dispatch matching StreamingMergeBuilder::build() — single-column primitive/string sorts stay on main's path, multi-column sorts use the new row-encoded path
    • AccumulationBuffer: incremental RowConverter::append as batches arrive
    • SortedRun: deferred materialization — stores original batches + sorted permutation + pre-encoded Rows
    • SortedRunStream: PartitionedStream that lazily materializes batches and provides pre-computed RowValues to the merge (no re-encoding)
    • Schema-dependent reconstruction: concat+take for StringArray, interleave for StringView/fixed-width
    • sort_coalesce_target_rows config (default 32768)

    Profiling results (sort utf8 tuple 1M, BATCH_SIZE=8192)

    Profiled with samply on sort utf8 tuple 1M (3-column StringArray, the worst-case schema for data movement). Compared four variants:

    main (per-batch sort, no coalescing):

    Self % Function
    3.6% _platform_memmove
    3.5% arrow_row::variable::encode_one (merge phase only)
    ~1.8% quicksort partition

    es5 (es4 approach: concat keys + RowConverter sort + take/interleave):

    Self % Function
    7.8% _platform_memmove (concat + take)
    7.6% arrow_row::variable::encode_one (sort phase)
    7.6% RowConverter::append (sort phase)
    3.7% take_impl

    sort_rows v1 (encode-once, interleave from originals):

    Self % Function
    19% _platform_memmove
    12% interleave_bytes (flush + merge — double interleave)
    4.4% encode_one (accumulation phase only)
    14% RowConverter::append (accumulation — the one encode)

    sort_rows v2 (encode-once, concat+take for StringArray):

    Self % Function
    21% _platform_memmove
    4% take_bytes (from concat'd batch)
    6.1% interleave_bytes (merge output only — unavoidable)
    6% encode_one
    0.4% concat_batches

    What we learned

    The encode-once goal was achieved. RowConverter encoding dropped from double (es5: sort + merge) to single (sort_rows: accumulation only). The merge uses pre-computed Rows via SortedRunStream, completely bypassing RowCursorStream's per-batch encoding.

    But the reconstruction cost dominates. At 1M rows, moving string bytes from unsorted batches to sorted output is the bottleneck, not encoding. We tried three reconstruction strategies:

    1. Interleave from originals (v1): 12% interleave_bytes — scatter-gather across 4 source batch string buffers
    2. Concat + take (v2): 0.4% concat + 4% take_bytes — sequential concat into one buffer, then indexed access
    3. StringView intermediate (discussed, not implemented): would have same random-access pattern as (1) when converting back to StringArray

    All strategies move the same total bytes. The profile shapes are nearly identical — the work just shifts between functions. The fundamental cost is O(N × avg_string_length) string copying in permutation order, regardless of approach.

    The double interleave problem. sort_rows v1 had two interleave passes:

    1. SortedRun::next_chunk → interleave from originals (materializing the run)
    2. BatchBuilder::build_record_batch → interleave from runs (merge output)

    Deferred materialization (not eagerly building RecordBatches in flush) didn't help — the lazy materialization in next_chunk does the same work, just later. The merge's BatchBuilder interleave exists on main too and is unavoidable.

    The crossover point is scale-dependent. At 1M rows / 8192 batch size, main produces 122 runs. Our approach produces ~30 runs. The merge fan-in reduction (122→30) saves merge comparison work, but the coalescing overhead (concat or interleave + Rows encoding + Rows gathering) exceeds the savings. At TPC-H SF10 (~60M rows, 7300→1800 runs), the merge savings dominate and the approach wins (+15% total query time).

    Fundamental insight

    For StringArray-heavy schemas at moderate scale, the sort pipeline's cost is dominated by data movement (physically copying string bytes), not comparison (sort/merge key encoding). The encode-once architecture successfully eliminates redundant encoding work, but encoding was only ~15% of the total. The remaining 85% is moving bytes around, and no intermediate representation change (StringView, Rows, concat) reduces the total bytes that must move.

    Main's approach — sort each batch independently within its own buffer, then merge — minimizes cross-batch data movement. Each batch's sort is entirely within one contiguous buffer (cache-friendly). The only cross-batch operation is the merge output, which is unavoidable.

    Coalescing adds cross-batch data movement (concat, interleave, or take-from-concat) in exchange for fewer merge runs. The tradeoff only pays off when the merge cost (proportional to fan-in) dominates the data movement cost (proportional to data size). This happens at large scale (TPC-H SF10) but not at moderate scale (1M rows).

    Conclusions and Recommendations

    The core finding

    ExternalSorter's per-batch sort + merge architecture on main is already well-suited for moderate-scale workloads (up to ~10M rows) at realistic batch sizes (8192). The bottleneck at this scale is data movement (physically copying string bytes), not comparison or encoding overhead. No intermediate representation change (Rows, StringView, concat) reduces the total bytes that must move.

    At large scale (TPC-H SF10, ~60M rows), merge fan-in becomes the bottleneck. Reducing fan-in from 7300 to 1800 runs gives a 15% speedup. But the cheapest way to reduce fan-in is to increase input batch size, not to coalesce inside the sort.

    For DataFusion

    No ExternalSorter changes recommended at this time. The current implementation is efficient. The encode-once architecture (branch sort_rows) successfully eliminates double RowConverter encoding in the merge, but the savings (~15% of pipeline time) are offset by the additional data movement cost of coalescing at moderate scale.

    When lexsort_radix lands in arrow-rs (#9683), it provides 1.5-2.5x faster sorting on Rows at 32K+ rows. The encode-once architecture on sort_rows provides the infrastructure to adopt it (incremental RowConverter::append, SortedRun with permutation, SortedRunStream for merge handoff). The radix sort may shift the break-even point for coalescing to a lower scale, making the approach viable for moderate workloads. This should be revisited when the kernel is available.

    For Comet

    Increase batch size rather than redesign the sort. Comet sees performance wins at 32768 batch size (up from 8192 default). At 60M rows, this reduces merge fan-in from 7300 to 1800 — the same improvement that es4's coalescing achieved, but without any sort-internal overhead.

    The simplest approach: have Comet's shuffle writer produce larger batches when feeding a CometSort. The planner knows the plan topology and can set batch size based on the downstream operator. Sort wants large batches (fewer runs); other operators may prefer smaller batches for pipelining.

    Alternatively: increase Comet's default batch size globally to 32768 and validate that nothing regresses. The sort and merge both handle 32K batches efficiently.

    A Comet-specific sort operator may still be worthwhile for other reasons (working directly with Spark's row format, spilling in shuffle format, JNI-aware memory), but the fan-in problem is solvable without one.

  11. gratus00 commented on Apr 23, 2026

    @gratus00

    @mbutrovich I will take a deeper look into this, and will add the PR when I get the chance this week

    But yeah seems like a trade off one way or another so far

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions