This is an automated email from the ASF dual-hosted git repository.
alamb pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-rs.git
The following commit(s) were added to refs/heads/main by this push:
new 9b190c9794 perf(arrow-ord): Avoid full index materialization for
small-limit lexsorts (#9991)
9b190c9794 is described below
commit 9b190c9794c6948314bdaee774d0e3d4fc355fc5
Author: pchintar <[email protected]>
AuthorDate: Mon Jun 29 21:13:13 2026 +0530
perf(arrow-ord): Avoid full index materialization for small-limit lexsorts
(#9991)
# Which issue does this PR close?
- Closes #9990 .
# Rationale for this change
In `arrow-ord/src/sort.rs`, `lexsort_to_indices` currently materializes
row indices for the full input before applying the requested lexsort
limit.
For example, with:
```text
row_count = 4096
limit = 10
```
the current implementation still allocates and initializes indices for
all 4096 rows:
```rust
let mut value_indices = (0..row_count).collect::<Vec<usize>>();
```
even though only the first 10 sorted indices are returned.
This PR reduces allocation and sorting work for small-limit lexsorts by
avoiding full index materialization when the requested limit is a small
fraction(limit < 1/10 th of row count) of the input size.
# What changes are included in this PR?
This PR adds a bounded top-k heap path for small-limit lexsorts in
`arrow-ord/src/sort.rs`.
The new path:
* maintains only the current top-k row indices while scanning the input
* avoids allocating indices for every row
* sorts only the retained top-k indices before returning the result
The existing partial-sort implementation remains unchanged for larger
limits.
To support the bounded heap implementation, this PR also adds:
- `lexsort_topk_fixed`
- Uses the existing fixed-column comparator with the bounded heap path.
- `lexsort_topk`
- Keeps only the current top-k row indices in a bounded max-heap while
scanning rows.
- `sift_up_worst_heap`
- Moves a newly inserted row index upward until the heap order is
restored.
- `sift_down_worst_heap`
- Moves the current worst retained row index downward after replacement.
# Are these changes tested?
Yes.
Existing tests pass:
```text
cargo test -p arrow-ord
cargo test -p arrow --features test_utils
```
My local Benchmark results from:
```text
cargo bench --package arrow --bench sort_kernel --features test_utils --
"limit"
```
show improvements for small-limit lexsorts such as limits 10, 100.
# Are there any user-facing changes?
No.
---
arrow-ord/src/sort.rs | 180 ++++++++++++++++++++++++++++++++++++++++++--------
1 file changed, 152 insertions(+), 28 deletions(-)
diff --git a/arrow-ord/src/sort.rs b/arrow-ord/src/sort.rs
index 114ca05bb7..3156063480 100644
--- a/arrow-ord/src/sort.rs
+++ b/arrow-ord/src/sort.rs
@@ -949,40 +949,56 @@ pub fn lexsort_to_indices(
));
};
- let mut value_indices = (0..row_count).collect::<Vec<usize>>();
- let mut len = value_indices.len();
+ let len = limit.unwrap_or(row_count).min(row_count);
- if let Some(limit) = limit {
- len = limit.min(len);
+ if len == 0 {
+ return Ok(UInt32Array::from(Vec::<u32>::new()));
}
- // Instantiate specialized versions of comparisons for small numbers
- // of columns as it helps the compiler generate better code.
- match columns.len() {
- 2 => {
- sort_fixed_column::<2>(columns, &mut value_indices, len)?;
- }
- 3 => {
- sort_fixed_column::<3>(columns, &mut value_indices, len)?;
- }
- 4 => {
- sort_fixed_column::<4>(columns, &mut value_indices, len)?;
- }
- 5 => {
- sort_fixed_column::<5>(columns, &mut value_indices, len)?;
- }
+ // The heap path avoids allocating and partially sorting all row indices
+ // when the requested limit is a small fraction of the input. For larger
+ // limits, the existing partial-sort path is preferred because heap
+ // maintenance costs grow with the requested limit.
+ let value_indices = match limit {
+ Some(limit) if limit <= row_count / 10 => match columns.len() {
+ 2 => lexsort_topk_fixed::<2>(columns, row_count, len)?,
+ 3 => lexsort_topk_fixed::<3>(columns, row_count, len)?,
+ 4 => lexsort_topk_fixed::<4>(columns, row_count, len)?,
+ 5 => lexsort_topk_fixed::<5>(columns, row_count, len)?,
+ _ => {
+ let lexicographical_comparator =
LexicographicalComparator::try_new(columns)?;
+ lexsort_topk(row_count, len, |a, b| {
+ lexicographical_comparator.compare(a, b)
+ })
+ }
+ },
_ => {
- let lexicographical_comparator =
LexicographicalComparator::try_new(columns)?;
- // uint32 can be sorted unstably
- sort_unstable_by(&mut value_indices, len, |a, b| {
- lexicographical_comparator.compare(*a, *b)
- });
+ let mut value_indices = (0..row_count).collect::<Vec<usize>>();
+
+ // Instantiate specialized versions of comparisons for small
numbers
+ // of columns as it helps the compiler generate better code.
+ match columns.len() {
+ 2 => sort_fixed_column::<2>(columns, &mut value_indices, len)?,
+ 3 => sort_fixed_column::<3>(columns, &mut value_indices, len)?,
+ 4 => sort_fixed_column::<4>(columns, &mut value_indices, len)?,
+ 5 => sort_fixed_column::<5>(columns, &mut value_indices, len)?,
+ _ => {
+ let lexicographical_comparator =
LexicographicalComparator::try_new(columns)?;
+ sort_unstable_by(&mut value_indices, len, |a, b| {
+ lexicographical_comparator.compare(*a, *b)
+ });
+ }
+ }
+
+ value_indices.truncate(len);
+ value_indices
}
- }
+ };
+
Ok(UInt32Array::from(
- value_indices[..len]
- .iter()
- .map(|i| *i as u32)
+ value_indices
+ .into_iter()
+ .map(|i| i as u32)
.collect::<Vec<_>>(),
))
}
@@ -1000,6 +1016,90 @@ fn sort_fixed_column<const N: usize>(
Ok(())
}
+// Uses the fixed-column comparator for the bounded heap path.
+fn lexsort_topk_fixed<const N: usize>(
+ columns: &[SortColumn],
+ row_count: usize,
+ limit: usize,
+) -> Result<Vec<usize>, ArrowError> {
+ let lexicographical_comparator =
FixedLexicographicalComparator::<N>::try_new(columns)?;
+ Ok(lexsort_topk(row_count, limit, |a, b| {
+ lexicographical_comparator.compare(a, b)
+ }))
+}
+
+// Keeps the smallest `limit` indices in a bounded max-heap.
+// The root is the largest retained index according to `compare`.
+fn lexsort_topk(
+ row_count: usize,
+ limit: usize,
+ mut compare: impl FnMut(usize, usize) -> Ordering,
+) -> Vec<usize> {
+ let mut heap = Vec::with_capacity(limit);
+
+ for idx in 0..row_count {
+ if heap.len() < limit {
+ heap.push(idx);
+ let pos = heap.len() - 1;
+ sift_up_worst_heap(&mut heap, pos, &mut compare);
+ } else if compare(idx, heap[0]) == Ordering::Less {
+ heap[0] = idx;
+ sift_down_worst_heap(&mut heap, 0, &mut compare);
+ }
+ }
+
+ heap.sort_unstable_by(|a, b| compare(*a, *b));
+ heap
+}
+
+// Moves a newly inserted index toward the root while it is larger than its
parent.
+fn sift_up_worst_heap(
+ heap: &mut [usize],
+ mut pos: usize,
+ compare: &mut impl FnMut(usize, usize) -> Ordering,
+) {
+ while pos > 0 {
+ let parent = (pos - 1) / 2;
+
+ if compare(heap[parent], heap[pos]) != Ordering::Less {
+ break;
+ }
+
+ heap.swap(parent, pos);
+ pos = parent;
+ }
+}
+
+// Moves the root down until both children are no larger than it.
+// The larger child is selected at each step so the worst retained row remains
+// at heap[0].
+fn sift_down_worst_heap(
+ heap: &mut [usize],
+ mut pos: usize,
+ compare: &mut impl FnMut(usize, usize) -> Ordering,
+) {
+ loop {
+ let left = pos * 2 + 1;
+ if left >= heap.len() {
+ break;
+ }
+
+ let right = left + 1;
+ let mut worst = left;
+
+ if right < heap.len() && compare(heap[left], heap[right]) ==
Ordering::Less {
+ worst = right;
+ }
+
+ if compare(heap[pos], heap[worst]) != Ordering::Less {
+ break;
+ }
+
+ heap.swap(pos, worst);
+ pos = worst;
+ }
+}
+
/// It's unstable_sort, may not preserve the order of equal elements
pub fn partial_sort<T, F>(v: &mut [T], limit: usize, mut is_less: F)
where
@@ -4483,6 +4583,30 @@ mod tests {
run_test!(LargeListArray);
}
+ #[test]
+ fn test_lex_sort_limit_paths() {
+ let input = vec![
+ SortColumn {
+ values: Arc::new(Int32Array::from_iter_values((0..100).rev()))
as ArrayRef,
+ options: None,
+ },
+ SortColumn {
+ values: Arc::new(Int32Array::from_iter_values(0..100)) as
ArrayRef,
+ options: None,
+ },
+ ];
+
+ // Exercise the bounded heap path.
+ let indices = lexsort_to_indices(&input, Some(10)).unwrap();
+ let expected = UInt32Array::from_iter_values((90..100).rev());
+ assert_eq!(indices, expected);
+
+ // Exercise the existing partial-sort path.
+ let indices = lexsort_to_indices(&input, Some(11)).unwrap();
+ let expected = UInt32Array::from_iter_values((89..100).rev());
+ assert_eq!(indices, expected);
+ }
+
#[test]
fn test_partial_sort() {
let mut before: Vec<&str> = vec![