Skip to content

Commit 1e7d83a

Browse files
committed
feat: add narrow group-contiguous source property
1 parent 372429f commit 1e7d83a

5 files changed

Lines changed: 307 additions & 6 deletions

File tree

‎datafusion/datasource/src/source.rs‎

Lines changed: 104 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -165,6 +165,19 @@ pub trait DataSource: Any + Send + Sync + Debug {
165165

166166
fn output_partitioning(&self) -> Partitioning;
167167
fn eq_properties(&self) -> EquivalenceProperties;
168+
169+
/// Expressions whose complete tuple is contiguous within each output
170+
/// partition.
171+
///
172+
/// See [`ExecutionPlan::group_contiguous_exprs`] for the full correctness
173+
/// contract. Expressions must refer to the schema returned by
174+
/// [`Self::eq_properties`]. A source rewrite such as [`Self::repartitioned`]
175+
/// or [`Self::try_swapping_with_projection`] must preserve, remap, or clear
176+
/// this assertion as appropriate for the rewritten output streams.
177+
fn group_contiguous_exprs(&self) -> &[Arc<dyn PhysicalExpr>] {
178+
&[]
179+
}
180+
168181
fn scheduling_type(&self) -> SchedulingType {
169182
SchedulingType::NonCooperative
170183
}
@@ -397,6 +410,10 @@ impl ExecutionPlan for DataSourceExec {
397410
&self.cache
398411
}
399412

413+
fn group_contiguous_exprs(&self) -> &[Arc<dyn PhysicalExpr>] {
414+
self.data_source.group_contiguous_exprs()
415+
}
416+
400417
fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
401418
Vec::new()
402419
}
@@ -705,3 +722,90 @@ where
705722
Self::new(Arc::new(source))
706723
}
707724
}
725+
726+
#[cfg(test)]
727+
mod tests {
728+
use super::*;
729+
730+
use arrow::datatypes::{DataType, Field, Schema, SchemaRef};
731+
use datafusion_physical_expr::expressions::Column;
732+
use datafusion_physical_plan::EmptyRecordBatchStream;
733+
734+
#[derive(Debug)]
735+
struct GroupContiguousSource {
736+
schema: SchemaRef,
737+
group_contiguous_exprs: Vec<Arc<dyn PhysicalExpr>>,
738+
}
739+
740+
impl DataSource for GroupContiguousSource {
741+
fn open(
742+
&self,
743+
_partition: usize,
744+
_context: Arc<TaskContext>,
745+
) -> Result<SendableRecordBatchStream> {
746+
Ok(Box::pin(EmptyRecordBatchStream::new(Arc::clone(
747+
&self.schema,
748+
))))
749+
}
750+
751+
fn fmt_as(&self, _t: DisplayFormatType, f: &mut Formatter) -> fmt::Result {
752+
write!(f, "GroupContiguousSource")
753+
}
754+
755+
fn output_partitioning(&self) -> Partitioning {
756+
Partitioning::UnknownPartitioning(1)
757+
}
758+
759+
fn eq_properties(&self) -> EquivalenceProperties {
760+
EquivalenceProperties::new(Arc::clone(&self.schema))
761+
}
762+
763+
fn group_contiguous_exprs(&self) -> &[Arc<dyn PhysicalExpr>] {
764+
&self.group_contiguous_exprs
765+
}
766+
767+
fn partition_statistics(
768+
&self,
769+
_partition: Option<usize>,
770+
) -> Result<Arc<Statistics>> {
771+
Ok(Arc::new(Statistics::new_unknown(&self.schema)))
772+
}
773+
774+
fn with_fetch(&self, _limit: Option<usize>) -> Option<Arc<dyn DataSource>> {
775+
None
776+
}
777+
778+
fn fetch(&self) -> Option<usize> {
779+
None
780+
}
781+
782+
fn try_swapping_with_projection(
783+
&self,
784+
_projection: &ProjectionExprs,
785+
) -> Result<Option<Arc<dyn DataSource>>> {
786+
Ok(None)
787+
}
788+
789+
fn apply_expressions(
790+
&self,
791+
_f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> Result<TreeNodeRecursion>,
792+
) -> Result<TreeNodeRecursion> {
793+
Ok(TreeNodeRecursion::Continue)
794+
}
795+
}
796+
797+
#[test]
798+
fn data_source_exec_exposes_group_contiguous_exprs() -> Result<()> {
799+
let schema =
800+
Arc::new(Schema::new(vec![Field::new("key", DataType::Int32, false)]));
801+
let key = Arc::new(Column::new("key", 0)) as Arc<dyn PhysicalExpr>;
802+
let exec = DataSourceExec::new(Arc::new(GroupContiguousSource {
803+
schema,
804+
group_contiguous_exprs: vec![Arc::clone(&key)],
805+
}));
806+
807+
assert_eq!(exec.group_contiguous_exprs().len(), 1);
808+
assert!(exec.group_contiguous_exprs()[0].eq(&key));
809+
Ok(())
810+
}
811+
}

