This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/main/pr-3314-e9c22bf7148a2fadd8dc94879371866dea89c883
in repository https://gitbox.apache.org/repos/asf/iceberg-rust.git

commit ffa864bf2fd3ba27ff221746f993d4bf7810440b
Author: Anoop Johnson <[email protected]>
AuthorDate: Fri Oct 2 13:57:07 2026 +0000

    fix(reader): page-index pruning drops rows for INT32/INT64 decimals (#3314)
    
    Page-index bounds for INT32/INT64 columns were built as Int/Long literals
    regardless of the Iceberg field type. For decimal fields, query datums use
    Int128, and Datum::partial_cmp only compares decimals when both sides are
    Int128, so every comparison returned None. This silently pruned all matching
    pages for inequality and IN predicates when page-index row selection was
    enabled.
    
    Convert INT32/INT64 page bounds to Int128 for decimal fields, matching the
    row-group statistics path.
    
    Follow-up: FIXED_LEN_BYTE_ARRAY decimals (precision > 18) still get no
    page-index pruning, even though the row-group stats path decodes them. This
    only over-selects pages (never drops rows) and is pre-existing. Decoding 
them
    here later must guard against truncated ColumnIndex bounds
    (column_index_truncate_length).
    
    Closes #3264
    
    Co-authored-by: Fokko Driesprong <[email protected]>
---
 .../src/expr/visitors/page_index_evaluator.rs      | 265 ++++++++++++++++++++-
 1 file changed, 259 insertions(+), 6 deletions(-)

diff --git a/crates/iceberg/src/expr/visitors/page_index_evaluator.rs 
b/crates/iceberg/src/expr/visitors/page_index_evaluator.rs
index 37d9c6a23..d494e5426 100644
--- a/crates/iceberg/src/expr/visitors/page_index_evaluator.rs
+++ b/crates/iceberg/src/expr/visitors/page_index_evaluator.rs
@@ -276,8 +276,8 @@ impl<'a> PageIndexEvaluator<'a> {
                 .zip(row_counts.iter())
                 .map(|((i, (min, max)), &row_count)| {
                     predicate(
-                        min.map(|&val| Datum::new(field_type.clone(), 
PrimitiveLiteral::Int(val))),
-                        max.map(|&val| Datum::new(field_type.clone(), 
PrimitiveLiteral::Int(val))),
+                        min.map(|&val| Self::int32_bound_to_datum(field_type, 
val)),
+                        max.map(|&val| Self::int32_bound_to_datum(field_type, 
val)),
                         PageNullCount::from_row_and_null_counts(row_count, 
idx.null_count(i)),
                     )
                 })
@@ -289,8 +289,8 @@ impl<'a> PageIndexEvaluator<'a> {
                 .zip(row_counts.iter())
                 .map(|((i, (min, max)), &row_count)| {
                     predicate(
-                        min.map(|&val| Datum::new(field_type.clone(), 
PrimitiveLiteral::Long(val))),
-                        max.map(|&val| Datum::new(field_type.clone(), 
PrimitiveLiteral::Long(val))),
+                        min.map(|&val| Self::int64_bound_to_datum(field_type, 
val)),
+                        max.map(|&val| Self::int64_bound_to_datum(field_type, 
val)),
                         PageNullCount::from_row_and_null_counts(row_count, 
idx.null_count(i)),
                     )
                 })
@@ -410,6 +410,30 @@ impl<'a> PageIndexEvaluator<'a> {
         Ok(Some(result?))
     }
 
+    /// Converts an `INT32` page bound into a [`Datum`] according to the 
field's
+    /// primitive type.
+    fn int32_bound_to_datum(field_type: &PrimitiveType, val: i32) -> Datum {
+        match field_type {
+            PrimitiveType::Decimal { .. } => Datum::new(
+                field_type.clone(),
+                PrimitiveLiteral::Int128(i128::from(val)),
+            ),
+            _ => Datum::new(field_type.clone(), PrimitiveLiteral::Int(val)),
+        }
+    }
+
+    /// Converts an `INT64` page bound into a [`Datum`] according to the 
field's
+    /// primitive type.
+    fn int64_bound_to_datum(field_type: &PrimitiveType, val: i64) -> Datum {
+        match field_type {
+            PrimitiveType::Decimal { .. } => Datum::new(
+                field_type.clone(),
+                PrimitiveLiteral::Int128(i128::from(val)),
+            ),
+            _ => Datum::new(field_type.clone(), PrimitiveLiteral::Long(val)),
+        }
+    }
+
     /// Converts a `BYTE_ARRAY` page bound into a [`Datum`] according to the
     /// field's primitive type. Parquet stores Iceberg `string` and `binary`
     /// bounds as `BYTE_ARRAY`.
@@ -832,7 +856,8 @@ mod tests {
     use std::sync::Arc;
 
     use arrow_array::{
-        ArrayRef, FixedSizeBinaryArray, Float32Array, LargeBinaryArray, 
RecordBatch, StringArray,
+        ArrayRef, Decimal128Array, FixedSizeBinaryArray, Float32Array, 
LargeBinaryArray,
+        RecordBatch, StringArray,
     };
     use arrow_schema::{DataType, Field, Schema as ArrowSchema};
     use parquet::arrow::ArrowWriter;
@@ -846,6 +871,7 @@ mod tests {
 
     use super::PageIndexEvaluator;
     use crate::expr::{Bind, Reference};
+    use crate::spec::decimal_utils::decimal_from_i128_with_scale;
     use crate::spec::{Datum, NestedField, PrimitiveType, Schema, Type};
     use crate::{ErrorKind, Result};
 
@@ -1036,6 +1062,55 @@ mod tests {
         Ok((metadata, temp_file))
     }
 
+    /// Creates a single-column `Decimal128(precision, scale)` parquet file.
+    /// Parquet encodes precision <= 9 as `INT32` and precision 10..=18 as
+    /// `INT64`, so `precision` selects which page-index encoding is exercised.
+    /// Writes one 1024-row page per value in `unscaled`, so page bounds
+    /// partition the value range.
+    fn create_decimal_parquet_file(
+        precision: u8,
+        scale: i8,
+        unscaled: &[i128],
+    ) -> Result<(Arc<ParquetMetaData>, NamedTempFile)> {
+        let arrow_schema = Arc::new(ArrowSchema::new(vec![Field::new(
+            "col_decimal",
+            DataType::Decimal128(precision, scale),
+            true,
+        )]));
+
+        let temp_file = NamedTempFile::new().unwrap();
+        let file = temp_file.reopen().unwrap();
+
+        let props = WriterProperties::builder()
+            .set_data_page_row_count_limit(1024)
+            .set_write_batch_size(512)
+            .build();
+
+        let mut writer = ArrowWriter::try_new(file, arrow_schema.clone(), 
Some(props)).unwrap();
+
+        for &unscaled in unscaled {
+            let array = Arc::new(
+                
Decimal128Array::from_iter_values(std::iter::repeat_n(unscaled, 1024))
+                    .with_precision_and_scale(precision, scale)
+                    .unwrap(),
+            ) as ArrayRef;
+            let batch = RecordBatch::try_new(arrow_schema.clone(), 
vec![array]).unwrap();
+            // Write rows one at a time so the writer splits into per-value 
pages.
+            for i in 0..batch.num_rows() {
+                writer.write(&batch.slice(i, 1)).unwrap();
+            }
+        }
+
+        writer.close().unwrap();
+
+        let file = temp_file.reopen().unwrap();
+        let options = 
ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
+        let reader = 
ParquetRecordBatchReaderBuilder::try_new_with_options(file, options).unwrap();
+        let metadata = reader.metadata().clone();
+
+        Ok((metadata, temp_file))
+    }
+
     /// Get the test metadata components for testing
     fn get_test_metadata(
         metadata: &ParquetMetaData,
@@ -1392,7 +1467,7 @@ mod tests {
         // A predicate that would prune every page if the bounds were decoded.
         let filter = Reference::new("col_decimal")
             .greater_than(Datum::decimal_with_precision(
-                
crate::spec::decimal_utils::decimal_from_i128_with_scale(99999, 2),
+                decimal_from_i128_with_scale(99999, 2),
                 10,
             )?)
             .bind(iceberg_schema.clone(), false)?;
@@ -1651,6 +1726,184 @@ mod tests {
         Ok(())
     }
 
+    #[test]
+    fn eval_inequality_prunes_int32_decimal_pages() -> Result<()> {
+        // precision 9 -> Parquet INT32 page bounds.
+        let (metadata, _temp_file) = create_decimal_parquet_file(9, 2, &[100, 
200, 300, 400])?;
+        let (column_index, offset_index, row_group_metadata) = 
get_test_metadata(&metadata);
+        let (iceberg_schema, field_id_map) = 
build_decimal_schema_and_field_map(9, 2)?;
+
+        // Pages hold 1.00, 2.00, 3.00, 4.00. `> 2.50` keeps the pages whose
+        // upper bound exceeds 2.50 (3.00 and 4.00).
+        let filter = Reference::new("col_decimal")
+            .greater_than(decimal_datum(250, 2, 9)?)
+            .bind(iceberg_schema.clone(), false)?;
+
+        let result = PageIndexEvaluator::eval(
+            &filter,
+            &column_index,
+            &offset_index,
+            row_group_metadata,
+            &field_id_map,
+            iceberg_schema.as_ref(),
+        )?;
+
+        assert_eq!(result, vec![
+            RowSelector::skip(2048),
+            RowSelector::select(2048)
+        ]);
+
+        Ok(())
+    }
+
+    #[test]
+    fn eval_in_prunes_int32_decimal_pages() -> Result<()> {
+        // precision 9 -> Parquet INT32 page bounds.
+        let (metadata, _temp_file) = create_decimal_parquet_file(9, 2, &[100, 
200, 300, 400])?;
+        let (column_index, offset_index, row_group_metadata) = 
get_test_metadata(&metadata);
+        let (iceberg_schema, field_id_map) = 
build_decimal_schema_and_field_map(9, 2)?;
+
+        // Pages hold 1.00, 2.00, 3.00, 4.00. IN (2.00, 4.00) keeps only the
+        // pages whose single value is one of the literals.
+        let filter = Reference::new("col_decimal")
+            .is_in([decimal_datum(200, 2, 9)?, decimal_datum(400, 2, 9)?])
+            .bind(iceberg_schema.clone(), false)?;
+
+        let result = PageIndexEvaluator::eval(
+            &filter,
+            &column_index,
+            &offset_index,
+            row_group_metadata,
+            &field_id_map,
+            iceberg_schema.as_ref(),
+        )?;
+
+        assert_eq!(result, vec![
+            RowSelector::skip(1024),
+            RowSelector::select(1024),
+            RowSelector::skip(1024),
+            RowSelector::select(1024),
+        ]);
+
+        Ok(())
+    }
+
+    #[test]
+    fn eval_inequality_prunes_int64_decimal_pages() -> Result<()> {
+        // precision 18 -> Parquet INT64 page bounds.
+        let (metadata, _temp_file) = create_decimal_parquet_file(18, 2, &[100, 
200, 300, 400])?;
+        let (column_index, offset_index, row_group_metadata) = 
get_test_metadata(&metadata);
+        let (iceberg_schema, field_id_map) = 
build_decimal_schema_and_field_map(18, 2)?;
+
+        // Pages hold 1.00, 2.00, 3.00, 4.00. `> 2.50` keeps the pages whose
+        // upper bound exceeds 2.50 (3.00 and 4.00).
+        let filter = Reference::new("col_decimal")
+            .greater_than(decimal_datum(250, 2, 18)?)
+            .bind(iceberg_schema.clone(), false)?;
+
+        let result = PageIndexEvaluator::eval(
+            &filter,
+            &column_index,
+            &offset_index,
+            row_group_metadata,
+            &field_id_map,
+            iceberg_schema.as_ref(),
+        )?;
+
+        assert_eq!(result, vec![
+            RowSelector::skip(2048),
+            RowSelector::select(2048)
+        ]);
+
+        Ok(())
+    }
+
+    #[test]
+    fn eval_in_prunes_int64_decimal_pages() -> Result<()> {
+        // precision 18 -> Parquet INT64 page bounds.
+        let (metadata, _temp_file) = create_decimal_parquet_file(18, 2, &[100, 
200, 300, 400])?;
+        let (column_index, offset_index, row_group_metadata) = 
get_test_metadata(&metadata);
+        let (iceberg_schema, field_id_map) = 
build_decimal_schema_and_field_map(18, 2)?;
+
+        // Pages hold 1.00, 2.00, 3.00, 4.00. IN (2.00, 4.00) keeps only the
+        // pages whose single value is one of the literals.
+        let filter = Reference::new("col_decimal")
+            .is_in([decimal_datum(200, 2, 18)?, decimal_datum(400, 2, 18)?])
+            .bind(iceberg_schema.clone(), false)?;
+
+        let result = PageIndexEvaluator::eval(
+            &filter,
+            &column_index,
+            &offset_index,
+            row_group_metadata,
+            &field_id_map,
+            iceberg_schema.as_ref(),
+        )?;
+
+        assert_eq!(result, vec![
+            RowSelector::skip(1024),
+            RowSelector::select(1024),
+            RowSelector::skip(1024),
+            RowSelector::select(1024),
+        ]);
+
+        Ok(())
+    }
+
+    #[test]
+    fn eval_inequality_prunes_negative_decimal_pages_at_boundary() -> 
Result<()> {
+        // precision 9 -> Parquet INT32 page bounds, spanning negative values.
+        let (metadata, _temp_file) = create_decimal_parquet_file(9, 2, &[-400, 
-200, 100, 300])?;
+        let (column_index, offset_index, row_group_metadata) = 
get_test_metadata(&metadata);
+        let (iceberg_schema, field_id_map) = 
build_decimal_schema_and_field_map(9, 2)?;
+
+        // Pages hold -4.00, -2.00, 1.00, 3.00. `>= -2.00` keeps the pages 
whose
+        // upper bound is at least -2.00, including page 1 whose bound equals 
it.
+        let filter = Reference::new("col_decimal")
+            .greater_than_or_equal_to(decimal_datum(-200, 2, 9)?)
+            .bind(iceberg_schema.clone(), false)?;
+
+        let result = PageIndexEvaluator::eval(
+            &filter,
+            &column_index,
+            &offset_index,
+            row_group_metadata,
+            &field_id_map,
+            iceberg_schema.as_ref(),
+        )?;
+
+        assert_eq!(result, vec![
+            RowSelector::skip(1024),
+            RowSelector::select(3072)
+        ]);
+
+        Ok(())
+    }
+
+    /// Builds a decimal query [`Datum`] from an `unscaled` mantissa, matching
+    /// the field's `scale` and `precision`.
+    fn decimal_datum(unscaled: i128, scale: u32, precision: u32) -> 
Result<Datum> {
+        Datum::decimal_with_precision(decimal_from_i128_with_scale(unscaled, 
scale), precision)
+    }
+
+    fn build_decimal_schema_and_field_map(
+        precision: u32,
+        scale: u32,
+    ) -> Result<(Arc<Schema>, HashMap<i32, usize>)> {
+        let iceberg_schema = Arc::new(
+            Schema::builder()
+                .with_fields([Arc::new(NestedField::new(
+                    1,
+                    "col_decimal",
+                    Type::Primitive(PrimitiveType::Decimal { precision, scale 
}),
+                    true,
+                ))])
+                .build()?,
+        );
+
+        Ok((iceberg_schema, HashMap::from_iter([(1, 0)])))
+    }
+
     fn build_iceberg_schema_and_field_map() -> Result<(Arc<Schema>, 
HashMap<i32, usize>)> {
         let iceberg_schema = Schema::builder()
             .with_fields([

Reply via email to