Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions crates/iceberg/public-api.txt
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@ impl iceberg::arrow::ArrowReaderBuilder
pub fn iceberg::arrow::ArrowReaderBuilder::build(self) -> iceberg::arrow::ArrowReader
pub fn iceberg::arrow::ArrowReaderBuilder::new(file_io: iceberg::io::FileIO, runtime: iceberg::Runtime) -> Self
pub fn iceberg::arrow::ArrowReaderBuilder::with_batch_size(self, batch_size: usize) -> Self
pub fn iceberg::arrow::ArrowReaderBuilder::with_bloom_filter_enabled(self, bloom_filter_enabled: bool) -> Self
pub fn iceberg::arrow::ArrowReaderBuilder::with_data_file_concurrency_limit(self, val: usize) -> Self
pub fn iceberg::arrow::ArrowReaderBuilder::with_metadata_size_hint(self, metadata_size_hint: usize) -> Self
pub fn iceberg::arrow::ArrowReaderBuilder::with_range_coalesce_bytes(self, range_coalesce_bytes: u64) -> Self
Expand Down Expand Up @@ -1363,6 +1364,7 @@ pub fn iceberg::scan::TableScanBuilder<'a>::select_all(self) -> Self
pub fn iceberg::scan::TableScanBuilder<'a>::select_empty(self) -> Self
pub fn iceberg::scan::TableScanBuilder<'a>::snapshot_id(self, snapshot_id: i64) -> Self
pub fn iceberg::scan::TableScanBuilder<'a>::with_batch_size(self, batch_size: core::option::Option<usize>) -> Self
pub fn iceberg::scan::TableScanBuilder<'a>::with_bloom_filter_enabled(self, bloom_filter_enabled: bool) -> Self
pub fn iceberg::scan::TableScanBuilder<'a>::with_case_sensitive(self, case_sensitive: bool) -> Self
pub fn iceberg::scan::TableScanBuilder<'a>::with_concurrency_limit(self, limit: usize) -> Self
pub fn iceberg::scan::TableScanBuilder<'a>::with_data_file_concurrency_limit(self, limit: usize) -> Self
Expand Down
19 changes: 19 additions & 0 deletions crates/iceberg/src/arrow/reader/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ pub struct ArrowReaderBuilder {
concurrency_limit_data_files: usize,
row_group_filtering_enabled: bool,
row_selection_enabled: bool,
bloom_filter_enabled: bool,
parquet_read_options: ParquetReadOptions,
runtime: Runtime,
}
Expand All @@ -72,6 +73,7 @@ impl ArrowReaderBuilder {
concurrency_limit_data_files: num_cpus,
row_group_filtering_enabled: true,
row_selection_enabled: false,
bloom_filter_enabled: false,
parquet_read_options: ParquetReadOptions::builder().build(),
runtime,
}
Expand Down Expand Up @@ -102,6 +104,21 @@ impl ArrowReaderBuilder {
self
}

/// Determines whether to enable bloom filter-based row group filtering.
///
/// When enabled, if a read is performed with an equality or IN predicate,
/// the bloom filter for relevant columns in each row group is read and
/// checked. Row groups where the bloom filter proves the value is absent
/// are skipped entirely.
///
/// Defaults to disabled. Each bloom filter is a separate read, and they are
/// issued serially — one round trip per relevant column per row group, before
/// any data is read. TODO(#3191)
pub fn with_bloom_filter_enabled(mut self, bloom_filter_enabled: bool) -> Self {
self.bloom_filter_enabled = bloom_filter_enabled;
self
}

