[
https://issues.apache.org/jira/browse/SPARK-59697?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
David Mollitor updated SPARK-59697:
-----------------------------------
Component/s: Spark Core
(was: SQL)
> Skip the redundant record comparator in the external-sort spill merge when
> the key prefix is a total order
> ----------------------------------------------------------------------------------------------------------
>
> Key: SPARK-59697
> URL: https://issues.apache.org/jira/browse/SPARK-59697
> Project: Spark
> Issue Type: Improvement
> Components: Spark Core
> Affects Versions: 4.1.0
> Reporter: David Mollitor
> Priority: Minor
>
> h2. Summary
> The external-sort spill merge ({{{}UnsafeSorterSpillMerger{}}} /
> {{{}UnsafeSorterBoundedSpillMerger{}}}, reached through
> {{{}UnsafeExternalSorter.getSortedIterator(){}}}) orders spill-run heads by
> the 8-byte key prefix and, on a prefix tie, falls back to the full
> {{{}RecordComparator{}}}, which decodes and byte-compares the records:
> {code:java}
> int prefixCmp = prefixComparator.compare(left.getKeyPrefix(),
> right.getKeyPrefix());
> if (prefixCmp == 0) {
> return recordComparator.compare(/* decode + byte-compare both records */);
> }
> return prefixCmp;
> {code}
> When the sort qualifies for radix sort ({{{}canUseRadixSort{}}}: a single,
> prefix-sortable key), the 8-byte prefix is a lossless, order-preserving
> _total order_ for that key. Equal prefixes are then equal keys, so the
> {{RecordComparator}} tie-break always returns 0 – it is dead work on every
> prefix collision.
> This change threads the already-computed {{canUseRadixSort}} flag (today only
> forwarded to {{UnsafeInMemorySorter}} and then dropped) into both merge paths
> – the single-round {{UnsafeSorterSpillMerger}} and the multi-round
> {{UnsafeSorterBoundedSpillMerger}} – so the record-comparator tie-break is
> skipped on prefix ties in that case. It is the merge-side analogue of the
> in-memory radix-sort optimization.
> h2. Why
> {{UnsafeExternalSorter}} backs {{SortExec}} and key-based aggregation /
> window / join (via {{{}UnsafeKVExternalSorter{}}}). On a spilled sort with a
> single prefix-sortable key, every prefix collision in the merge currently
> pays for a full record decode and comparison that provably returns 0.
> Removing it speeds up the merge; the benefit is largest for low-cardinality
> or duplicate-heavy keys where prefix ties are frequent, and neutral otherwise.
> h2. Correctness
> Skipping the tie-break is safe because {{canUseRadixSort}} guarantees the
> prefix fully determines the sort order (equal prefix == equal key), so the
> record comparator would have returned 0 anyway.
> This is the same precondition the in-memory radix sort already relies on:
> each spilled run was itself ordered by prefix alone, so merging by prefix
> alone is consistent. {{canUseRadixSort}} is {{{}enableRadixSort &&
> sortOrder.length == 1 && SortPrefixUtils.canSortFullyWithPrefix(...){}}}, and
> {{canSortFullyWithPrefix}} restricts it to types whose 8-byte prefix is a
> lossless, order-preserving encoding: integral / date / time / interval, float
> / double, and decimals with {{{}precision <= Decimal.MAX_LONG_DIGITS{}}}.
> Sorts with {{canUseRadixSort = false}} (multi-key sorts, strings / binary,
> large decimals) keep the record-comparator tie-break unchanged.
> h2. Scope
> Core only ({{{}org.apache.spark.util.collection.unsafe.sort{}}}). No SQL-side
> changes are needed – {{canUseRadixSort}} already reaches
> {{UnsafeExternalSorter.create()}} from the SQL callers. No new configuration
> and no user-facing behavior change.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]