mbutrovich opened a new issue, #3298: URL: https://github.com/apache/iceberg-rust/issues/3298
## Apache Iceberg Rust version main (deab0692f6df211a396d6fc0ad21811f8b6c8491) ## Describe the bug A column added by schema evolution can carry an [`initial-default`](https://github.com/apache/iceberg/blob/48330b8dacab6662242d252b39c8444190979bb2/format/spec.md?plain=1#L325-L334), the value that every row written before the column existed must read as. The spec's [column projection rules](https://github.com/apache/iceberg/blob/48330b8dacab6662242d252b39c8444190979bb2/format/spec.md?plain=1#L401-L409) say a reader returns that default for a field ID that isn't in the data file. `ArrowReader` gets projection right. `RecordBatchTransformer` [fills the default in](https://github.com/apache/iceberg-rust/blob/deab0692f6df211a396d6fc0ad21811f8b6c8491/crates/iceberg/src/arrow/record_batch_transformer.rs#L815-L820). Filtering gets it wrong. When a scan predicate references a column that is absent from the file, two places treat the column as all null and ignore `initial-default`: 1. The Arrow row filter built in `predicate_visitor.rs`. Every "A missing column, treating it as null" branch, for example [`is_null`](https://github.com/apache/iceberg-rust/blob/deab0692f6df211a396d6fc0ad21811f8b6c8491/crates/iceberg/src/arrow/reader/predicate_visitor.rs#L355-L369) and [`eq`](https://github.com/apache/iceberg-rust/blob/deab0692f6df211a396d6fc0ad21811f8b6c8491/crates/iceberg/src/arrow/reader/predicate_visitor.rs#L500-L518), returns a constant that assumes null. This runs on every read with a predicate. 2. `PageIndexEvaluator`, which prunes rows using the Parquet page index (per-page min/max statistics) when row selection is enabled. For a column missing from the file, [`calc_row_selection` exits early](https://github.com/apache/iceberg-rust/blob/deab0692f6df211a396d6fc0ad21811f8b6c8491/crates/iceberg/src/expr/visitors/page_index_evaluator.rs#L139-L148) and skips every row for [`not_null`](https://github.com/apache/iceberg-rust/blob/deab0692f6df211a396d6fc0ad21811f8b6c8491/crates/iceberg/src/expr/visitors/page_index_evaluator.rs#L512-L524), [`eq`](https://github.com/apache/iceberg-rust/blob/deab0692f6df211a396d6fc0ad21811f8b6c8491/crates/iceberg/src/expr/visitors/page_index_evaluator.rs#L585-L616), and [`in`](https://github.com/apache/iceberg-rust/blob/deab0692f6df211a396d6fc0ad21811f8b6c8491/crates/iceberg/src/expr/visitors/page_index_evaluator.rs#L771-L814). Row selection is off by default. The result is a silent wrong answer. With a column `b` added with `initial-default` 7, an unfiltered scan returns `b = 7` on every old row, but `WHERE b = 7` returns none of them and `WHERE b IS NULL` returns all of them. Any data file written before the column was added is affected, including files that are still live after later appends. Engines that push down filters, such as DataFusion, hit this on ordinary queries. The other pruning layers don't drop these rows, as far as I can tell from reading them. `RowGroupMetricsEvaluator`, `InclusiveMetricsEvaluator`, and `BloomFilterEvaluator` all treat a column with no statistics as "might match". Iceberg Java had the same bug in its Parquet row-group filter and fixed it in apache/iceberg#16692 (issue apache/iceberg#16690). The fix in [`ParquetMetricsRowGroupFilter.predicate`](https://github.com/apache/iceberg/blob/48330b8dacab6662242d252b39c8444190979bb2/parquet/src/main/java/org/apache/iceberg/parquet/ParquetMetricsRowGroupFilter.java#L134-L155) evaluates the predicate against `initialDefault` when the column is missing from the file. For a nested field it also allows null, because a nested field reads as null whenever an ancestor struct is null. Java's row-level filtering was already correct, because its Parquet readers materialize the default before the engine applies the residual filter. This is a gap for the default-value audit in #3250, which doesn't cover filtering yet. ## To Reproduce The two tests below fail on main. They go in the `tests` module of `crates/iceberg/src/arrow/reader/row_filter.rs` and reuse its `field_with_id`, `write_row_groups`, and `read_once` helpers. They need these imports added to that module: ```rust use arrow_array::types::Int64Type; use parquet::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions}; use parquet::file::metadata::PageIndexPolicy; use crate::spec::Literal; ``` ```rust /// A file written before column `b` was added, read with a schema in which `b` /// has `initial-default` 7. Every row reads back as `b = 7`, so each predicate /// must keep or drop all three rows as if `b` were stored as 7. fn setup_absent_column_with_initial_default() -> (Arc<Schema>, String, TempDir) { let schema = Arc::new( Schema::builder() .with_schema_id(1) .with_fields(vec![ NestedField::required(1, "a", Type::Primitive(PrimitiveType::Long)).into(), NestedField::optional(2, "b", Type::Primitive(PrimitiveType::Long)) .with_initial_default(Literal::long(7)) .into(), ]) .build() .unwrap(), ); let tmp_dir = TempDir::new().unwrap(); let file_path = format!("{}/1.parquet", tmp_dir.path().to_str().unwrap()); let arrow_schema = Arc::new(ArrowSchema::new(vec![field_with_id( "a", DataType::Int64, 1, )])); let batch = RecordBatch::try_new(arrow_schema.clone(), vec![Arc::new(Int64Array::from( vec![1, 2, 3], ))]) .unwrap(); write_row_groups(&file_path, arrow_schema, vec![batch], false); (schema, file_path, tmp_dir) } #[tokio::test] async fn test_predicate_on_absent_column_uses_initial_default() { let (schema, file_path, _tmp_dir) = setup_absent_column_with_initial_default(); // Projection applies the default. let (batches, _) = read_once( &file_path, schema.clone(), vec![1, 2], Predicate::AlwaysTrue, false, ) .await; assert_eq!( batches[0].column(1).as_primitive::<Int64Type>().values(), &[7, 7, 7] ); let cases = [ ("b = 7", Reference::new("b").equal_to(Datum::long(7)), 3), ("b IS NOT NULL", Reference::new("b").is_not_null(), 3), ( "b IN (7, 8)", Reference::new("b").is_in([Datum::long(7), Datum::long(8)]), 3, ), ("b > 5", Reference::new("b").greater_than(Datum::long(5)), 3), ("b = 8", Reference::new("b").equal_to(Datum::long(8)), 0), ("b IS NULL", Reference::new("b").is_null(), 0), ]; let mut actual = Vec::new(); for (name, predicate, _) in &cases { let (batches, _) = read_once( &file_path, schema.clone(), vec![1, 2], predicate.clone(), false, ) .await; actual.push(( *name, batches.iter().map(RecordBatch::num_rows).sum::<usize>(), )); } let expected: Vec<_> = cases.iter().map(|(name, _, rows)| (*name, *rows)).collect(); assert_eq!(actual, expected); } #[tokio::test] async fn test_page_index_on_absent_column_uses_initial_default() { let (schema, file_path, _tmp_dir) = setup_absent_column_with_initial_default(); let file = File::open(&file_path).unwrap(); let metadata = ArrowReaderMetadata::load( &file, ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required), ) .unwrap() .metadata() .clone(); let field_id_map = HashMap::from([(1, 0)]); let cases = [ ("b = 7", Reference::new("b").equal_to(Datum::long(7))), ("b IS NOT NULL", Reference::new("b").is_not_null()), ( "b IN (7, 8)", Reference::new("b").is_in([Datum::long(7), Datum::long(8)]), ), ]; let mut actual = Vec::new(); for (name, predicate) in &cases { let selection = ArrowReader::get_row_selection_for_filter_predicate( &predicate.clone().bind(schema.clone(), true).unwrap(), &metadata, &None, &field_id_map, &schema, ) .unwrap() .expect("the file has a page index"); actual.push((*name, selection.row_count())); } let expected: Vec<_> = cases.iter().map(|(name, _)| (*name, 3)).collect(); assert_eq!(actual, expected); } ``` Running `cargo test -p iceberg --lib -- absent_column_uses_initial_default` on main gives the following. The projection assertion passes, and the filter assertions fail: ``` ---- arrow::reader::row_filter::tests::test_predicate_on_absent_column_uses_initial_default stdout ---- assertion `left == right` failed left: [("b = 7", 0), ("b IS NOT NULL", 0), ("b IN (7, 8)", 0), ("b > 5", 0), ("b = 8", 0), ("b IS NULL", 3)] right: [("b = 7", 3), ("b IS NOT NULL", 3), ("b IN (7, 8)", 3), ("b > 5", 3), ("b = 8", 0), ("b IS NULL", 0)] ---- arrow::reader::row_filter::tests::test_page_index_on_absent_column_uses_initial_default stdout ---- assertion `left == right` failed left: [("b = 7", 0), ("b IS NOT NULL", 0), ("b IN (7, 8)", 0)] right: [("b = 7", 3), ("b IS NOT NULL", 3), ("b IN (7, 8)", 3)] ``` The page index test calls `get_row_selection_for_filter_predicate` directly. Through `ArrowReader` both bugs drop the same rows, so an end-to-end read can't tell them apart. ## Expected behavior A predicate on a column that is missing from a data file should be evaluated as if every row held the field's `initial-default`, which is the value projection already returns. When the field has no `initial-default`, the current null handling is correct. Page index pruning should keep the rows whenever the default can match the predicate. A nested field with an `initial-default` can also read as null when an ancestor struct is null, so a predicate on it should keep rows that match either the default or null, as Java does. #3261 tracks a related gap where `RecordBatchTransformer` doesn't apply `initial-default` to nested fields. ## Willingness to contribute I can contribute a fix for this bug independently -- 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]