/// Provide a hint as to the number of bytes to prefetch for parsing the Parquet metadata
///
/// This hint can help reduce the number of fetch requests. For more details see the
Expand Down Expand Up @@ -141,6 +158,7 @@ impl ArrowReaderBuilder {
concurrency_limit_data_files: self.concurrency_limit_data_files,
row_group_filtering_enabled: self.row_group_filtering_enabled,
row_selection_enabled: self.row_selection_enabled,
bloom_filter_enabled: self.bloom_filter_enabled,
parquet_read_options: self.parquet_read_options,
}
}
Expand All @@ -158,5 +176,6 @@ pub struct ArrowReader {

row_group_filtering_enabled: bool,
row_selection_enabled: bool,
bloom_filter_enabled: bool,
parquet_read_options: ParquetReadOptions,
}
102 changes: 102 additions & 0 deletions crates/iceberg/src/arrow/reader/pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,10 @@ use crate::arrow::record_batch_transformer::RecordBatchTransformerBuilder;
use crate::arrow::scan_metrics::{CountingFileRead, ScanMetrics, ScanResult};
use crate::encryption::StandardKeyMetadata;
use crate::error::Result;
use crate::expr::BoundPredicate;
use crate::expr::visitors::bloom_filter_evaluator::{
BloomFilterEvaluator, ColumnBloomFilter, collect_bloom_filter_field_ids,
};
use crate::io::{FileIO, FileMetadata, FileRead};
use crate::metadata_columns::{
RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER, RESERVED_COL_NAME_POS,
Expand All @@ -70,6 +74,7 @@ impl ArrowReader {
.with_scan_metrics(scan_metrics.clone()),
row_group_filtering_enabled: self.row_group_filtering_enabled,
row_selection_enabled: self.row_selection_enabled,
bloom_filter_enabled: self.bloom_filter_enabled,
parquet_read_options: self.parquet_read_options,
scan_metrics: scan_metrics.clone(),
};
Expand Down Expand Up @@ -123,6 +128,7 @@ struct FileScanTaskReader {
delete_file_loader: CachingDeleteFileLoader,
row_group_filtering_enabled: bool,
row_selection_enabled: bool,
bloom_filter_enabled: bool,
parquet_read_options: ParquetReadOptions,
scan_metrics: ScanMetrics,
}
Expand Down Expand Up @@ -632,6 +638,30 @@ impl FileScanTaskReader {
};
}

if self.bloom_filter_enabled {
let all_rgs;
let candidate_rgs = match &selected_row_group_indices {
Some(indices) => indices.as_slice(),
None => {
all_rgs = (0..record_batch_stream_builder.metadata().num_row_groups())
.collect::<Vec<_>>();
&all_rgs
}
};

let bloom_filtered = Self::filter_row_groups_by_bloom_filter(
&predicate,
&mut record_batch_stream_builder,
candidate_rgs,
&field_id_map,
)
.await?;

if bloom_filtered.len() < candidate_rgs.len() {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Since bloom_filtered is an order-preserving subset of candidate_rgs, is bloom_filtered.len() < candidate_rgs.len() ever false when the contents differ? Wondering whether the guard is needed.

selected_row_group_indices = Some(bloom_filtered);
}
}

