comphead commented on code in PR #3302:
URL: https://github.com/apache/iceberg-rust/pull/3302#discussion_r4139977182


##########
crates/iceberg/src/arrow/reader/predicate_visitor.rs:
##########
@@ -195,6 +197,238 @@ impl BoundPredicateVisitor for CollectFieldIdVisitor {
     }
 }
 
+/// Returns the residual of `predicate` for one data file: each leaf on a 
top-level field that is
+/// missing from the file becomes `AlwaysTrue` or `AlwaysFalse`, based on the 
value projection
+/// returns for that field. That value is the identity partition value, 
otherwise the field's
+/// `initial-default`, and every row of the file holds it. Leaves on missing 
fields with neither
+/// keep the null handling of the row filter and the page index evaluator.
+///
+/// Leaves on nested fields are kept, because a nested field also reads as 
null in any row where
+/// an ancestor struct is null.
+pub(super) fn residual_for_missing_fields(
+    predicate: BoundPredicate,
+    predicate_field_ids: &HashSet<i32>,
+    field_id_map: &HashMap<i32, usize>,
+    schema: &Schema,
+    partition_spec: Option<&PartitionSpec>,
+    partition: Option<&Struct>,
+) -> Result<BoundPredicate> {
+    if predicate_field_ids
+        .iter()
+        .all(|id| field_id_map.contains_key(id))
+    {
+        return Ok(predicate);
+    }
+
+    let partition_constants = match (partition_spec, partition) {
+        (Some(spec), Some(data)) => constants_map(spec, data, schema)?,
+        _ => HashMap::new(),
+    };
+
+    let mut field_ids = HashSet::new();
+    let row: Struct = schema
+        .as_struct()
+        .fields()
+        .iter()
+        .map(|field| {
+            if !predicate_field_ids.contains(&field.id) || 
field_id_map.contains_key(&field.id) {
+                return None;
+            }
+            let value = match partition_constants.get(&field.id) {
+                Some(datum) => 
Some(Literal::Primitive(datum.literal().clone())),
+                None => field.initial_default.clone(),
+            };
+            if value.is_some() {

Review Comment:
   Question, not for this PR. The null handling that leaves on a field with no 
value fall back to is not uniform. In `PredicateConverter`, `less_than` and 
`less_than_or_eq` return `build_always_true()` for a missing column, while 
`greater_than`, `greater_than_or_eq`, and `eq` return false. So for a column 
added without a default, `WHERE d < 5` keeps every row of an older file and 
`WHERE d > 5` keeps none. Java's `ParquetMetricsRowGroupFilter` treats a 
missing column as all null, and `TestMetricsRowGroupFilter.testColumnNotInFile` 
expects `lt` and `le` to be unable to match. The issue calls the current null 
handling correct, so I wanted to check whether the `<` and `<=` branches are 
intended. If not, a separate issue would be enough.



##########
crates/iceberg/src/arrow/reader/predicate_visitor.rs:
##########
@@ -195,6 +197,238 @@ impl BoundPredicateVisitor for CollectFieldIdVisitor {
     }
 }
 
+/// Returns the residual of `predicate` for one data file: each leaf on a 
top-level field that is
+/// missing from the file becomes `AlwaysTrue` or `AlwaysFalse`, based on the 
value projection
+/// returns for that field. That value is the identity partition value, 
otherwise the field's
+/// `initial-default`, and every row of the file holds it. Leaves on missing 
fields with neither
+/// keep the null handling of the row filter and the page index evaluator.
+///
+/// Leaves on nested fields are kept, because a nested field also reads as 
null in any row where
+/// an ancestor struct is null.
+pub(super) fn residual_for_missing_fields(
+    predicate: BoundPredicate,
+    predicate_field_ids: &HashSet<i32>,
+    field_id_map: &HashMap<i32, usize>,
+    schema: &Schema,
+    partition_spec: Option<&PartitionSpec>,
+    partition: Option<&Struct>,
+) -> Result<BoundPredicate> {
+    if predicate_field_ids
+        .iter()
+        .all(|id| field_id_map.contains_key(id))
+    {
+        return Ok(predicate);
+    }
+
+    let partition_constants = match (partition_spec, partition) {
+        (Some(spec), Some(data)) => constants_map(spec, data, schema)?,
+        _ => HashMap::new(),
+    };
+
+    let mut field_ids = HashSet::new();
+    let row: Struct = schema
+        .as_struct()
+        .fields()
+        .iter()
+        .map(|field| {
+            if !predicate_field_ids.contains(&field.id) || 
field_id_map.contains_key(&field.id) {
+                return None;
+            }
+            let value = match partition_constants.get(&field.id) {
+                Some(datum) => 
Some(Literal::Primitive(datum.literal().clone())),
+                None => field.initial_default.clone(),
+            };
+            if value.is_some() {
+                field_ids.insert(field.id);
+            }
+            value
+        })
+        .collect();
+
+    if field_ids.is_empty() {
+        return Ok(predicate);
+    }
+    visit(
+        &mut MissingFieldResidualVisitor {
+            row: &row,
+            field_ids: &field_ids,
+        },
+        &predicate,
+    )
+}
+
+/// Replaces leaves on the fields in `field_ids` with their result on `row`, 
which holds each
+/// top-level field's value at its position in the schema.
+struct MissingFieldResidualVisitor<'a> {
+    row: &'a Struct,
+    field_ids: &'a HashSet<i32>,
+}
+
+impl MissingFieldResidualVisitor<'_> {
+    fn residual(
+        &self,
+        reference: &BoundReference,
+        predicate: &BoundPredicate,
+    ) -> Result<BoundPredicate> {
+        if !self.field_ids.contains(&reference.field().id) {
+            return Ok(predicate.clone());
+        }
+        if visit(&mut ExpressionEvaluatorVisitor::new(self.row), predicate)? {
+            Ok(BoundPredicate::AlwaysTrue)
+        } else {
+            Ok(BoundPredicate::AlwaysFalse)
+        }
+    }
+}
+
+impl BoundPredicateVisitor for MissingFieldResidualVisitor<'_> {
+    type T = BoundPredicate;
+
+    fn always_true(&mut self) -> Result<BoundPredicate> {
+        Ok(BoundPredicate::AlwaysTrue)
+    }
+
+    fn always_false(&mut self) -> Result<BoundPredicate> {
+        Ok(BoundPredicate::AlwaysFalse)
+    }
+
+    fn and(&mut self, lhs: BoundPredicate, rhs: BoundPredicate) -> 
Result<BoundPredicate> {
+        Ok(lhs.and(rhs))
+    }
+
+    fn or(&mut self, lhs: BoundPredicate, rhs: BoundPredicate) -> 
Result<BoundPredicate> {
+        Ok(lhs.or(rhs))
+    }
+
+    fn not(&mut self, inner: BoundPredicate) -> Result<BoundPredicate> {

Review Comment:
   `not` does not normally see a `Not`. `with_filter` applies `rewrite_not`, 
and equality delete predicates are built without one. If one did arrive, 
rebuilding it here would give `PageIndexEvaluator::not` a node it rejects (`NOT 
unsupported at this point. NOT-rewrite should be performed first`). 
`inner.negate()` is what `RewriteNotVisitor::not` does. It would also let 
`LogicalExpression::new` stay private and turn the `NOT (b = 7)` test case into 
`AlwaysFalse`.



##########
crates/iceberg/src/arrow/reader/row_filter.rs:
##########
@@ -1757,4 +1763,182 @@ mod tests {
         .await;
         assert_eq!(rows(&on), 1);
     }
+
+    /// Writes a file that stores only field 1 (`a`), with values 1, 2, 3.
+    fn write_file_without_b() -> (String, TempDir) {
+        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);
+        (file_path, tmp_dir)
+    }
+
+    fn schema_with_b(b: NestedField) -> SchemaRef {
+        Arc::new(
+            Schema::builder()
+                .with_schema_id(1)
+                .with_fields(vec![
+                    NestedField::required(1, "a", 
Type::Primitive(PrimitiveType::Long)).into(),
+                    b.into(),
+                ])
+                .build()
+                .unwrap(),
+        )
+    }
+
+    /// Reads the file from [`write_file_without_b`], in which `b` reads as 7 
on every row through
+    /// `schema` or `partition`, and asserts the rows each predicate keeps. 
Each predicate must keep
+    /// or drop all three rows as if `b` were stored as 7.
+    async fn assert_absent_b_reads_as_7(
+        schema: SchemaRef,
+        partition: Option<(Arc<PartitionSpec>, Struct)>,
+        cases: Vec<(Predicate, usize)>,
+    ) {
+        let (file_path, _tmp_dir) = write_file_without_b();
+        let read = async |predicate: Predicate| {
+            let task = FileScanTask::builder()
+                
.with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
+                .with_start(0)
+                .with_length(0)
+                .with_data_file_path(file_path.clone())
+                .with_data_file_format(DataFileFormat::Parquet)
+                .with_schema(schema.clone())
+                .with_project_field_ids(vec![1, 2])
+                .with_predicate(Some(predicate.bind(schema.clone(), 
true).unwrap()))
+                .with_partition_spec(partition.as_ref().map(|(spec, _)| 
spec.clone()))
+                .with_partition(partition.as_ref().map(|(_, data)| 
data.clone()))
+                .with_case_sensitive(false)
+                .build()
+                .unwrap();
+            let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as 
FileScanTaskStream;
+            ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current())
+                .build()
+                .read(tasks)
+                .unwrap()
+                .stream()
+                .try_collect::<Vec<RecordBatch>>()
+                .await
+                .unwrap()
+        };
+
+        let batches = read(Predicate::AlwaysTrue).await;
+        assert_eq!(
+            batches[0].column(1).as_primitive::<Int64Type>().values(),
+            &[7, 7, 7]
+        );
+
+        let mut actual = Vec::new();
+        for (predicate, _) in &cases {
+            let batches = read(predicate.clone()).await;
+            let rows = 
batches.iter().map(RecordBatch::num_rows).sum::<usize>();
+            actual.push((predicate.to_string(), rows));
+        }
+        let expected: Vec<_> = cases
+            .iter()
+            .map(|(predicate, rows)| (predicate.to_string(), *rows))
+            .collect();
+        assert_eq!(actual, expected);
+    }
+
+    /// A file written before column `b` was added with `initial-default` 7.
+    #[tokio::test]
+    async fn test_predicate_on_absent_column_uses_initial_default() {
+        let schema = schema_with_b(
+            NestedField::optional(2, "b", Type::Primitive(PrimitiveType::Long))
+                .with_initial_default(Literal::long(7)),
+        );
+        let b = || Reference::new("b");
+        assert_absent_b_reads_as_7(schema, None, vec![
+            (b().equal_to(Datum::long(7)), 3),
+            (b().is_not_null(), 3),
+            (b().is_in([Datum::long(7), Datum::long(8)]), 3),
+            (b().greater_than(Datum::long(5)), 3),
+            (b().equal_to(Datum::long(8)), 0),
+            (b().is_null(), 0),
+        ])
+        .await;
+    }
+
+    /// A file that doesn't store its identity partition column `b`, as after 
a Hive migration or
+    /// `add_files`, with partition value 7.
+    #[tokio::test]
+    async fn 
test_predicate_on_absent_identity_partition_column_uses_partition_value() {
+        let schema = schema_with_b(NestedField::optional(
+            2,
+            "b",
+            Type::Primitive(PrimitiveType::Long),
+        ));
+        let partition_spec = PartitionSpec::builder(schema.clone())
+            .add_partition_field("b", "b", Transform::Identity)
+            .unwrap()
+            .build()
+            .unwrap();
+        let partition = Struct::from_iter([Some(Literal::long(7))]);
+        let b = || Reference::new("b");
+        assert_absent_b_reads_as_7(schema, Some((Arc::new(partition_spec), 
partition)), vec![
+            (b().equal_to(Datum::long(7)), 3),
+            (b().is_not_null(), 3),
+            (b().equal_to(Datum::long(8)), 0),
+            (b().is_null(), 0),
+        ])
+        .await;
+    }
+
+    #[tokio::test]
+    async fn test_page_index_on_absent_column_uses_initial_default() {

Review Comment:
   This test applies `residual_for_missing_fields` by hand, so it would still 
pass if the call were dropped from or moved within `pipeline.rs`. After the 
residual all three predicates are `AlwaysTrue`, which leaves the page index 
evaluator little to decide. The `row_selection_enabled = true` leg of 
`test_scan_filter_on_initial_default_column_absent_from_file` already covers 
this path end to end. I would drop this test. If you want row selection covered 
for the identity partition case, `assert_absent_b_reads_as_7` could loop over 
`with_row_selection_enabled(false)` and `(true)`.



##########
crates/iceberg/src/arrow/reader/pipeline.rs:
##########
@@ -603,6 +604,14 @@ impl FileScanTaskReader {
                 &predicate,
                 use_position_fallback,
             )?;
+            let predicate = residual_for_missing_fields(

Review Comment:
   The description says this also covers the predicate built from equality 
deletes, but no test exercises it, and it is the case with the largest effect. 
The keep predicate for a delete row `b = 7` is `b IS NULL OR b != 7` 
(`parse_equality_deletes_record_batch_stream`). On `main`, a file written 
before `b` was added hits the missing-column branch, where `is_null` is always 
true. Every row is kept and the delete never applies. With this change both 
leaves fold to false and the rows are deleted. I traced this by reading the 
code and did not run it.
   
   Could you add a case? The cheap version is two rows in 
`test_residual_for_missing_fields` with `b` defaulting to 7. `b IS NULL OR b != 
7` should fold to false and `b IS NULL OR b != 5` to true. A reader-level 
version could follow the position-delete test in `row_filter.rs` that passes a 
`FileScanTaskDeleteFile`, using an equality delete file with `equality_ids: 
Some(vec![2])`.



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