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]