if self.row_selection_enabled {
row_selection = ArrowReader::get_row_selection_for_filter_predicate(
&predicate,
Expand Down Expand Up @@ -692,6 +722,78 @@ impl FileScanTaskReader {

Ok(Box::pin(record_batch_stream) as ArrowRecordBatchStream)
}

/// Reads bloom filters for relevant columns and evaluates the predicate
/// against them to filter out row groups that definitely don't match.
async fn filter_row_groups_by_bloom_filter(
predicate: &BoundPredicate,
builder: &mut ParquetRecordBatchStreamBuilder<ArrowFileReader>,

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

mut reference is because get_row_group_column_bloom_filter requires it.

candidate_row_groups: &[usize],
field_id_map: &HashMap<i32, usize>,
) -> Result<Vec<usize>> {
// Only collect field IDs from eq/in predicates — the only types
// bloom filters can help with. Skip columns not in the parquet schema.
let bloom_filter_field_ids: Vec<i32> = collect_bloom_filter_field_ids(predicate)?
.into_iter()
.filter(|id| field_id_map.contains_key(id))
.collect();

if bloom_filter_field_ids.is_empty() {
return Ok(candidate_row_groups.to_vec());
}

let mut result = Vec::with_capacity(candidate_row_groups.len());

for &rg_idx in candidate_row_groups {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Every (row_group, column) pair here ends in a separate sequential .await on get_row_group_column_bloom_filter. For 100 row groups and 3 predicate columns that's 300 round trips in series before a single group is skipped — on object storage that latency can easily swamp whatever I/O the pruning saves, which undercuts the reason to enable this at all.

I'd fetch the filters concurrently — join_all over the (rg, col) pairs, or at least per-column within a row group — so the round trips overlap. Thoughts?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

join_all won't work here since get_row_group_column_bloom_filter requires a mutable reference to self. Same reason a buffered stream won't work.

For what it's worth DataFusion does the same as us https://github.com/apache/datafusion/blob/55.0.0/datafusion/datasource-parquet/src/opener/mod.rs#L1266-L1292

So I'd rather do this as a follow up optimization rather than block this PR? WDYT?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Ah, https://github.com/apache/arrow-rs/pull/8462/changes is what we need. I'll create an issue to track a follow-up

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not a new comment, a response to your question about deferring.

Deferring seems right to me. get_row_group_column_bloom_filter takes &mut self, so this isn't a case of picking the wrong combinator, and arrow-rs#8462 is the actual unblock. The DataFusion reference is fair, and since the option is off by default, a user who opts in is accepting the current cost rather than silently paying it.

The thing that would make me comfortable is the doc on with_bloom_filter_enabled being explicit that the reads are serialized per row group per column, so the cost model is visible at the call site instead of only in the tracking issue. Right now it says extra I/O per column per row group, which reads as a volume cost rather than a latency one, and latency is what dominates on object storage.

let mut bloom_filters: HashMap<i32, ColumnBloomFilter> = HashMap::new();

for &field_id in &bloom_filter_field_ids {
let col_idx = field_id_map[&field_id];
let col_meta = builder.metadata().row_group(rg_idx).column(col_idx);

// Only attempt to load if this column chunk actually has a bloom filter
if col_meta.bloom_filter_offset().is_none() {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is this bloom_filter_offset().is_none() pre-check needed? get_row_group_column_bloom_filter appears to make the same check and return Ok(None) before any I/O, which the Ok(None) arm at L771 already handles.

continue;
}

let physical_type = col_meta.column_type();
let type_length = col_meta.column_descr().type_length();

match builder
.get_row_group_column_bloom_filter(rg_idx, col_idx)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

These fetches look sequential, one awaited get_row_group_column_bloom_filter at a time, nested inside the per-row-group loop. A file with 100 row groups and 3 predicate columns would be up to 300 serialized round trips. On high-latency object storage, could that cost more than the pruning saves? try_buffer_unordered against a concurrency limit is used in this file at L100.

Separately, is there a way for a user to tell whether the option helped? Nothing appears to record row groups pruned, and ScanMetrics only exposes bytes_read, so that may be better as a follow-up than something for this PR.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Answered the first part #2398 (comment)

Happy to track metrics in a follow-up!

.await
{
Ok(Some(sbbf)) => {
bloom_filters.insert(
field_id,
ColumnBloomFilter::new(sbbf, physical_type, type_length),
);
}
Ok(None) => {}
Err(e) => {
// Left absent from the map, so the evaluator treats the column
// as might-match and the row group survives.
tracing::debug!(
"Bloom filter for field {field_id} in row group {rg_idx} could not be read: {e}"
);
}
}
}

match BloomFilterEvaluator::eval(predicate, &bloom_filters) {
Ok(true) => result.push(rg_idx),
Ok(false) => { /* Row group pruned by bloom filter */ }
Err(e) => {
tracing::debug!(
"Bloom filter evaluation failed for row group {rg_idx}, including it: {e}"
);
result.push(rg_idx);
}
}
}

Ok(result)
}
}

impl ArrowReader {
Expand Down
Loading
Loading