This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg-rust.git
The following commit(s) were added to refs/heads/main by this push:
new ffa864bf2 fix(reader): page-index pruning drops rows for INT32/INT64
decimals (#3314)
ffa864bf2 is described below
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([