From baf714e202bbc1d5920b7ef65ba091cdc1d98579 Mon Sep 17 00:00:00 2001 From: wenn-id Date: Wed, 29 Jul 2026 09:02:38 +0700 Subject: [PATCH] fix(parquet): check narrow Arrow integers in bloom filters --- parquet/src/arrow/arrow_writer/mod.rs | 40 +++++-- parquet/src/arrow/bloom_filter.rs | 153 ++++++++++++++++++++++++++ parquet/src/arrow/mod.rs | 1 + 3 files changed, 182 insertions(+), 12 deletions(-) create mode 100644 parquet/src/arrow/bloom_filter.rs diff --git a/parquet/src/arrow/arrow_writer/mod.rs b/parquet/src/arrow/arrow_writer/mod.rs index 207b4c3327ef..e2aa6319c80c 100644 --- a/parquet/src/arrow/arrow_writer/mod.rs +++ b/parquet/src/arrow/arrow_writer/mod.rs @@ -1524,6 +1524,30 @@ impl ArrowColumnWriterFactory { } } +pub(crate) fn coerce_narrow_integer_to_int32(array: &dyn arrow_array::Array) -> Result { + macro_rules! coerce { + ($ty:ty) => {{ + let array = array.as_any().downcast_ref::<$ty>().ok_or_else(|| { + ParquetError::General(format!( + "Arrow array has data type {} but an incompatible representation", + array.data_type() + )) + })?; + Ok(array.unary(|value| value as i32)) + }}; + } + + match array.data_type() { + ArrowDataType::Int8 => coerce!(arrow_array::Int8Array), + ArrowDataType::Int16 => coerce!(arrow_array::Int16Array), + ArrowDataType::UInt8 => coerce!(arrow_array::UInt8Array), + ArrowDataType::UInt16 => coerce!(arrow_array::UInt16Array), + data_type => Err(ParquetError::General(format!( + "Cannot coerce Arrow data type {data_type} to Parquet INT32" + ))), + } +} + fn write_leaf( writer: &mut ColumnWriter<'_>, column: &dyn arrow_array::Array, @@ -1539,23 +1563,15 @@ fn write_leaf( let array = Int32Array::new_null(column.len()); write_primitive(typed, array.values(), levels) } - ArrowDataType::Int8 => { - let array: Int32Array = column.as_primitive::().unary(|x| x as i32); - write_primitive(typed, array.values(), levels) - } - ArrowDataType::Int16 => { - let array: Int32Array = column.as_primitive::().unary(|x| x as i32); + ArrowDataType::Int8 | ArrowDataType::Int16 => { + let array = coerce_narrow_integer_to_int32(column)?; write_primitive(typed, array.values(), levels) } ArrowDataType::Int32 => { write_primitive(typed, column.as_primitive::().values(), levels) } - ArrowDataType::UInt8 => { - let array: Int32Array = column.as_primitive::().unary(|x| x as i32); - write_primitive(typed, array.values(), levels) - } - ArrowDataType::UInt16 => { - let array: Int32Array = column.as_primitive::().unary(|x| x as i32); + ArrowDataType::UInt8 | ArrowDataType::UInt16 => { + let array = coerce_narrow_integer_to_int32(column)?; write_primitive(typed, array.values(), levels) } ArrowDataType::UInt32 => { diff --git a/parquet/src/arrow/bloom_filter.rs b/parquet/src/arrow/bloom_filter.rs new file mode 100644 index 000000000000..1eb0ab91cb50 --- /dev/null +++ b/parquet/src/arrow/bloom_filter.rs @@ -0,0 +1,153 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Arrow-aware Parquet bloom filter kernels. + +use crate::arrow::arrow_writer::coerce_narrow_integer_to_int32; +use crate::bloom_filter::Sbbf; +use crate::errors::{ParquetError, Result}; +use arrow_array::{Array, BooleanArray}; +use arrow_schema::DataType; + +/// Checks an Arrow array against a Parquet bloom filter. +/// +/// The input uses its logical Arrow representation. Supported narrow integer +/// arrays (`Int8`, `Int16`, `UInt8`, and `UInt16`) are widened to the Parquet +/// `INT32` physical representation used by [`super::ArrowWriter`] before each +/// non-null value is checked. Null inputs produce null outputs. +/// +/// # Example +/// +/// ``` +/// use arrow_array::{BooleanArray, Int8Array}; +/// use parquet::arrow::bloom_filter::check_bloom_filter; +/// use parquet::bloom_filter::Sbbf; +/// +/// let mut filter = Sbbf::new_with_num_of_bytes(32); +/// filter.insert(&7_i32); // Arrow Int8 is stored as Parquet INT32 +/// let values = Int8Array::from(vec![Some(7), None]); +/// let actual = check_bloom_filter(&filter, &values)?; +/// assert_eq!(actual, BooleanArray::from(vec![Some(true), None])); +/// # Ok::<(), parquet::errors::ParquetError>(()) +/// ``` +pub fn check_bloom_filter(filter: &Sbbf, array: &dyn Array) -> Result { + if !matches!( + array.data_type(), + DataType::Int8 | DataType::Int16 | DataType::UInt8 | DataType::UInt16 + ) { + return Err(ParquetError::General(format!( + "Arrow data type {} is not supported by the bloom filter kernel", + array.data_type() + ))); + } + + let array = coerce_narrow_integer_to_int32(array)?; + Ok(BooleanArray::from_iter( + array + .iter() + .map(|value| value.map(|value| filter.check(&value))), + )) +} + +#[cfg(test)] +mod tests { + use super::check_bloom_filter; + use crate::arrow::ArrowWriter; + use crate::bloom_filter::Sbbf; + use crate::file::properties::{ReaderProperties, WriterProperties}; + use crate::file::reader::{FileReader, SerializedFileReader}; + use crate::file::serialized_reader::ReadOptionsBuilder; + use arrow_array::{ + ArrayRef, BooleanArray, Int8Array, Int16Array, RecordBatch, UInt8Array, UInt16Array, + }; + use arrow_schema::{DataType, Field, Schema}; + use bytes::Bytes; + use std::sync::Arc; + + fn write_bloom_filter(array: ArrayRef) -> Sbbf { + let schema = Arc::new(Schema::new(vec![Field::new( + "value", + array.data_type().clone(), + true, + )])); + let batch = RecordBatch::try_new(Arc::clone(&schema), vec![array]).unwrap(); + let props = WriterProperties::builder() + .set_bloom_filter_enabled(true) + .build(); + let mut buffer = Vec::new(); + let mut writer = ArrowWriter::try_new(&mut buffer, schema, Some(props)).unwrap(); + writer.write(&batch).unwrap(); + writer.close().unwrap(); + + let options = ReadOptionsBuilder::new() + .with_reader_properties( + ReaderProperties::builder() + .set_read_bloom_filter(true) + .build(), + ) + .build(); + let reader = SerializedFileReader::new_with_options(Bytes::from(buffer), options).unwrap(); + reader + .get_row_group(0) + .unwrap() + .get_column_bloom_filter(0) + .unwrap() + .clone() + } + + macro_rules! test_narrow_integer { + ($name:ident, $array:ty, $values:expr, $missing:expr) => { + #[test] + fn $name() { + let array = <$array>::from($values); + let bloom_filter = write_bloom_filter(Arc::new(array.clone())); + let query = + <$array>::from(vec![array.value(0), $missing, array.value(array.len() - 1)]); + let actual = check_bloom_filter(&bloom_filter, &query).unwrap(); + let expected = BooleanArray::from(vec![true, false, true]); + assert_eq!(actual, expected); + + let actual = check_bloom_filter(&bloom_filter, &array).unwrap(); + let expected = BooleanArray::from(vec![Some(true), None, Some(true)]); + assert_eq!(actual, expected); + } + }; + } + + test_narrow_integer!(check_int8, Int8Array, vec![Some(-128), None, Some(127)], 0); + test_narrow_integer!( + check_int16, + Int16Array, + vec![Some(-32_768), None, Some(32_767)], + 0 + ); + test_narrow_integer!(check_uint8, UInt8Array, vec![Some(0), None, Some(255)], 127); + test_narrow_integer!( + check_uint16, + UInt16Array, + vec![Some(0), None, Some(65_535)], + 32_767 + ); + + #[test] + fn rejects_unsupported_arrow_type() { + let array = BooleanArray::from(vec![true]); + let bloom_filter = write_bloom_filter(Arc::new(array.clone())); + let error = check_bloom_filter(&bloom_filter, &array).unwrap_err(); + assert!(error.to_string().contains(&DataType::Boolean.to_string())); + } +} diff --git a/parquet/src/arrow/mod.rs b/parquet/src/arrow/mod.rs index b89db361eda1..55edda063d50 100644 --- a/parquet/src/arrow/mod.rs +++ b/parquet/src/arrow/mod.rs @@ -182,6 +182,7 @@ experimental!(mod array_reader); pub mod arrow_reader; pub mod arrow_writer; +pub mod bloom_filter; mod buffer; mod decoder;