-
Notifications
You must be signed in to change notification settings - Fork 587
feat: bloom filter pushdown #2398
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
1c013de
0671bf1
87146af
73c4fdd
1559571
58ad776
e23e733
29072aa
1c4a8c1
5e48213
310b374
32302b7
eed3e03
6c67826
db5daf8
a1c6b92
3d80060
8adcd72
cfbf233
351f1dd
1bf11bd
631461d
053ddb3
1a1540e
76aca80
1f54c40
f3b179c
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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, | ||
|
|
@@ -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(), | ||
| }; | ||
|
|
@@ -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, | ||
| } | ||
|
|
@@ -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() { | ||
| selected_row_group_indices = Some(bloom_filtered); | ||
| } | ||
| } | ||
|
|
||
| if self.row_selection_enabled { | ||
| row_selection = ArrowReader::get_row_selection_for_filter_predicate( | ||
| &predicate, | ||
|
|
@@ -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>, | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. mut reference is because |
||
| 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 { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Every I'd fetch the filters concurrently —
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. join_all won't work here since 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?
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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. The thing that would make me comfortable is the doc on |
||
| 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() { | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Is this |
||
| 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) | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. These fetches look sequential, one awaited Separately, is there a way for a user to tell whether the option helped? Nothing appears to record row groups pruned, and
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 { | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Since
bloom_filteredis an order-preserving subset ofcandidate_rgs, isbloom_filtered.len() < candidate_rgs.len()ever false when the contents differ? Wondering whether the guard is needed.