zhuqi-lucas opened a new issue, #25268:
URL: https://github.com/apache/datafusion/issues/25268

   ### Describe the bug
   
   With `pushdown_filters = true`, a filter conjunct that plan time marks as 
pushed down (removing it from the parent `FilterExec`) can be **silently 
dropped at runtime** when the per-file `RowFilter` is rebuilt against the 
physical file schema — the query then returns rows that violate the `WHERE` 
clause. No error, no metric, only a `debug!` line.
   
   The safety invariant is "the plan-time pushdown checker must be strictly 
more conservative than runtime row-filter candidate building". Adapter-inserted 
casts validated by arrow's permissive `can_cast_types` break it:
   
   1. **Plan time says Yes**: `try_pushdown_filters` checks 
`can_expr_be_pushed_down_with_schemas` against the **table** schema 
(`datafusion/datasource-parquet/src/source.rs`). With table schema `a: Utf8`, 
`a = 'keep'` is a primitive-column predicate → `PushedDown::Yes` → the conjunct 
is removed from `FilterExec`.
   2. **Runtime adapter inserts a *valid* cast**: for a file whose physical 
schema stores `a` as `List<Int32>` (schema evolution), the physical expr 
rewriter validates with arrow's `can_cast_types`, which allows `List → Utf8` → 
a `CastExpr` is inserted with no error (`schema_rewriter.rs` → 
`nested_struct.rs::validate_data_type_compatibility`, `_` fallback arm).
   3. **Runtime pushdown checker rejects the shape**: 
`prebuild_row_filter_candidates` → `PushdownChecker::check_single_column` sees 
the physical `List` type; `=` is not in `supports_list_predicates` (only 
`IsNull`/`IsNotNull`/`array_has*`) → the candidate comes back `Ok(None)`.
   4. **Silent drop**: the `None` vanishes in `.flatten()` in 
