comphead commented on code in PR #6433:
URL: https://github.com/apache/datafusion-comet/pull/6433#discussion_r4185464748
##########
native/core/src/execution/planner.rs:
##########
@@ -2540,14 +2540,36 @@ impl PhysicalPlanner {
}
}
- /// Keep the Spark filter's metric identity when its reader is replaced
for an execution.
- fn prepare_probe_filter_for_runtime_reader(plan: Arc<SparkPlan>) ->
Arc<SparkPlan> {
- let Some(filter) = plan.native_plan.downcast_ref::<FilterExec>() else {
- return plan;
- };
+ /// Keep Spark metric identities for the small set of nodes whose children
+ /// runtime-filter placement can replace. Stop at native/Spark tree
boundaries.
+ fn prepare_probe_filter_for_runtime_reader(
+ plan: Arc<SparkPlan>,
+ ) -> Result<Arc<SparkPlan>, ExecutionError> {
+ let native = &plan.native_plan;
+ if !native.is::<FilterExec>() && !native.is::<ProjectionExec>() {
+ return Ok(plan);
+ }
let mut prepared = plan.as_ref().clone();
- prepared.native_plan =
Arc::new(CometFilterExec::from_datafusion(filter.clone()));
- Arc::new(prepared)
+ let mut native = Arc::clone(native);
+ if let [child] = plan.children.as_slice() {
+ if native.children().len() == 1 &&
Arc::ptr_eq(native.children()[0], &child.native_plan)
Review Comment:
`native.children()` builds a `Vec` on each call and runs twice in this
condition. The node kind is also checked with `is` at the top and again with
`downcast_ref` after `replace_children`.
Would it be simpler to build the `CometFilterExec` or `CometProjectionExec`
once up front and call `with_new_children` on it when the child is replaced?
Both wrappers implement it, and a single `if let [only] =
native.children().as_slice()` would evaluate `children()` once. I have not
compiled this.
##########
native/core/src/execution/operators/dynamic_filter/mod.rs:
##########
@@ -122,12 +140,14 @@ impl ExecutionPlan for DynamicFilterExec {
if children.len() != 1 {
return internal_err!("CometDynamicFilterExec requires one child");
}
- Ok(Arc::new(Self::new(
+ let mut replaced = Self::new(
children.remove(0),
Arc::clone(&self.predicate),
ExecutionPlanMetricsSet::new(),
self.metric_prefix,
- )))
+ );
+ replaced.adaptive = self.adaptive;
+ Ok(Arc::new(replaced))
Review Comment:
`adaptive` has to be copied by hand here and in `reset_state`, because
`Self::new` starts it at `false`. That is easy to lose when the struct gains
another field. #6431 hits the same hazard on `CometFilterExec`. It patches four
constructors and adds `runtime_filter_permission_survives_rebuilding` to pin
them.
Would `#[derive(Clone)]` plus `Self { input, ..self.clone() }` in
`with_execution_input`, `replace_children` and `reset_state` be simpler? Either
way, a small test that `adaptive` survives all three rebuild paths would stop a
later refactor from turning the early consumer into a plain one without any
failure.
##########
native/core/src/execution/operators/dynamic_filter/join/tests/early_benchmark.rs:
##########
@@ -0,0 +1,147 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Reproducible component benchmark for selective and all-matching join
chains.
+//!
+//! Run this ignored test with an optimized build and identical settings on
both
+//! revisions. It covers decoded native batches, not Spark or Parquet I/O.
+
+use super::*;
+use std::time::Instant;
+
+fn benchmark_input(rows: usize, keys: usize, batch_rows: usize) -> Arc<dyn
ExecutionPlan> {
+ let schema = Arc::new(Schema::new(vec![
+ Field::new("key", DataType::Int32, false),
+ Field::new("payload", DataType::Int32, false),
+ ]));
+ let batches = (0..rows)
+ .step_by(batch_rows)
+ .map(|offset| {
+ let end = (offset + batch_rows).min(rows);
+ RecordBatch::try_new(
+ Arc::clone(&schema),
+ vec![
+ Arc::new(Int32Array::from_iter_values(
+ (offset..end).map(|i| (i % keys) as i32),
+ )),
+ Arc::new(Int32Array::from_iter_values(
+ (offset..end).map(|i| i as i32),
+ )),
+ ],
+ )
+ .unwrap()
+ })
+ .collect::<Vec<_>>();
+ memory_exec(batches)
+}
+
+fn benchmark_join(
+ build: Arc<dyn ExecutionPlan>,
+ probe: Arc<dyn ExecutionPlan>,
+ probe_key: usize,
+ config: &ConfigOptions,
+) -> Arc<dyn ExecutionPlan> {
+ let join = HashJoinExec::try_new(
+ build,
+ probe,
+ vec![(
+ Arc::new(Column::new("key", 0)),
+ Arc::new(Column::new("key", probe_key)),
+ )],
+ None,
+ &JoinType::Inner,
+ None,
+ PartitionMode::Partitioned,
+ NullEquality::NullEqualsNothing,
+ false,
+ )
+ .unwrap();
+ PhysicalPlanner::apply_join_dynamic_filter(Arc::new(join), true,
config).unwrap()
+}
Review Comment:
`benchmark_join` repeats the `HashJoinExec::try_new` argument list from
`chain_join` in `tests.rs`. If this file stays, could it call
`chain_join(build, probe, 0, probe_key)` and keep only the
`apply_join_dynamic_filter` call?
On moving it to criterion, as suggested on the `#[ignore]` test below:
`benches/` targets only see the public API of the `comet` crate.
`execution::planner` and `DynamicFilterJoinExec` are `pub(crate)`, so that move
would also need them exposed. `CometTopKBenchmark` is a closer precedent for
this feature. It toggles the sibling TopK filter over the real planner and a
Parquet scan, so it would also cover the decode cost that this component
benchmark leaves out.
##########
native/core/src/execution/operators/runtime_filter_projection.rs:
##########
@@ -0,0 +1,211 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! A DataFusion projection whose metrics remain owned by its Spark plan node.
+//!
+//! Comet normally keeps a one-to-one Spark/native plan tree so native metric
+//! handles map back to the corresponding Spark operator. Some execution-local
+//! rewrites need to replace a projection's child. DataFusion gives that
replacement
+//! a new private metric set, so this adapter owns the stable metric set and
+//! registers the handles from the projection that actually executes.
+
+use std::fmt::Formatter;
+use std::sync::Arc;
+
+use datafusion::common::tree_node::TreeNodeRecursion;
+use datafusion::common::{internal_err, Result, Statistics};
+use datafusion::execution::TaskContext;
+use datafusion::physical_expr::PhysicalExpr;
+use datafusion::physical_plan::execution_plan::{
+ CardinalityEffect, ChildrenPropertiesMode, ReplaceChildrenOptions,
+};
+use datafusion::physical_plan::metrics::{ExecutionPlanMetricsSet, MetricsSet};
+use datafusion::physical_plan::projection::ProjectionExec;
+use datafusion::physical_plan::statistics::{ChildStats, StatisticsArgs};
+use datafusion::physical_plan::{
+ DisplayAs, DisplayFormatType, ExecutionPlan, PlanProperties,
SendableRecordBatchStream,
+};
+
+#[derive(Debug)]
+pub(crate) struct CometProjectionExec {
+ projection: ProjectionExec,
+ metrics: ExecutionPlanMetricsSet,
+}
Review Comment:
Following up on the thread at the top of this file. `filter.rs` has two unit
tests, `expression_visitor_preserves_predicate_root_and_recursion` and
`statistics_match_datafusion_for_each_partition`, that pin the delegations this
copy repeats. This file has no tests. If the copy stays, could it get the same
two?
Otherwise, would one generic wrapper over `FilterExec` and `ProjectionExec`
be simpler than two copies and two test sets? I have not compiled that idea, so
there may be a catch.
##########
native/core/src/execution/operators/dynamic_filter/early.rs:
##########
@@ -0,0 +1,135 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Place a live ancestor filter before an intermediate join's probe work.
+//!
+//! This is decoded-batch filtering only. It leaves scan/schema conversion and
+//! arbitrary expressions in place, and never propagates into an intermediate
+//! build side. The downstream join remains the authority for matching rows.
+
+use std::any::Any;
+use std::sync::Arc;
+
+use datafusion::common::{internal_err, Result};
+use datafusion::physical_expr::expressions::{Column,
DynamicFilterPhysicalExpr};
+use datafusion::physical_expr::PhysicalExpr;
+use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet;
+use datafusion::physical_plan::ExecutionPlan;
+use datafusion_comet_operators::CometFilterExec;
+
+use super::parquet_reader::is_direct_column_null_checks;
+use super::{DynamicFilterExec, DynamicFilterJoinExec};
+use crate::execution::operators::CometProjectionExec;
+
+pub(super) fn place_early_filter(
+ input: &Arc<dyn ExecutionPlan>,
+ predicate: Arc<DynamicFilterPhysicalExpr>,
+ metrics: &ExecutionPlanMetricsSet,
+) -> Result<Arc<dyn ExecutionPlan>> {
+ Ok(place(input, predicate, metrics, false)?.unwrap_or_else(||
Arc::clone(input)))
+}
+
+fn remap(
+ predicate: Arc<DynamicFilterPhysicalExpr>,
+ column: Arc<dyn PhysicalExpr>,
+) -> Result<Arc<DynamicFilterPhysicalExpr>> {
+ // Derived expressions share producer updates. Taking current() here would
+ // capture the initial TRUE placeholder instead of the completed build
domain.
+ let mapped: Arc<dyn Any + Send + Sync> =
predicate.with_new_children(vec![column])?;
+ mapped.downcast::<DynamicFilterPhysicalExpr>().map_err(|_| {
+ datafusion::common::DataFusionError::Internal(
+ "Dynamic filter remapping changed type".into(),
+ )
+ })
+}
+
+fn place(
+ input: &Arc<dyn ExecutionPlan>,
+ predicate: Arc<DynamicFilterPhysicalExpr>,
+ metrics: &ExecutionPlanMetricsSet,
+ crossed_join: bool,
+) -> Result<Option<Arc<dyn ExecutionPlan>>> {
+ let children = predicate.children();
+ let [key] = children.as_slice() else {
+ return internal_err!("Early join filtering requires one key");
+ };
+ let Some(key) = key.downcast_ref::<Column>() else {
+ return internal_err!("Early join filtering requires a column key");
+ };
+ if input.fetch().is_none() {
Review Comment:
`place()` relies on `fetch()` to stop at limits, and the join arm below is
new. `DynamicFilterJoinExec` does not override `fetch()`, so it reports `None`
even when its `HashJoinExec` template carries one (DF 55 has
`HashJoinExec::with_fetch`). Comet's planner does not set a fetch on a join
today, so I do not expect this to be reachable now.
`TopKReaderFilterExec` forwards `self.template.fetch()`. Could
`DynamicFilterJoinExec` do the same, or could this arm check
`template.fetch()`? That keeps the guard true if a join limit is ever pushed
down, since filtering below a join that has a fetch changes which rows it emits.
##########
spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala:
##########
@@ -431,6 +431,60 @@ class CometJoinSuite extends CometTestBase {
}
}
+ test("join dynamic filter rejects rows before an intermediate join") {
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1",
+ SQLConf.LEAF_NODE_DEFAULT_PARALLELISM.key -> "1",
+ CometConf.COMET_BATCH_SIZE.key -> "128") {
+ withParquetTable((0 until 1000).map(i => (i, i.toLong)), "early_fact") {
+ withParquetTable((0 until 2000).map(i => (i % 1000, i)),
"early_dimension") {
+ withParquetTable(Seq((42, 1), (42, 2)), "early_selection") {
+ for (buildLeft <- Seq(false, true); enabled <- Seq(false, true)) {
+ val from = if (buildLeft) {
+ "early_dimension d JOIN early_fact f"
+ } else {
+ "early_fact f JOIN early_dimension d"
+ }
+ val query = "SELECT /*+ BROADCAST(d), BROADCAST(s) */ " +
+ "f._1 AS selected_key, f._2 AS payload, d._2 AS detail, s._2
AS selection " +
+ s"FROM $from ON f._1 = d._1 " +
+ "JOIN early_selection s ON f._1 = s._1"
+ withSQLConf(
+ CometConf.COMET_EXEC_JOIN_DYNAMIC_FILTER_ENABLED.key ->
enabled.toString) {
+ val (_, plan) = checkSparkAnswerAndOperator(
+ sql(query),
+ Seq(classOf[CometBroadcastHashJoinExec]))
+ checkAnswer(
+ sql(query),
+ Seq(
+ Row(42, 42L, 42, 1),
+ Row(42, 42L, 42, 2),
+ Row(42, 42L, 1042, 1),
+ Row(42, 42L, 1042, 2)))
Review Comment:
The metric assertions have to stay in Scala, but the result checks could
also run as a SQL file test. That would be a cheap way to get the flag-on
coverage asked for in the thread at the top of this test, since
`CometSqlFileTestSuite` already runs in CI.
`sql-tests/join/sort_merge_join_binary.sql` shows `-- Config:` and `--
ConfigMatrix:` in a join file. A file with `-- ConfigMatrix:
spark.comet.exec.join.dynamicFilter.enabled=false,true` could compare more
shapes with Spark: both build sides, a three-join chain, duplicate keys and
NULL keys.
A SQL file cannot assert that the early path ran, so this test would keep
the metric checks. It could then drop the hard-coded `checkAnswer` rows and the
second execution of `query`, because `checkSparkAnswerAndOperator` already
compares with Spark.
--
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]