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


##########
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:
   Thanks for keeping the traversal this conservative, the error and limit 
boundaries look right to me. One interaction I'd like to check is with Spark's 
inferred `IS NOT NULL` filters. Since `place()` descends through the null-check 
`CometFilterExec`, the early consumer ends up directly above the scan and sees 
null keys. Null keys are rejected by the domain, so any batch with a null key 
resets `unselective_batches` in `DynamicFilterExec::execute`. On a fact table 
with nullable foreign keys (TPC-DS fact FKs contain nulls), almost every batch 
has one. So the bypass never triggers and an unselective ancestor filter is 
evaluated twice on every row. Descending through the null check also doesn't 
buy anything here, because the early consumer never attaches to a reader.
   
   Would it work to place the terminal consumer above a null-check filter once 
a join has been crossed? Another option is to judge selectivity against the 
non-null key count. Could you add a test with a nullable key and a 
non-selective ancestor filter that shows the bypass kicking in? The current 
benchmark uses a non-nullable key with no null-check filter, so it wouldn't 
catch this.



##########
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:
   This new section ends up nested under "Task-Level Input Metrics on Spark 
4.1+", which is unrelated. The behavioral description (eligibility, where 
placement stops, the adaptive bypass) seems to belong in `tuning/operators.md` 
under "Join Runtime Filters". That section still says the domain only filters 
probe batches before the hash probe. `metrics.md` could then keep just the 
table rows and a link.



##########
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:
   The PR description mentions a three-join regression that passed locally. 
Could you commit it? Stacked early consumers from two ancestors are the only 
thing that exercises the `DynamicFilterExec` arm of `place()` and the reader 
attachment underneath it.
   
   Also, CI runs everything with `join.dynamicFilter.enabled=false`, so this 
test and the native tests are the only coverage of the new path. Could you 
share a TPC-DS run with the flag on, both the result verification and the 
`dynamic_filter_early_*` metrics? Star-join chains are exactly the target 
workload, and that would show the benefit holds beyond the in-memory component 
benchmark.



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

Review Comment:
   This is nearly a line-for-line copy of `CometFilterExec` in 
`native/operators/src/filter.rs`, including the module docs. Could it live next 
to that one, for example as `native/operators/src/projection.rs`? Then the two 
metric-owning adapters stay together and a later change to one is less likely 
to miss the other. The file name `runtime_filter_projection.rs` also doesn't 
match the type name.



##########
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:
   The bypass condition is "two consecutive batches with zero rows removed". 
That means a filter keeping 99% of rows never bypasses and pays double 
evaluation for the whole input. It also means a key-clustered input whose first 
two batches all match gives up permanently. Have you considered a ratio 
threshold, or re-sampling periodically after bypassing? If the current rule is 
deliberate, a short comment on why exact-zero was chosen would help future 
readers.



##########
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:
   `prepare_probe_filter_for_runtime_reader` now recursively wraps projections 
too, and it serves early placement as well as reader attachment. Could you 
rename it to something like `prepare_probe_for_runtime_filters`? The call-site 
comment in the `HashJoin` arm (it still talks about "the probe filter's child") 
should be updated to match.



##########
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:
   This branch is what keeps the intermediate join's own reader pruning working 
once an ancestor's early consumer sits above its scan. Without it, enabling the 
flag would quietly turn off row-group pruning for the intermediate join. Could 
the new `CometJoinSuite` test (or a native test) assert that the intermediate 
join still reports `dynamic_filter_join_filters_attached` when the feature is 
on?



##########
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:
   Comet keeps benchmarks under `native/core/benches` with criterion. As an 
`#[ignore]` unit test, nothing in CI ever runs this, so it'll drift. Could it 
move to a criterion bench? Or could it be dropped from this PR, with the 
numbers kept in the description?



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