Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions datafusion/physical-plan/src/aggregates/aggregate_stream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -109,9 +109,9 @@ impl AggregateStreamInner {
};

let mut predicates: Vec<Arc<dyn PhysicalExpr>> =
Vec::with_capacity(filter_state.supported_accumulators_info.len());
Vec::with_capacity(filter_state.accumulator_dyn_filter_info.len());

for acc_info in &filter_state.supported_accumulators_info {
for acc_info in &filter_state.accumulator_dyn_filter_info {
// Skip if we don't yet have a meaningful bound
let bound = {
let guard = acc_info.shared_bound.lock();
Expand Down Expand Up @@ -171,7 +171,7 @@ impl AggregateStreamInner {

let mut bounds_changed = false;

for acc_info in &filter_state.supported_accumulators_info {
for acc_info in &filter_state.accumulator_dyn_filter_info {
let acc =
self.accumulators
.get_mut(acc_info.aggr_index)
Expand Down
40 changes: 23 additions & 17 deletions datafusion/physical-plan/src/aggregates/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -755,10 +755,11 @@ impl From<StreamType> for SendableRecordBatchStream {
///
/// ## Enable Condition
/// - No grouping (no `GROUP BY` clause in the sql, only a single global group to aggregate)
/// - The aggregate expression must be `min`/`max`, and evaluate directly on columns.
/// Note multiple aggregate expressions that satisfy this requirement are allowed,
/// and a dynamic filter will be constructed combining all applicable expr's
/// states. See more in the following example with dynamic filter on multiple columns.
/// - Every aggregate expression must be `min`/`max`, and evaluate directly on a
/// column. If any aggregate expression is unsupported, dynamic filtering is
/// disabled for the entire [`AggregateExec`]. Multiple supported aggregate
/// expressions are combined into one dynamic filter. See the following example
/// with a dynamic filter on multiple columns.
///
/// ## Filter Construction
/// The filter is kept in the `DataSourceExec`, and it will gets update during execution,
Expand All @@ -778,11 +779,11 @@ struct AggrDynFilter {
/// The current bounds for the dynamic filter, updates during the execution to
/// tighten the bound for more effective pruning.
///
/// Each vector element is for the accumulators that support dynamic filter.
/// e.g. This `AggregateExec` has accumulator:
/// min(a), avg(a), max(b)
/// And this field stores [PerAccumulatorDynFilter(min(a)), PerAccumulatorDynFilter(min(b))]
supported_accumulators_info: Vec<PerAccumulatorDynFilter>,
/// Each vector element corresponds to one aggregate expression. Dynamic filtering
/// is enabled only when every aggregate expression is supported, so this vector
/// contains an entry for every accumulator. For example, `min(a), max(b)` produces
/// entries for `min(a)` and `max(b)`.
accumulator_dyn_filter_info: Vec<PerAccumulatorDynFilter>,
}

// ---- Aggregate Dynamic Filter Utility Structs ----
Expand Down Expand Up @@ -1157,7 +1158,7 @@ impl AggregateExec {
};

// Validate that the filter is compatible with the aggregation columns.
let cols = self.cols_for_dynamic_filter(&dyn_filter.supported_accumulators_info);
let cols = self.cols_for_dynamic_filter(&dyn_filter.accumulator_dyn_filter_info);
if cols.len() != filter.children().len() {
return internal_err!(
"Dynamic filter expression is incompatible with aggregate due to mismatched number of columns"
Expand All @@ -1174,7 +1175,7 @@ impl AggregateExec {
// Overwrite our filter
self.dynamic_filter = Some(Arc::new(AggrDynFilter {
filter,
supported_accumulators_info: dyn_filter.supported_accumulators_info.clone(),
accumulator_dyn_filter_info: dyn_filter.accumulator_dyn_filter_info.clone(),
}));
Ok(self)
}
Expand Down Expand Up @@ -1821,7 +1822,7 @@ impl AggregateExec {
return;
}

// Collect supported accumulators
// Collect dynamic filter metadata for every accumulator
// It is assumed the order of aggregate expressions are not changed from `AggregateExec`
// to `AggregateStream`
let mut aggr_dyn_filters = Vec::new();
Expand Down Expand Up @@ -1852,23 +1853,28 @@ impl AggregateExec {
aggr_index: i,
shared_bound: Arc::new(Mutex::new(ScalarValue::Null)),
});
} else {
// An incomplete filter could prune rows that still improve an
// unsupported aggregate, so every aggregate must be represented.
// TODO: Derive safe predicates for expressions such as `min(col + literal)`.
return;
}
}

if !aggr_dyn_filters.is_empty() {
self.dynamic_filter = Some(Arc::new(AggrDynFilter {
filter: Arc::new(DynamicFilterPhysicalExpr::new(all_cols, lit(true))),
supported_accumulators_info: aggr_dyn_filters,
accumulator_dyn_filter_info: aggr_dyn_filters,
}))
}
}

// Collect column references for the dynamic filter expression from the supported accumulators.
// Collect column references for the dynamic filter expression from the accumulators.
fn cols_for_dynamic_filter(
&self,
supported_accumulators_info: &[PerAccumulatorDynFilter],
accumulator_dyn_filter_info: &[PerAccumulatorDynFilter],
) -> Vec<Arc<dyn PhysicalExpr>> {
let all_cols: Vec<Arc<dyn PhysicalExpr>> = supported_accumulators_info
let all_cols: Vec<Arc<dyn PhysicalExpr>> = accumulator_dyn_filter_info
.iter()
.filter_map(|info| {
// This should always be true due to how the supported accumulators
Expand All @@ -1881,7 +1887,7 @@ impl AggregateExec {
None
})
.collect();
debug_assert_eq!(all_cols.len(), supported_accumulators_info.len());
debug_assert_eq!(all_cols.len(), accumulator_dyn_filter_info.len());
all_cols
}

Expand Down
35 changes: 25 additions & 10 deletions datafusion/sqllogictest/test_files/push_down_filter_regression.slt
Original file line number Diff line number Diff line change
Expand Up @@ -393,41 +393,56 @@ statement ok
drop table agg_dyn_two_col;

# --- mixed expressions: MIN(a), MAX(a), MAX(b), MIN(c+1) ---
# Supported aggregates (MIN(a), MAX(a), MAX(b)) should drive a filter;
# MIN(c+1) is unsupported and must not contribute.
# Every file shares the same per-file min(a)=1, max(a)=8 and max(b)=12 so the
# DynamicFilter content is deterministic regardless of publish order (see #22621).
# MIN(c+1) cannot contribute a dynamic-filter predicate. Ignoring it could prune
# rows that still improve MIN(c+1), so the aggregate must not produce a filter.
# Each file starts with a row group that establishes all supported bounds,
# followed by a row group containing the minimum c. With two rows per batch,
# the old filter prunes the required rows in both files regardless of which
# partition runs first.

statement ok
set datafusion.execution.batch_size = 2;

statement ok
COPY (
SELECT * FROM (VALUES (1, 12, 100), (8, 4, 70)) AS v(a, b, c)
SELECT * FROM (VALUES (1, 12, 100), (8, 4, 100), (1, 6, 70), (8, 12, 110)) AS v(a, b, c)
) TO 'test_files/scratch/push_down_filter_regression/agg_dyn_mixed/file_0.parquet'
STORED AS PARQUET;
STORED AS PARQUET
OPTIONS ('format.max_row_group_size' '2');

statement ok
COPY (
SELECT * FROM (VALUES (1, 6, 90), (8, 12, 110)) AS v(a, b, c)
SELECT * FROM (VALUES (1, 12, 100), (8, 4, 100), (1, 6, 70), (8, 12, 110)) AS v(a, b, c)
) TO 'test_files/scratch/push_down_filter_regression/agg_dyn_mixed/file_1.parquet'
STORED AS PARQUET;
STORED AS PARQUET
OPTIONS ('format.max_row_group_size' '2');

statement ok
CREATE EXTERNAL TABLE agg_dyn_mixed (a INT, b INT, c INT)
STORED AS PARQUET
LOCATION 'test_files/scratch/push_down_filter_regression/agg_dyn_mixed/';

# -> DynamicFilter [ a < 1 OR a > 8 OR b > 12 ] (MIN(c+1) dropped as unsupported)
query IIII
SELECT MIN(a), MAX(a), MAX(b), MIN(c + 1) FROM agg_dyn_mixed;
----
1 8 12 71

# No dynamic filter because not every aggregate has a safe predicate.
query TT
EXPLAIN ANALYZE SELECT MIN(a), MAX(a), MAX(b), MIN(c + 1) FROM agg_dyn_mixed;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we add an execution-result assertion here as well? The current EXPLAIN ANALYZE assertion confirms that the dynamic filter is absent, but it does not directly verify the correctness issue this change is intended to fix: the unsupported aggregate still needs to see every relevant row.

For this fixture, the expected result for MIN(a), MAX(a), MAX(b), MIN(c + 1) is 1, 8, 12, 71.

It would also be good to make the fixture or scan order deterministic enough that the old behavior actually prunes the row needed by MIN(c + 1). With the current two-file layout, that can depend on which file publishes its bounds first. Ideally, this regression test should fail with the old implementation because it produces the wrong result, rather than only because the plan contains a dynamic filter.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. I've added a result check and adjusted the test data so the old implementation produces the wrong result. The test passes with the fix.

----
Plan with Metrics
01)AggregateExec: mode=Final, gby=[], aggr=[min(agg_dyn_mixed.a), max(agg_dyn_mixed.a), max(agg_dyn_mixed.b), min(agg_dyn_mixed.c + Int64(1))], metrics=[]
02)--CoalescePartitionsExec, metrics=[]
03)----AggregateExec: mode=Partial, gby=[], aggr=[min(agg_dyn_mixed.a), max(agg_dyn_mixed.a), max(agg_dyn_mixed.b), min(agg_dyn_mixed.c + Int64(1))], metrics=[]
04)------DataSourceExec: file_groups={2 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_regression/agg_dyn_mixed/file_0.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_regression/agg_dyn_mixed/file_1.parquet]]}, projection=[a, b, c], file_type=parquet, predicate=DynamicFilter [ a@0 < 1 OR a@0 > 8 OR b@1 > 12 ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_min@0 < 1 OR a_null_count@1 != row_count@2 AND a_max@3 > 8 OR b_null_count@5 != row_count@2 AND b_max@4 > 12, required_guarantees=[], metrics=[]
04)------DataSourceExec: file_groups={2 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_regression/agg_dyn_mixed/file_0.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_regression/agg_dyn_mixed/file_1.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[]

statement ok
drop table agg_dyn_mixed;

statement ok
reset datafusion.execution.batch_size;

# --- all-NULLs input: filter should stay `true` (no meaningful bound) ---

statement ok
Expand Down