sunchao commented on code in PR #6433:
URL: https://github.com/apache/datafusion-comet/pull/6433#discussion_r4189081789


##########
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() {
+        if let Some(join) = input.downcast_ref::<DynamicFilterJoinExec>() {
+            let template = join.template();
+            // This wrapper is already limited to ordinary inner, single-key 
joins.
+            // Do not skip fallible join residuals, or guess an embedded 
projection.
+            if template.filter().is_none() && !template.contains_projection() {
+                let build_columns = template.left().schema().fields().len();
+                if let Some(index) = key.index().checked_sub(build_columns) {
+                    if let Some(field) = 
template.right().schema().fields().get(index) {
+                        let mapped = remap(
+                            Arc::clone(&predicate),
+                            Arc::new(Column::new(field.name(), index)),
+                        )?;
+                        if let Some(probe) = place(template.right(), mapped, 
metrics, true)? {
+                            return Ok(Some(join.with_execution_probe(probe)?));
+                        }
+                    }
+                }
+            }
+        } else if let Some(projection) = 
input.downcast_ref::<CometProjectionExec>() {
+            let exprs = projection.projection().expr();
+            // Even an unrelated computed expression may fail or be stateful.
+            if exprs.iter().all(|expr| expr.expr.is::<Column>()) {
+                if let Some(expr) = exprs.get(key.index()) {
+                    let mapped = remap(Arc::clone(&predicate), 
Arc::clone(&expr.expr))?;
+                    if let Some(child) = place(projection.input(), mapped, 
metrics, crossed_join)? {
+                        return 
Ok(Some(projection.with_execution_input(child)?));
+                    }
+                }
+            }
+        } else if let Some(filter) = input.downcast_ref::<CometFilterExec>() {

Review Comment:
   Placed the terminal early consumer above inferred null checks after crossing 
a join. The nullable-FK regression has a null in every decoded batch and 
verifies that the ancestor evaluates the first two surviving batches, prunes 
zero rows, then bypasses the remaining three; results match the unfiltered join 
plan.



##########
native/core/src/execution/operators/dynamic_filter/mod.rs:
##########
@@ -172,8 +194,18 @@ impl ExecutionPlan for DynamicFilterExec {
         // add its input/output counts or elapsed time to the join's existing 
metrics.
         let eval_time = MetricBuilder::new(&self.metrics)
             .subset_time(format!("{}_eval_time", self.metric_prefix), 
partition);
+        // Early filtering duplicates the final consumer. Stop that extra work 
after
+        // two nonempty evaluated batches remove nothing. The downstream join 
still
+        // verifies every row, so later selectivity changes only lose an 
optimization.
+        // Keep this decision per stream; an inactive TRUE placeholder is not 
a sample.
+        let adaptive = self.adaptive;

Review Comment:
   Kept exact-zero adaptation and documented the tradeoff: even a small 
reduction may save substantial intermediate join work, so this avoids choosing 
a workload-specific ratio. Permanent bypass bounds duplicate checks; clustered 
inputs can lose later pruning, while the final consumer preserves correctness. 
Resampling/ratio policies need workload evidence.



##########
native/core/src/execution/operators/dynamic_filter/parquet_reader.rs:
##########
@@ -62,6 +62,14 @@ pub(super) fn try_attach_parquet_reader_filter(
         log::debug!("Join dynamic filter reader pushdown skipped: probe has a 
fetch limit");
         return Ok(None);
     }
+    // An ancestor join may already filter decoded batches below this join. 
Keep

Review Comment:
   Added intermediate-join `dynamic_filter_join_filters_attached > 0` 
assertions to both the two-join and committed three-join Scala regressions, 
with the runtime-filter flag enabled and each intermediate build side.



##########
native/operators/src/projection.rs:
##########
@@ -0,0 +1,211 @@
+// Licensed to the Apache Software Foundation (ASF) under one

Review Comment:
   Moved the metric-owning CometProjectionExec adapter to 
`native/operators/src/projection.rs`, beside CometFilterExec, and updated 
imports/module exports. Its module documentation describes projection metric 
ownership.



##########
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(

Review Comment:
   Renamed the helper to `prepare_probe_for_runtime_filters` and updated the 
HashJoin call-site comment to describe recursive filter/projection preparation 
for both reader attachment and early placement.



##########
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()
+}
+
+#[tokio::test]
+#[ignore = "component benchmark; run explicitly with an optimized build"]

Review Comment:
   Removed the ignored component benchmark. The PR description now reports 
current validation and avoids carrying old-head timing numbers forward.



##########
docs/source/user-guide/latest/metrics.md:
##########
@@ -271,3 +275,26 @@ If you compare Comet's `bytesRead` against vanilla Spark's 
on Spark 4.1+ (via th
 the REST API), expect Comet's number to be substantially larger for small 
files, and closer to
 Spark's for large files in that workload. Neither metric should be interpreted 
as complete
 filesystem or network traffic accounting.
+
+### Early filtering in join chains

Review Comment:
   Moved eligibility, placement, and adaptive-bypass behavior beneath Join 
Runtime Filters in `tuning/operators.md`. The metrics guide retains the metric 
rows and links to that section.



-- 
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