apache / apache/datafusion

Wrong results: filter marked as pushed down at plan time can be silently dropped per-file when schema evolution makes the runtime pushdown checker reject the adapted predicate

Open
#25,268 0 comments 0 reactions 1 assignee Claimed by @zhuqi-lucas View on GitHub
bug
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

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

Full repro test

```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 (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");

}
```

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

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.