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