Skip to content

Commit 6718bbf

Browse files
committed
feat: stream aggregates for group-contiguous input
1 parent f951020 commit 6718bbf

27 files changed

Lines changed: 1383 additions & 138 deletions

‎datafusion/core/tests/physical_optimizer/combine_partial_final_agg.rs‎

Lines changed: 47 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,12 +35,14 @@ use datafusion_physical_expr::expressions::{col, lit};
3535
use datafusion_physical_expr_common::physical_expr::PhysicalExpr;
3636
use datafusion_physical_optimizer::PhysicalOptimizerRule;
3737
use datafusion_physical_optimizer::combine_partial_final_agg::CombinePartialFinalAggregate;
38-
use datafusion_physical_plan::ExecutionPlan;
3938
use datafusion_physical_plan::aggregates::{
4039
AggregateExec, AggregateMode, LimitOptions, PhysicalGroupBy,
4140
};
4241
use datafusion_physical_plan::displayable;
42+
use datafusion_physical_plan::execution_plan::EmissionType;
4343
use datafusion_physical_plan::repartition::RepartitionExec;
44+
use datafusion_physical_plan::test::TestMemoryExec;
45+
use datafusion_physical_plan::{ExecutionPlan, ExecutionPlanProperties};
4446

4547
/// Runs the CombinePartialFinalAggregate optimizer and asserts the plan against the expected
4648
macro_rules! assert_optimized {
@@ -233,6 +235,50 @@ fn aggregations_with_group_combined() -> datafusion_common::Result<()> {
233235
Ok(())
234236
}
235237

238+
#[test]
239+
fn partition_disjoint_aggregation_combines_and_streams() -> datafusion_common::Result<()>
240+
{
241+
let schema = schema();
242+
let source = TestMemoryExec::try_new(&[vec![]], Arc::clone(&schema), None)?
243+
.try_with_group_contiguous_keys(vec![col("c", &schema)?])?;
244+
let aggr_expr = vec![Arc::new(
245+
AggregateExprBuilder::new(sum_udaf(), vec![col("b", &schema)?])
246+
.schema(Arc::clone(&schema))
247+
.alias("Sum(b)")
248+
.build()?,
249+
)];
250+
let partial = partial_aggregate_exec(
251+
Arc::new(source),
252+
PhysicalGroupBy::new_single(vec![(col("c", &schema)?, "c".to_string())]),
253+
aggr_expr.clone(),
254+
);
255+
let final_group_by = PhysicalGroupBy::new_single(vec![(
256+
col("c", &partial.schema())?,
257+
"c".to_string(),
258+
)]);
259+
let plan = final_aggregate_exec(partial, final_group_by, aggr_expr);
260+
261+
let optimized =
262+
CombinePartialFinalAggregate::default().optimize(plan, &ConfigOptions::new())?;
263+
let aggregate = optimized
264+
.downcast_ref::<AggregateExec>()
265+
.expect("adjacent partial and final aggregates should combine");
266+
assert_eq!(aggregate.mode(), &AggregateMode::Single);
267+
assert_eq!(
268+
aggregate.input_order_mode(),
269+
&datafusion_physical_plan::InputOrderMode::Linear
270+
);
271+
assert_eq!(optimized.pipeline_behavior(), EmissionType::Incremental);
272+
273+
let display = displayable(optimized.as_ref()).indent(true).to_string();
274+
assert_snapshot!(display.trim(), @r"
275+
AggregateExec: mode=Single, gby=[c@2 as c], aggr=[Sum(b)], group_completion_mode=Full
276+
DataSourceExec: partitions=1, partition_sizes=[0], group_contiguous=[c@2]
277+
");
278+
279+
Ok(())
280+
}
281+
236282
#[test]
237283
fn aggregations_with_limit_combined() -> datafusion_common::Result<()> {
238284
let schema = schema();

‎datafusion/core/tests/physical_optimizer/enforce_distribution.rs‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,7 @@ use datafusion_physical_plan::joins::utils::JoinOn;
7272
use datafusion_physical_plan::limit::{GlobalLimitExec, LocalLimitExec};
7373
use datafusion_physical_plan::projection::{ProjectionExec, ProjectionExpr};
7474
use datafusion_physical_plan::sorts::sort_preserving_merge::SortPreservingMergeExec;
75+
use datafusion_physical_plan::test::TestMemoryExec;
7576
use datafusion_physical_plan::union::UnionExec;
7677
use datafusion_physical_plan::{
7778
ChildrenPropertiesMode, DisplayAs, DisplayFormatType, ExecutionPlanProperties,
@@ -2617,6 +2618,29 @@ fn added_repartition_to_single_partition() -> Result<()> {
26172618
Ok(())
26182619
}
26192620

2621+
#[test]
2622+
fn partition_disjoint_aggregate_preserves_source_partitioning() -> Result<()> {
2623+
let schema = schema();
2624+
let input = TestMemoryExec::try_new(&[vec![]], Arc::clone(&schema), None)?
2625+
.try_with_group_contiguous_keys(vec![col("a", &schema)?])?;
2626+
let aggregate = Arc::new(AggregateExec::try_new(
2627+
AggregateMode::Partial,
2628+
PhysicalGroupBy::new_single(vec![(col("a", &schema)?, "a".to_string())]),
2629+
vec![],
2630+
vec![],
2631+
Arc::new(input),
2632+
schema,
2633+
)?) as Arc<dyn ExecutionPlan>;
2634+
2635+
let plan = TestConfig::default().to_plan(aggregate, &DISTRIB_DISTRIB_SORT);
2636+
assert_plan!(plan, @r"
2637+
AggregateExec: mode=Partial, gby=[a@0 as a], aggr=[], group_completion_mode=Full
2638+
DataSourceExec: partitions=1, partition_sizes=[0], group_contiguous=[a@0]
2639+
");
2640+
2641+
Ok(())
2642+
}
2643+
26202644
#[test]
26212645
fn repartition_deepest_node() -> Result<()> {
26222646
let alias = vec![("a".to_string(), "a".to_string())];

‎datafusion/datasource/src/source.rs‎

Lines changed: 29 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -165,6 +165,20 @@ pub trait DataSource: Any + Send + Sync + Debug {
165165

166166
fn output_partitioning(&self) -> Partitioning;
167167
fn eq_properties(&self) -> EquivalenceProperties;
168+
169+
/// Components of one composite key whose tuple values occur in one
170+
/// contiguous range within each output stream.
171+
///
172+
/// See [`ExecutionPlan::group_contiguous_exprs`] for the full contract.
173+
/// Implementations must return expressions in terms of the schema returned
174+
/// by [`Self::eq_properties`]. Any implementation that returns a rewritten
175+
/// source from [`Self::repartitioned`] or [`Self::try_swapping_with_projection`]
176+
/// must remap, preserve, or clear these expressions so the contract remains
177+
/// valid for the rewritten output streams.
178+
fn group_contiguous_exprs(&self) -> &[Arc<dyn PhysicalExpr>] {
179+
&[]
180+
}
181+
168182
fn scheduling_type(&self) -> SchedulingType {
169183
SchedulingType::NonCooperative
170184
}
@@ -384,7 +398,20 @@ impl DisplayAs for DataSourceExec {
384398
}
385399
DisplayFormatType::TreeRender => {}
386400
}
387-
self.data_source.fmt_as(t, f)
401+
self.data_source.fmt_as(t, f)?;
402+
if matches!(t, DisplayFormatType::Default | DisplayFormatType::Verbose)
403+
&& !self.group_contiguous_exprs().is_empty()
404+
{
405+
write!(
406+
f,
407+
", group_contiguous=[{}]",
408+
self.group_contiguous_exprs()
409+
.iter()
410+
.map(ToString::to_string)
411+
.join(", ")
412+
)?;
413+
}
414+
Ok(())
388415
}
389416
}
390417

@@ -674,6 +701,7 @@ impl DataSourceExec {
674701
EmissionType::Incremental,
675702
Boundedness::Bounded,
676703
)
704+
.with_group_contiguous_exprs(data_source.group_contiguous_exprs().to_vec())
677705
.with_scheduling_type(data_source.scheduling_type())
678706
}
679707

‎datafusion/ffi/src/execution_plan.rs‎

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -615,6 +615,17 @@ pub mod tests {
615615
self.dynamic_expressions = dynamic_expressions;
616616
self
617617
}
618+
619+
pub fn with_group_contiguous_exprs(
620+
mut self,
621+
group_contiguous_exprs: Vec<Arc<dyn PhysicalExpr>>,
622+
) -> Self {
623+
self.props = Arc::new(
624+
PlanProperties::clone(&self.props)
625+
.with_group_contiguous_exprs(group_contiguous_exprs),
626+
);
627+
self
628+
}
618629
}
619630

620631
impl DisplayAs for EmptyExec {
@@ -787,6 +798,30 @@ pub mod tests {
787798
Ok(())
788799
}
789800

801+
#[test]
802+
fn test_ffi_execution_plan_group_contiguous_exprs() -> Result<()> {
803+
use datafusion_physical_expr::expressions::col;
804+
805+
let schema = Arc::new(arrow::datatypes::Schema::new(vec![
806+
arrow::datatypes::Field::new("key", arrow::datatypes::DataType::Int32, false),
807+
]));
808+
let expression = col("key", &schema)?;
809+
let original_plan = Arc::new(
810+
EmptyExec::new(schema)
811+
.with_group_contiguous_exprs(vec![Arc::clone(&expression)]),
812+
);
813+
814+
let mut ffi_plan = FFI_ExecutionPlan::new(original_plan, None);
815+
ffi_plan.library_marker_id = crate::mock_foreign_marker_id;
816+
let foreign_plan: Arc<dyn ExecutionPlan> = (&ffi_plan).try_into()?;
817+
818+
let group_contiguous_exprs = foreign_plan.group_contiguous_exprs();
819+
assert_eq!(group_contiguous_exprs.len(), 1);
820+
assert!(group_contiguous_exprs[0].eq(&expression));
821+
822+
Ok(())
823+
}
824+
790825
#[test]
791826
fn test_ffi_execution_plan_children() -> Result<()> {
792827
let schema = Arc::new(arrow::datatypes::Schema::new(vec![

‎datafusion/ffi/src/plan_properties.rs‎

Lines changed: 35 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ use datafusion_physical_plan::execution_plan::{Boundedness, EmissionType};
2828
use stabby::vec::Vec as SVec;
2929

3030
use crate::arrow_wrappers::WrappedSchema;
31+
use crate::physical_expr::FFI_PhysicalExpr;
3132
use crate::physical_expr::partitioning::FFI_Partitioning;
3233
use crate::physical_expr::sort::FFI_PhysicalSortExpr;
3334
use crate::util::FFI_Option;
@@ -63,6 +64,13 @@ pub struct FFI_PlanProperties {
6364
/// the foreign interface. See [`crate::get_library_marker_id`] and
6465
/// the crate's `README.md` for more information.
6566
pub library_marker_id: extern "C" fn() -> usize,
67+
68+
/// Return the components of the plan's group-contiguous composite key.
69+
///
70+
/// See [`datafusion_physical_plan::ExecutionPlan::group_contiguous_exprs`]
71+
/// for the correctness contract.
72+
pub group_contiguous_exprs:
73+
unsafe extern "C" fn(plan: &Self) -> SVec<FFI_PhysicalExpr>,
6674
}
6775

6876
struct PlanPropertiesPrivateData {
@@ -109,6 +117,18 @@ unsafe extern "C" fn output_ordering_fn_wrapper(
109117
ordering.into()
110118
}
111119

120+
unsafe extern "C" fn group_contiguous_exprs_fn_wrapper(
121+
properties: &FFI_PlanProperties,
122+
) -> SVec<FFI_PhysicalExpr> {
123+
properties
124+
.inner()
125+
.group_contiguous_exprs()
126+
.iter()
127+
.cloned()
128+
.map(FFI_PhysicalExpr::from)
129+
.collect()
130+
}
131+
112132
unsafe extern "C" fn schema_fn_wrapper(properties: &FFI_PlanProperties) -> WrappedSchema {
113133
let schema: SchemaRef = Arc::clone(properties.inner().eq_properties.schema());
114134
schema.into()
@@ -145,6 +165,7 @@ impl From<&PlanProperties> for FFI_PlanProperties {
145165
release: release_fn_wrapper,
146166
private_data: Box::into_raw(private_data) as *mut c_void,
147167
library_marker_id: crate::get_library_marker_id,
168+
group_contiguous_exprs: group_contiguous_exprs_fn_wrapper,
148169
}
149170
}
150171
}
@@ -186,12 +207,16 @@ impl TryFrom<FFI_PlanProperties> for PlanProperties {
186207
let boundedness: Boundedness =
187208
unsafe { (ffi_props.boundedness)(&ffi_props).into() };
188209

189-
Ok(PlanProperties::new(
190-
eq_properties,
191-
partitioning,
192-
emission_type,
193-
boundedness,
194-
))
210+
let group_contiguous_exprs =
211+
unsafe { (ffi_props.group_contiguous_exprs)(&ffi_props) }
212+
.iter()
213+
.map(<Arc<dyn datafusion_physical_expr::PhysicalExpr>>::from)
214+
.collect();
215+
216+
Ok(
217+
PlanProperties::new(eq_properties, partitioning, emission_type, boundedness)
218+
.with_group_contiguous_exprs(group_contiguous_exprs),
219+
)
195220
}
196221
}
197222

@@ -273,16 +298,16 @@ mod tests {
273298
let schema =
274299
Arc::new(Schema::new(vec![Field::new("a", DataType::Float32, false)]));
275300

301+
let column = datafusion::physical_plan::expressions::col("a", &schema)?;
276302
let mut eqp = EquivalenceProperties::new(Arc::clone(&schema));
277-
let _ = eqp.reorder([PhysicalSortExpr::new_default(
278-
datafusion::physical_plan::expressions::col("a", &schema)?,
279-
)]);
303+
let _ = eqp.reorder([PhysicalSortExpr::new_default(Arc::clone(&column))]);
280304
Ok(PlanProperties::new(
281305
eqp,
282306
Partitioning::RoundRobinBatch(3),
283307
EmissionType::Incremental,
284308
Boundedness::Bounded,
285-
))
309+
)
310+
.with_group_contiguous_exprs(vec![column]))
286311
}
287312

288313
fn create_range_test_props() -> Result<PlanProperties> {

‎datafusion/ffi/src/session/mod.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -976,7 +976,7 @@ mod tests {
976976
let physical_plan = foreign_session.create_physical_plan(&logical_plan).await?;
977977
assert_eq!(
978978
format!("{physical_plan:?}"),
979-
"EmptyExec { schema: Schema { fields: [], metadata: {} }, partitions: 1, cache: PlanProperties { eq_properties: EquivalenceProperties { eq_group: EquivalenceGroup { map: {}, classes: [] }, oeq_class: OrderingEquivalenceClass { orderings: [] }, oeq_cache: OrderingEquivalenceCache { normal_cls: OrderingEquivalenceClass { orderings: [] }, leading_map: {} }, constraints: Constraints { inner: [] }, schema: Schema { fields: [], metadata: {} } }, partitioning: UnknownPartitioning(1), emission_type: Incremental, boundedness: Bounded, evaluation_type: Lazy, scheduling_type: Cooperative, output_ordering: None } }"
979+
"EmptyExec { schema: Schema { fields: [], metadata: {} }, partitions: 1, cache: PlanProperties { eq_properties: EquivalenceProperties { eq_group: EquivalenceGroup { map: {}, classes: [] }, oeq_class: OrderingEquivalenceClass { orderings: [] }, oeq_cache: OrderingEquivalenceCache { normal_cls: OrderingEquivalenceClass { orderings: [] }, leading_map: {} }, constraints: Constraints { inner: [] }, schema: Schema { fields: [], metadata: {} } }, partitioning: UnknownPartitioning(1), emission_type: Incremental, boundedness: Bounded, evaluation_type: Lazy, scheduling_type: Cooperative, output_ordering: None, group_contiguous_exprs: [] } }"
980980
);
981981

982982
assert_eq!(

‎datafusion/ffi/src/tests/mod.rs‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -178,8 +178,9 @@ extern "C" fn construct_table_provider_factory(
178178

179179
pub(crate) extern "C" fn create_empty_exec() -> FFI_ExecutionPlan {
180180
let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Float32, false)]));
181-
182-
let plan = Arc::new(EmptyExec::new(schema));
181+
let expression = datafusion_physical_expr::expressions::col("a", &schema).unwrap();
182+
let plan =
183+
Arc::new(EmptyExec::new(schema).with_group_contiguous_exprs(vec![expression]));
183184
FFI_ExecutionPlan::new(plan, None)
184185
}
185186

‎datafusion/ffi/tests/ffi_execution_plan.rs‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -109,6 +109,21 @@ mod tests {
109109
Ok(())
110110
}
111111

112+
#[test]
113+
fn test_ffi_execution_plan_group_contiguous_exprs_cross_library()
114+
-> Result<(), DataFusionError> {
115+
let module = get_module()?;
116+
let plan = (module.create_empty_exec)();
117+
let plan: Arc<dyn ExecutionPlan> = (&plan).try_into()?;
118+
assert!(plan.is::<ForeignExecutionPlan>());
119+
120+
let group_contiguous_exprs = plan.group_contiguous_exprs();
121+
assert_eq!(group_contiguous_exprs.len(), 1);
122+
assert_eq!(group_contiguous_exprs[0].to_string(), "a@0");
123+
124+
Ok(())
125+
}
126+
112127
#[test]
113128
fn test_ffi_execution_plan_new_sets_runtimes_on_children()
114129
-> Result<(), DataFusionError> {

‎datafusion/physical-optimizer/src/ensure_coop.rs‎

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -152,6 +152,42 @@ mod tests {
152152
");
153153
}
154154

155+
#[test]
156+
fn test_cooperative_exec_preserves_partition_disjoint_aggregation() -> Result<()> {
157+
use arrow::datatypes::{DataType, Field, Schema};
158+
use datafusion_physical_expr::expressions::col;
159+
use datafusion_physical_plan::ExecutionPlanProperties;
160+
use datafusion_physical_plan::aggregates::{
161+
AggregateExec, AggregateMode, PhysicalGroupBy,
162+
};
163+
use datafusion_physical_plan::execution_plan::EmissionType;
164+
use datafusion_physical_plan::test::TestMemoryExec;
165+
166+
let schema =
167+
Arc::new(Schema::new(vec![Field::new("key", DataType::Int32, false)]));
168+
let input = TestMemoryExec::try_new(&[vec![]], Arc::clone(&schema), None)?
169+
.try_with_group_contiguous_keys(vec![col("key", &schema)?])?;
170+
let aggregate = Arc::new(AggregateExec::try_new(
171+
AggregateMode::Partial,
172+
PhysicalGroupBy::new_single(vec![(col("key", &schema)?, "key".to_string())]),
173+
vec![],
174+
vec![],
175+
Arc::new(input),
176+
schema,
177+
)?);
178+
179+
let optimized =
180+
EnsureCooperative::new().optimize(aggregate, &ConfigOptions::new())?;
181+
assert_eq!(optimized.pipeline_behavior(), EmissionType::Incremental);
182+
assert_snapshot!(displayable(optimized.as_ref()).indent(true).to_string(), @r"
183+
AggregateExec: mode=Partial, gby=[key@0 as key], aggr=[], group_completion_mode=Full
184+
CooperativeExec
185+
DataSourceExec: partitions=1, partition_sizes=[0], group_contiguous=[key@0]
186+
");
187+
188+
Ok(())
189+
}
190+
155191
#[tokio::test]
156192
async fn test_optimizer_is_idempotent() {
157193
// Comprehensive idempotency test: verify f(f(...f(x))) = f(x)

0 commit comments

Comments
 (0)