kosiew commented on code in PR #25358:
URL: https://github.com/apache/datafusion/pull/25358#discussion_r4183441601
##########
datafusion/physical-plan/src/sorts/sort.rs:
##########
@@ -957,6 +1001,141 @@ impl Debug for ExternalSorter {
}
}
+fn staged_sort_batch_chunked(
+ batch: &RecordBatch,
+ expressions: &LexOrdering,
+ batch_size: usize,
+ group_size: usize,
+) -> Result<Vec<RecordBatch>> {
+ let indices = staged_sort_indices(batch, expressions, group_size)?;
+ IncrementalSortIterator::new(batch.clone(), expressions.clone(),
batch_size)
+ .with_sorted_indices(indices)
+ .collect()
+}
+
+/// Compute a complete permutation by evaluating successive key groups only
for ties.
+fn staged_sort_indices(
+ batch: &RecordBatch,
+ expressions: &LexOrdering,
+ group_size: usize,
+) -> Result<UInt32Array> {
+ let mut indices = Vec::new();
+ // Ranges address positions in `indices`, whose values address the input
batch.
+ let mut unresolved: Vec<_> =
std::iter::once(0..batch.num_rows()).collect();
+ let stage_count = expressions.len().div_ceil(group_size);
+ for (stage, keys) in expressions.chunks(group_size).enumerate() {
+ if unresolved.is_empty() {
+ break;
+ }
+ unresolved = refine_tied_groups(
+ batch,
+ keys,
+ &mut indices,
+ &unresolved,
+ stage == 0,
+ stage + 1 < stage_count,
+ )?;
+ }
+
+ Ok(UInt32Array::from(indices))
+}
+
+/// Evaluate keys once for all unresolved rows, then sort each tied group
independently.
+fn refine_tied_groups(
+ batch: &RecordBatch,
+ keys: &[PhysicalSortExpr],
+ indices: &mut Vec<u32>,
+ unresolved: &[Range<usize>],
+ first_stage: bool,
+ find_ties: bool,
+) -> Result<Vec<Range<usize>>> {
+ let selected = if first_stage {
+ batch.clone()
+ } else {
+ // Gather the original row indices for all tied groups into one buffer.
+ let mut selection =
Vec::with_capacity(unresolved.iter().map(Range::len).sum());
+ for range in unresolved {
+ selection.extend_from_slice(&indices[range.clone()]);
+ }
+ arrow::compute::take_record_batch(batch,
&UInt32Array::from(selection))?
+ };
+ let sort_columns = keys
Review Comment:
Staged sorting can skip later fallible or volatile ORDER BY expressions when
earlier keys already resolve the order, so queries can behave differently
depending on tie distribution. Please either document this
experimental/error/volatile-evaluation contract and add tests that pin it,
including that disabled staging stays eager, or preserve eager evaluation.
##########
datafusion/physical-plan/src/sorts/sort.rs:
##########
@@ -957,6 +1001,141 @@ impl Debug for ExternalSorter {
}
}
+fn staged_sort_batch_chunked(
Review Comment:
The new staged sorting path does not have focused behavioral tests. Please
add coverage for staged versus eager equivalence, successive tie groups, null
ordering, chunked output, the zero-tie/error case, and both single-batch and
multiple in-memory-run paths.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]