sunchao commented on code in PR #6433:
URL: https://github.com/apache/datafusion-comet/pull/6433#discussion_r4189184697
##########
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") {
Review Comment:
Committed the three-join regression and verified both intermediate build
orientations with the flag off/on, two active ancestor consumers, and retained
intermediate reader attachment. The SQL file also runs false/true matrices with
duplicate and NULL keys.
On current head ff72f70d47, SF1 TPC-DS q7, q19, and q42 each matched Spark
with the flag off and on (6/6 checks, Spark 4.1.3, AQE off). Enabled metrics
summed across joins:
| Query | Early evaluated | Early pruned | Early bypassed | Reader
attachments |
| --- | ---: | ---: | ---: | ---: |
| q7 | 3,188,398 | 2,123,647 | 528,869 | 5 |
| q19 | 2,654,022 | 2,605,984 | 94,270 | 5 |
| q42 | 2,750,838 | 2,696,363 | 0 | 5 |
These are selected-query result/counter checks, not the full TPC-DS suite or
timing measurements. Detailed scope and JNI provenance are in the updated PR
description.
##########
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:
DynamicFilterJoinExec now forwards its HashJoinExec template fetch. A native
regression creates an intermediate join with fetch=1 and verifies that early
placement leaves that boundary intact.
##########
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:
Removed the ignored component benchmark and its duplicated join constructor.
The committed Scala/SQL regressions exercise the real planner and Parquet path;
current TPC-DS result/metric evidence is in the PR description.
##########
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:
Added `dynamic_filter_chain.sql` with the false/true ConfigMatrix, both
broadcast build sides, two- and three-join chains, duplicates, and NULL keys.
Scala retains metric assertions and no longer repeats the query for hard-coded
result rows.
##########
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:
Derived Clone and used struct updates for all three rebuild paths.
Execution-local replacement preserves metric identity, while ordinary child
replacement/reset retain fresh metrics and reset clears the old predicate. A
regression pins adaptive mode across all three paths.
##########
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:
The helper now constructs the filter/projection adapter once and evaluates
children() once before recursive preparation. It uses replace_children with
explicit recomputation because DataFusion deprecates with_new_children.
##########
native/operators/src/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:
Added projection-specific delegation tests for expression-root
visitation/Stop recursion and statistics at whole-plan and individual-partition
scope. Both tests pass; the adapter is now beside filter.rs.
--
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]