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]
