Skip to content

Commit fef01ed

Browse files
committed
test: cover unsorted contiguous groups in one partition
1 parent 1b6dc92 commit fef01ed

1 file changed

Lines changed: 77 additions & 0 deletions

File tree

  • datafusion/physical-plan/src/aggregates

‎datafusion/physical-plan/src/aggregates/mod.rs‎

Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4833,6 +4833,83 @@ mod tests {
48334833
Ok(())
48344834
}
48354835

4836+
#[tokio::test]
4837+
async fn unsorted_contiguous_groups_use_final_emission() -> Result<()> {
4838+
let schema = Arc::new(Schema::new(vec![
4839+
Field::new("key", DataType::Int32, false),
4840+
Field::new("time_bin", DataType::Int64, false),
4841+
Field::new("value", DataType::Int64, false),
4842+
]));
4843+
// Two sorted logical runs are emitted as batches in one DataFusion
4844+
// partition. Every distinct grouping tuple occupies one contiguous range,
4845+
// but tuple order resets at the batch boundary, so (key, time_bin) is not
4846+
// globally sorted.
4847+
let input_batches = vec![
4848+
RecordBatch::try_new(
4849+
Arc::clone(&schema),
4850+
vec![
4851+
Arc::new(Int32Array::from(vec![1, 1, 2, 2])),
4852+
Arc::new(Int64Array::from(vec![20, 20, 20, 20])),
4853+
Arc::new(Int64Array::from(vec![10, 20, 30, 40])),
4854+
],
4855+
)?,
4856+
RecordBatch::try_new(
4857+
Arc::clone(&schema),
4858+
vec![
4859+
Arc::new(Int32Array::from(vec![1, 1, 2, 2])),
4860+
Arc::new(Int64Array::from(vec![0, 0, 0, 0])),
4861+
Arc::new(Int64Array::from(vec![50, 60, 70, 80])),
4862+
],
4863+
)?,
4864+
];
4865+
let group_by = PhysicalGroupBy::new_single(vec![
4866+
(col("key", &schema)?, "key".to_string()),
4867+
(col("time_bin", &schema)?, "time_bin".to_string()),
4868+
]);
4869+
let aggr_expr = Arc::new(
4870+
AggregateExprBuilder::new(sum_udaf(), vec![col("value", &schema)?])
4871+
.schema(Arc::clone(&schema))
4872+
.alias("SUM(value)")
4873+
.build()?,
4874+
);
4875+
let input: Arc<dyn ExecutionPlan> =
4876+
TestMemoryExec::try_new_exec(&[input_batches], Arc::clone(&schema), None)?;
4877+
assert_eq!(input.output_partitioning().partition_count(), 1);
4878+
4879+
let aggregate = AggregateExec::try_new(
4880+
AggregateMode::Single,
4881+
group_by,
4882+
vec![aggr_expr],
4883+
vec![None],
4884+
input,
4885+
schema,
4886+
)?;
4887+
4888+
assert_eq!(aggregate.input_order_mode(), &InputOrderMode::Linear);
4889+
// This captures the behavior before #24438. When the source can declare
4890+
// `(key, time_bin)` group-contiguous, the corresponding case can use
4891+
// `EmissionType::Incremental`.
4892+
assert_eq!(aggregate.cache().emission_type, EmissionType::Final);
4893+
4894+
let task_ctx = new_migrated_hash_ctx(1024);
4895+
let stream = aggregate.execute_typed(0, &task_ctx)?;
4896+
assert!(matches!(stream, StreamType::SingleHash(_)));
4897+
let stream: SendableRecordBatchStream = stream.into();
4898+
let output = collect(stream).await?;
4899+
assert_snapshot!(batches_to_sort_string(&output), @r"
4900+
+-----+----------+------------+
4901+
| key | time_bin | SUM(value) |
4902+
+-----+----------+------------+
4903+
| 1 | 0 | 110 |
4904+
| 1 | 20 | 30 |
4905+
| 2 | 0 | 150 |
4906+
| 2 | 20 | 70 |
4907+
+-----+----------+------------+
4908+
");
4909+
4910+
Ok(())
4911+
}
4912+
48364913
/// Ensures for ordered input, `OrderedPartialAggregateStream` is used.
48374914
#[tokio::test]
48384915
async fn ordered_partial_aggregate_planning() -> Result<()> {

0 commit comments

Comments
 (0)