`build_row_filter`/`prebuild_row_filter_candidates` (`row_filter.rs`, the doc 
comment even states "Conjuncts that cannot be evaluated as an ArrowPredicate 
are ignored"). No `RowFilter` is installed for that file; row-group/page 
pruning only proves "may contain matches" and never filters rows → the 
predicate is applied **nowhere** for that file.
   
   The cast itself is perfectly evaluable — the `pushdown_filters=false` path 
proves it (the retained `FilterExec` evaluates the adapted expression 
correctly). The runtime *checker* is just more conservative than the adapter's 
validator.
   
   A second reachable instance of the same shape: physical `FixedSizeList(T, 
1)` vs logical primitive `T` (arrow allows `FSL(1) → T`, the checker rejects 
`FixedSizeList` for non-list predicates).
   
   There is also a latent adjacent hazard: if `prebuild_row_filter_candidates` 
returns `Err`, the entire `RowFilter` is dropped with only a `debug!` 
(`push_decoder.rs`, `Err(e) => { debug!(...); None }`). I could not construct a 
reachable input for that leg, but it silently discards ALL pushed conjuncts if 
ever hit.
   
   ### To Reproduce
   
   Failing test (appended to 
`datafusion/core/tests/parquet/schema_coercion.rs`; two parquet files in one 
listing table — `a: Utf8` in f1, `a: List<Int32>` in f2, table schema declares 
`a Utf8`):
   
   Result with `pushdown_filters = true` (physical plan is a bare 
`DataSourceExec ... predicate=a@1 = keep`, no `FilterExec`):
   
   ```
   | id | a      |
   |----|--------|
   | 1  | keep   |
   | 3  | [1, 2] |   <- violates WHERE a = 'keep'
   | 4  | [3]    |   <- violates WHERE a = 'keep'
   ```
   
   Same query with `pushdown_filters = false` correctly returns only `id = 1`.
   
   <details>
   <summary>Full repro test</summary>
   
   ```rust
   #[tokio::test]
   async fn pushdown_filter_dropped_conjunct_returns_wrong_rows() {
       use arrow::array::{AsArray, Int32Array, ListArray};
       use arrow::buffer::OffsetBuffer;
       use datafusion::datasource::file_format::parquet::ParquetFormat;
       use datafusion::datasource::listing::ListingOptions;
       use datafusion::prelude::SessionConfig;
   
       let dir = tempfile::tempdir().unwrap();
   
       // file 1: a is Utf8 (matches table schema)
       let batch1 = RecordBatch::try_from_iter(vec![
           (
               "id",
               Arc::new(Int64Array::from(vec![1_i64, 2])) as ArrayRef,
           ),
           (
               "a",
               Arc::new(StringArray::from(vec!["keep", "drop"])) as ArrayRef,
           ),
       ])
       .unwrap();
   
       // file 2: a is physically List<Int32> (schema evolution)
       let list = ListArray::new(
           Arc::new(Field::new("item", DataType::Int32, true)),
           OffsetBuffer::from_lengths([2, 1]),
           Arc::new(Int32Array::from(vec![1, 2, 3])),
           None,
       );
       let batch2 = RecordBatch::try_from_iter(vec![
           (
               "id",
               Arc::new(Int64Array::from(vec![3_i64, 4])) as ArrayRef,
           ),
           ("a", Arc::new(list) as ArrayRef),
       ])
       .unwrap();
   
       for (name, batch) in [("f1.parquet", batch1), ("f2.parquet", batch2)] {
           let file = std::fs::File::create(dir.path().join(name)).unwrap();
           let mut w = ArrowWriter::try_new(file, batch.schema(), 
None).unwrap();
           w.write(&batch).unwrap();
           w.close().unwrap();
       }
   
       let mut cfg = SessionConfig::new();
       cfg.options_mut().execution.parquet.pushdown_filters = true;
       let ctx = SessionContext::new_with_config(cfg);
   
       let table_schema = Arc::new(Schema::new(vec![
           Field::new("id", DataType::Int64, true),
           Field::new("a", DataType::Utf8, true),
       ]));
       let opts = ListingOptions::new(Arc::new(ParquetFormat::default()));
       ctx.register_listing_table(
           "t",
           dir.path().to_str().unwrap(),
           opts,
           Some(table_schema),
           None,
       )
       .await
       .unwrap();
   
       // Show the physical plan (the FilterExec should have been removed /
       // the predicate reported as handled by the scan)
       let plan = ctx
           .sql("EXPLAIN SELECT id, a FROM t WHERE a = 'keep'")
           .await
           .unwrap()
           .collect()
           .await
           .unwrap();
       println!(
           "{}",
           datafusion_common::test_util::batches_to_string(&plan)
       );
   
       let results = ctx
           .sql("SELECT id, a FROM t WHERE a = 'keep' ORDER BY id")
           .await
           .unwrap()
           .collect()
           .await
           .unwrap();
   
       println!(
           "{}",
           datafusion_common::test_util::batches_to_string(&results)
       );
   
       // Sanity: with pushdown_filters=false the same query is answered 
correctly
       let ctx2 = SessionContext::new();
       let opts2 = ListingOptions::new(Arc::new(ParquetFormat::default()));
       ctx2.register_listing_table(
           "t",
           dir.path().to_str().unwrap(),
           opts2,
           Some(Arc::new(Schema::new(vec![
               Field::new("id", DataType::Int64, true),
               Field::new("a", DataType::Utf8, true),
           ]))),
           None,
       )
       .await
       .unwrap();
       let no_pushdown = ctx2
           .sql("SELECT id, a FROM t WHERE a = 'keep' ORDER BY id")
           .await
           .unwrap()
           .collect()
           .await
           .unwrap();
       println!(
           "no-pushdown result:\n{}",
           datafusion_common::test_util::batches_to_string(&no_pushdown)
       );
       let no_pushdown_rows: usize = no_pushdown.iter().map(|b| 
b.num_rows()).sum();
       assert_eq!(no_pushdown_rows, 1, "without pushdown only id=1 matches");
   
   }
   ```
   
   </details>
   
   ### Expected behavior
   
   Either the query returns correct results (runtime accepts the evaluable 
adapted cast), or the scan fails loudly. Never silent extra rows.
   
   Fix directions (not mutually exclusive):
   1. **Make the drop loud**: when a conjunct that plan time reported as 
`PushedDown::Yes` produces no candidate at runtime, return an error (or fall 
back to evaluating the residual conjuncts post-decode) instead of ignoring it; 
at minimum bump a metric.
   2. **Close the checker/adapter gap**: let the runtime `PushdownChecker` 
accept cast-over-nested-column shapes the adapter can produce and evaluate.
   
   ### Additional context
   
   - Verified present on current `main` (load-bearing files `row_filter.rs`, 
`push_decoder.rs`, `projection_read_plan.rs` byte-identical to tip at time of 
writing).
   - Related history: #24109 was the same bug class for struct columns and was 
fixed by `try_narrow_struct_cast`; missing columns, primitive type evolution, 
and explicit user struct casts are all guarded correctly (I checked each).
   - This directly affects the pushdown-by-default work (#24426): the blast 
radius of the silent drop grows if `pushdown_filters` becomes the default.
   - Found during an aggregation/scan audit; repro + investigation notes 
available. I plan to follow up with a fix PR.
   


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