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]

Reply via email to