zhuqi-lucas commented on code in PR #23599:
URL: https://github.com/apache/datafusion/pull/23599#discussion_r4195962249
##########
datafusion/core/tests/physical_optimizer/window_topn.rs:
##########
@@ -724,6 +724,153 @@ async fn partitioned_topk_exec_exposes_metrics() ->
Result<()> {
.expect("output_batches metric")
.as_usize();
assert_eq!(output_batches, 5);
+ Ok(())
+}
+
+/// Regression: FilterExec that carries an embedded projection (e.g. from
+/// an earlier filter/projection pushdown pass) used to make the rule bail
+/// out entirely. The rule now captures the projection, applies the
+/// PartitionedTopKExec rewrite, and wraps the result in a ProjectionExec
+/// that reproduces the original output schema.
+#[test]
+fn filter_with_projection_still_rewrites() -> Result<()> {
+ let s = schema();
+ let input: Arc<dyn ExecutionPlan> =
Arc::new(PlaceholderRowExec::new(Arc::clone(&s)));
+
+ let ordering = LexOrdering::new(vec![
+ PhysicalSortExpr::new_default(col("pk", &s)?).asc(),
+ PhysicalSortExpr::new_default(col("val", &s)?).asc(),
+ ])
+ .unwrap();
+ let sort: Arc<dyn ExecutionPlan> =
+ Arc::new(SortExec::new(ordering,
input).with_preserve_partitioning(true));
+ let partition_by = vec![col("pk", &s)?];
+ let order_by = vec![PhysicalSortExpr::new_default(col("val", &s)?).asc()];
+ let window_expr = Arc::new(StandardWindowExpr::new(
+ create_udwf_window_expr(
+ &row_number_udwf(),
+ &[],
+ &s,
+ "row_number".to_string(),
+ false,
+ )?,
+ &partition_by,
+ &order_by,
+ Arc::new(WindowFrame::new_bounds(
+ WindowFrameUnits::Rows,
+ WindowFrameBound::Preceding(ScalarValue::UInt64(None)),
+ WindowFrameBound::CurrentRow,
+ )),
+ ));
+ let window: Arc<dyn ExecutionPlan> =
Arc::new(BoundedWindowAggExec::try_new(
+ vec![window_expr],
+ sort,
+ InputOrderMode::Sorted,
+ true,
+ )?);
+
+ // Filter: row_number@2 <= 3, with an embedded projection that keeps
+ // only [pk, val] (drops the row_number column) — the shape produced
+ // by filter/projection pushdown when downstream doesn't need the
+ // window column.
+ let rn_col = Arc::new(Column::new("row_number", 2));
+ let limit_lit = lit(ScalarValue::UInt64(Some(3)));
+ let predicate = Arc::new(BinaryExpr::new(rn_col, Operator::LtEq,
limit_lit));
+ let filter: Arc<dyn ExecutionPlan> = Arc::new(
+ FilterExecBuilder::new(predicate, window)
+ .apply_projection(Some(vec![0, 1]))?
Review Comment:
Thanks @kosiew — added in 7c3bc01 as
`reordered_duplicated_projection_preserves_filter_schema`: a `[1, 0, 1]`
projection (reorder + duplicate) over a schema with both field-level and
schema-level metadata, asserting `optimized.schema() == filter.schema()`.
It passes — the metadata does survive. Two guards so it can't pass for the
wrong reason: it checks the rewrite actually fired (rather than the rule
declining and handing back the input), and that `FilterExec` itself carries the
metadata being compared (otherwise both sides would be empty).
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]