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]

Reply via email to