‎datafusion/physical-plan/src/coop.rs‎

Lines changed: 20 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -261,6 +261,10 @@ impl ExecutionPlan for CooperativeExec {
261261
&self.properties
262262
}
263263

264+
fn group_contiguous_exprs(&self) -> &[Arc<dyn PhysicalExpr>] {
265+
self.input.group_contiguous_exprs()
266+
}
267+
264268
fn maintains_input_order(&self) -> Vec<bool> {
265269
vec![true; self.children().len()]
266270
}
@@ -477,7 +481,9 @@ pub fn make_cooperative(stream: SendableRecordBatchStream) -> SendableRecordBatc
477481
mod tests {
478482
use super::*;
479483

480-
use arrow_schema::SchemaRef;
484+
use crate::test::TestMemoryExec;
485+
use arrow_schema::{DataType, Field, Schema, SchemaRef};
486+
use datafusion_physical_expr::expressions::col;
481487

482488
use futures::stream;
483489

@@ -497,6 +503,19 @@ mod tests {
497503
Box::pin(RecordBatchStreamAdapter::new(schema, s))
498504
}
499505

506+
#[test]
507+
fn cooperative_exec_preserves_group_contiguity() -> Result<()> {
508+
let schema =
509+
Arc::new(Schema::new(vec![Field::new("key", DataType::Int32, false)]));
510+
let source = TestMemoryExec::try_new(&[vec![]], Arc::clone(&schema), None)?
511+
.try_with_group_contiguous_exprs(vec![col("key", &schema)?])?;
512+
let cooperative = CooperativeExec::new(Arc::new(source));
513+
514+
assert_eq!(cooperative.group_contiguous_exprs().len(), 1);
515+
assert!(cooperative.group_contiguous_exprs()[0].eq(&col("key", &schema)?));
516+
Ok(())
517+
}
518+
500519
#[tokio::test]
501520
async fn yield_less_than_threshold() -> Result<()> {
502521
let count = TASK_BUDGET - 10;

‎datafusion/physical-plan/src/execution_plan.rs‎

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -159,6 +159,32 @@ pub trait ExecutionPlan: Any + Debug + DisplayAs + Send + Sync {
159159
/// trait, which is implemented for all `ExecutionPlan`s.
160160
fn properties(&self) -> &Arc<PlanProperties>;
161161

162+
/// Expressions whose complete tuple is contiguous within each output
163+
/// partition.
164+
///
165+
/// For every distinct tuple of expression values, all rows with that tuple
166+
/// must occur in at most one contiguous range in each output stream. Once a
167+
/// stream produces a different tuple, the previous tuple must never occur
168+
/// again. Tuple equality follows `GROUP BY` semantics, and tuple values do
169+
/// not need to be sorted.
170+
///
171+
/// This property does not imply any particular output ordering or
172+
/// distribution. It is also deliberately fail-closed: the default is no
173+
/// guarantee, and operators must explicitly preserve it. Currently,
174+
/// [`ProjectionExec`] preserves the complete tuple when every expression
175+
/// can be projected, and [`crate::coop::CooperativeExec`] preserves it
176+
/// because it does not change rows. Other operators use the default empty
177+
/// value.
178+
///
179+
/// # Correctness
180+
///
181+
/// This is a correctness contract. An invalid declaration can cause a
182+
/// streaming aggregate to emit a group before all of its rows have been
183+
/// observed, producing incorrect results.
184+
fn group_contiguous_exprs(&self) -> &[Arc<dyn PhysicalExpr>] {
185+
&[]
186+
}
187+
162188
/// Returns an error if this individual node does not conform to its invariants.
163189
/// These invariants are typically only checked in debug mode.
164190
///

‎datafusion/physical-plan/src/projection.rs‎

Lines changed: 107 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -83,6 +83,8 @@ pub struct ProjectionExec {
8383
metrics: ExecutionPlanMetricsSet,
8484
/// Cache holding plan properties like equivalences, output partitioning etc.
8585
cache: Arc<PlanProperties>,
86+
/// Complete group-contiguous tuple projected from the input, if any.
87+
group_contiguous_exprs: Vec<Arc<dyn PhysicalExpr>>,
8688
}
8789

8890
impl ProjectionExec {
@@ -208,6 +210,8 @@ impl ProjectionExec {
208210
// Construct a map from the input expressions to the output expression of the Projection
209211
let projection_mapping =
210212
projector.projection().projection_mapping(&input.schema())?;
213+
let group_contiguous_exprs =
214+
Self::project_group_contiguous_exprs(&input, &projection_mapping);
211215
let cache = Self::compute_properties(
212216
&input,
213217
&projection_mapping,
@@ -219,9 +223,26 @@ impl ProjectionExec {
219223
input,
220224
metrics: ExecutionPlanMetricsSet::new(),
221225
cache: Arc::new(cache),
226+
group_contiguous_exprs,
222227
})
223228
}
224229

230+
/// Projects the complete group-contiguous tuple, dropping it when any
231+
/// component cannot be mapped to the output schema.
232+
fn project_group_contiguous_exprs(
233+
input: &Arc<dyn ExecutionPlan>,
234+
projection_mapping: &ProjectionMapping,
235+
) -> Vec<Arc<dyn PhysicalExpr>> {
236+
input
237+
.equivalence_properties()
238+
.project_expressions(
239+
input.group_contiguous_exprs().iter(),
240+
projection_mapping,
241+
)
242+
.collect::<Option<Vec<_>>>()
243+
.unwrap_or_default()
244+
}
245+
225246
/// The projection expressions stored as tuples of (expression, output column name)
226247
pub fn expr(&self) -> &[ProjectionExpr] {
227248
self.projector.projection().as_ref()
@@ -345,6 +366,10 @@ impl ExecutionPlan for ProjectionExec {
345366
&self.cache
346367
}
347368

369+
fn group_contiguous_exprs(&self) -> &[Arc<dyn PhysicalExpr>] {
370+
&self.group_contiguous_exprs
371+
}
372+
348373
fn maintains_input_order(&self) -> Vec<bool> {
349374
// Tell optimizer this operator doesn't reorder its input
350375
vec![true]
@@ -386,11 +411,21 @@ impl ExecutionPlan for ProjectionExec {
386411
) -> Result<Arc<dyn ExecutionPlan>> {
387412
validate_child_count!(self, children);
388413
match options.children_properties {
389-
ChildrenPropertiesMode::Keep => Ok(Arc::new(Self {
390-
input: children.swap_remove(0),
391-
metrics: ExecutionPlanMetricsSet::new(),
392-
..Self::clone(&*self)
393-
})),
414+
ChildrenPropertiesMode::Keep => {
415+
let input = children.swap_remove(0);
416+
let projection_mapping = self
417+
.projector
418+
.projection()
419+
.projection_mapping(&input.schema())?;
420+
let group_contiguous_exprs =
421+
Self::project_group_contiguous_exprs(&input, &projection_mapping);
422+
Ok(Arc::new(Self {
423+
input,
424+
metrics: ExecutionPlanMetricsSet::new(),
425+
group_contiguous_exprs,
426+
..Self::clone(&*self)
427+
}))
428+
}
394429
ChildrenPropertiesMode::Recompute => {
395430
// `Keep` above requires the child's properties to be unchanged
396431
// outright. A rule that introduces a sort below this projection
@@ -631,6 +666,8 @@ impl ExecutionPlan for ProjectionExec {
631666
metrics: _,
632667
// Derived plan properties, recomputed on decode.
633668
cache: _,
669+
// Derived from the input assertion and projection expressions.
670+
group_contiguous_exprs: _,
634671
} = self;
635672
let projection_exprs = projector.projection().as_ref();
636673
let input = ctx.encode_child(input)?;
@@ -1515,6 +1552,7 @@ mod tests {
15151552
use crate::filter_pushdown::PushedDown;
15161553
use crate::statistics::{StatisticsArgs, StatisticsContext};
15171554
use crate::test;
1555+
use crate::test::TestMemoryExec;
15181556
use crate::test::exec::StatisticsExec;
15191557

15201558
use arrow::datatypes::{DataType, Field, Schema};
@@ -1526,6 +1564,70 @@ mod tests {
15261564
BinaryExpr, Column, DynamicFilterPhysicalExpr, Literal, binary, col, lit,
15271565
};
15281566

1567+
#[test]
1568+
fn group_contiguous_projection_is_all_or_nothing() -> Result<()> {
1569+
let schema = Arc::new(Schema::new(vec![
1570+
Field::new("key", DataType::Int32, false),
1571+
Field::new("time", DataType::Int32, false),
1572+
]));
1573+
let time_bin = binary(
1574+
col("time", &schema)?,
1575+
Operator::Divide,
1576+
lit(ScalarValue::Int32(Some(10))),
1577+
&schema,
1578+
)?;
1579+
let plain_source = TestMemoryExec::try_new(&[vec![]], Arc::clone(&schema), None)?;
1580+
let source = plain_source.clone().try_with_group_contiguous_exprs(vec![
1581+
col("key", &schema)?,
1582+
Arc::clone(&time_bin),
1583+
])?;
1584+
1585+
let projection = ProjectionExec::try_new(
1586+
[
1587+
ProjectionExpr::new(col("key", &schema)?, "key"),
1588+
ProjectionExpr::new(time_bin, "time_bin"),
1589+
],
1590+
Arc::new(source.clone()),
1591+
)?;
1592+
let projected_schema = projection.schema();
1593+
let expected = [
1594+
col("key", &projected_schema)?,
1595+
col("time_bin", &projected_schema)?,
1596+
];
1597+
assert_eq!(projection.group_contiguous_exprs().len(), expected.len());
1598+
assert!(
1599+
projection
1600+
.group_contiguous_exprs()
1601+
.iter()
1602+
.zip(expected)
1603+
.all(|(actual, expected)| actual.eq(&expected))
1604+
);
1605+
1606+
// A strict subset is not sufficient: contiguity of `(key, time_bin)`
1607+
// does not imply that `key` alone is contiguous.
1608+
let partial_projection = ProjectionExec::try_new(
1609+
[ProjectionExpr::new(col("key", &schema)?, "key")],
1610+
Arc::new(source.clone()),
1611+
)?;
1612+
assert!(partial_projection.group_contiguous_exprs().is_empty());
1613+
1614+
// Row-preserving operators do not inherit the assertion unless they
1615+
// opt in explicitly.
1616+
let filter = FilterExec::try_new(lit(true), Arc::new(source.clone()))?;
1617+
assert!(filter.group_contiguous_exprs().is_empty());
1618+
1619+
// The assertion deliberately lives outside PlanProperties. Exercise
1620+
// the child-replacement fast path with identical cached properties and
1621+
// verify that ProjectionExec still recomputes it.
1622+
assert!(Arc::ptr_eq(source.properties(), plain_source.properties()));
1623+
let projection: Arc<dyn ExecutionPlan> = Arc::new(projection);
1624+
let replaced =
1625+
replace_children_if_necessary(projection, vec![Arc::new(plain_source)])?;
1626+
assert!(replaced.group_contiguous_exprs().is_empty());
1627+
1628+
Ok(())
1629+
}
1630+
15291631
#[test]
15301632
fn test_try_new_with_schema_metadata_only_replaces_metadata() -> Result<()> {
15311633
let input_schema = Arc::new(Schema::new(vec![Field::new(

0 commit comments

Comments
 (0)