Skip to content

Commit 73bb06b

Browse files
committed
Cleanup.
1 parent e0fb0a8 commit 73bb06b

1 file changed

Lines changed: 11 additions & 15 deletions

File tree

  • datafusion/physical-plan/src/sorts

‎datafusion/physical-plan/src/sorts/sort.rs‎

Lines changed: 11 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -347,9 +347,8 @@ impl ExternalSorter {
347347
/// Sorts a single coalesced batch and stores the result as a new run.
348348
///
349349
/// Uses radix sort for multi-column radix-eligible schemas when the
350-
/// batch is large enough to amortize RowConverter encoding overhead
351-
/// (more than `batch_size` rows). Falls back to lexsort for small
352-
/// batches (e.g. early flush under memory pressure).
350+
/// batch reached `sort_coalesce_target_rows`. Falls back to lexsort
351+
/// for smaller batches (e.g. the final partial flush).
353352
/// Both paths chunk output back to `batch_size`.
354353
fn sort_and_store_run(&mut self, batch: &RecordBatch) -> Result<()> {
355354
let use_radix =
@@ -1813,10 +1812,8 @@ mod tests {
18131812
);
18141813

18151814
// Build 200 partitions, each with one 100-row batch containing a
1816-
// Utf8 column (~42 bytes per value) and an Int32 column. The second
1817-
// column ensures the radix sort path is used (use_radix_sort
1818-
// requires > 1 sort column), which keeps the coalesce target at
1819-
// sort_coalesce_target_rows instead of falling back to batch_size.
1815+
// Utf8 column (~42 bytes per value) and an Int32 column. Two sort
1816+
// columns so use_radix_sort returns true and the radix path is exercised.
18201817
let schema = Arc::new(Schema::new(vec![
18211818
Field::new("i", DataType::Utf8, true),
18221819
Field::new("j", DataType::Int32, true),
@@ -1976,7 +1973,7 @@ mod tests {
19761973
SessionConfig::new().with_batch_size(batch_size), // Ensure we don't concat batches
19771974
));
19781975

1979-
// Two columns so the radix path is used and sorted runs are chunked.
1976+
// Two columns so use_radix_sort returns true.
19801977
let schema = Arc::new(Schema::new(vec![
19811978
Field::new("a", DataType::Int32, false),
19821979
Field::new("b", DataType::Int32, false),
@@ -2574,8 +2571,7 @@ mod tests {
25742571
batch_size_to_generate: usize,
25752572
create_task_ctx: impl Fn(&[RecordBatch]) -> TaskContext,
25762573
) -> Result<MetricsSet> {
2577-
// Two-column batches so use_radix_sort returns true and sorted runs
2578-
// are chunked to batch_size, which these tests depend on.
2574+
// Two-column batches so use_radix_sort returns true.
25792575
let schema = Arc::new(Schema::new(vec![
25802576
Field::new("i", DataType::Int32, true),
25812577
Field::new("j", DataType::Int32, true),
@@ -2853,7 +2849,7 @@ mod tests {
28532849
.build_arc()?;
28542850

28552851
let metrics_set = ExecutionPlanMetricsSet::new();
2856-
// Two columns so radix sort is used, matching the coalesce_target_rows config.
2852+
// Two columns so use_radix_sort returns true.
28572853
let schema = Arc::new(Schema::new(vec![
28582854
Field::new("x", DataType::Int32, false),
28592855
Field::new("y", DataType::Int32, false),
@@ -3279,7 +3275,7 @@ mod tests {
32793275
/// and chunked back to `batch_size` after sorting.
32803276
#[tokio::test]
32813277
async fn test_chunked_sort_radix_coalescing() -> Result<()> {
3282-
// Two sort columns required so use_radix_sort returns true.
3278+
// Two sort columns so use_radix_sort returns true.
32833279
let schema = Arc::new(Schema::new(vec![
32843280
Field::new("x", DataType::Int32, false),
32853281
Field::new("y", DataType::Int32, false),
@@ -3323,7 +3319,7 @@ mod tests {
33233319
/// the partial coalescer contents are flushed and sorted.
33243320
#[tokio::test]
33253321
async fn test_chunked_sort_partial_flush() -> Result<()> {
3326-
// Two sort columns required so use_radix_sort returns true.
3322+
// Two sort columns so use_radix_sort returns true.
33273323
let schema = Arc::new(Schema::new(vec![
33283324
Field::new("x", DataType::Int32, false),
33293325
Field::new("y", DataType::Int32, false),
@@ -3364,7 +3360,7 @@ mod tests {
33643360
/// Spilling writes one spill file per sorted run (no merge before spill).
33653361
#[tokio::test]
33663362
async fn test_spill_creates_one_file_per_run() -> Result<()> {
3367-
// Two sort columns required so use_radix_sort returns true.
3363+
// Two sort columns so use_radix_sort returns true.
33683364
let schema = Arc::new(Schema::new(vec![
33693365
Field::new("x", DataType::Int32, false),
33703366
Field::new("y", DataType::Int32, false),
@@ -3424,7 +3420,7 @@ mod tests {
34243420
/// merged into a single sorted stream before spilling to one file.
34253421
#[tokio::test]
34263422
async fn test_spill_merges_runs_with_headroom() -> Result<()> {
3427-
// Two sort columns required so use_radix_sort returns true.
3423+
// Two sort columns so use_radix_sort returns true.
34283424
let schema = Arc::new(Schema::new(vec![
34293425
Field::new("x", DataType::Int32, false),
34303426
Field::new("y", DataType::Int32, false),

0 commit comments

Comments
 (0)