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


##########
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:
   Done. `not` returns `inner.negate()`, `LogicalExpression::new` is private 
again, and the `NOT (b = 7)` case expects `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:
   I dropped it. `assert_absent_b_reads_as_7` now runs every filter and 
equality delete case with row selection off and on.



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