Skip to content
Open
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
7 changes: 7 additions & 0 deletions datafusion/common/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1133,6 +1133,13 @@ config_namespace! {
/// the hint, two reads will still be performed.
pub metadata_size_hint: Option<usize>, default = Some(512 * 1024)

/// (reading) If specified, the parquet reader will prefetch data for
/// subsequent row groups when the projected column chunks fit within
/// this many bytes. The required ranges for the current row group are
/// always read, even when they exceed this value. None disables data
/// prefetching.
pub prefetch_size: Option<usize>, default = None

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.

I agree this is an obvious and common strategy for reading from object store and we should add it

One thing I am wondering is if we can make this IO strategy more generic (for example, I can imagine a user wanting to support racing reads like the craziness described in #23492 (comment)

Where I am heading is "can we add an API / trait" that lets user customize the I/O behavior more?

For example, maybe the trait kind of like madvise that could have methods like

/// The reader will need the bytes in `range` at some point in the future
/// it will request the bytes in the same order as called by this function
/// though it may not read all of them
fn advise_bytes_needed(&self, range: Range<usize>)

And then we could include a default implementation that implemented memory limited prefetch 🤔 but users could provide their own implementations that did other things (like racing reads, etc)


/// (reading) If true, filter expressions are be applied during the parquet decoding operation to
/// reduce the number of rows decoded. This optimization is sometimes called "late materialization".
pub pushdown_filters: bool, default = false
Expand Down
3 changes: 3 additions & 0 deletions datafusion/common/src/file_options/parquet_writer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -235,6 +235,7 @@ impl ParquetOptions {
pruning: _,
skip_metadata: _,
metadata_size_hint: _,
prefetch_size: _,
pushdown_filters: _,
reorder_filters: _,
force_filter_selections: _, // not used for writer props
Expand Down Expand Up @@ -491,6 +492,7 @@ mod tests {
pruning: defaults.pruning,
skip_metadata: defaults.skip_metadata,
metadata_size_hint: defaults.metadata_size_hint,
prefetch_size: defaults.prefetch_size,
pushdown_filters: defaults.pushdown_filters,
reorder_filters: defaults.reorder_filters,
force_filter_selections: defaults.force_filter_selections,
Expand Down Expand Up @@ -610,6 +612,7 @@ mod tests {
pruning: global_options_defaults.pruning,
skip_metadata: global_options_defaults.skip_metadata,
metadata_size_hint: global_options_defaults.metadata_size_hint,
prefetch_size: global_options_defaults.prefetch_size,
pushdown_filters: global_options_defaults.pushdown_filters,
reorder_filters: global_options_defaults.reorder_filters,
force_filter_selections: global_options_defaults.force_filter_selections,
Expand Down
8 changes: 8 additions & 0 deletions datafusion/datasource-parquet/src/opener/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -252,6 +252,8 @@ pub(super) struct ParquetMorselizer {
/// Optional hint for how large the initial request to read parquet metadata
/// should be
pub metadata_size_hint: Option<usize>,
/// Maximum number of bytes buffered by speculative row-group prefetching.
pub prefetch_size: Option<usize>,
/// Metrics for reporting
pub metrics: ExecutionPlanMetricsSet,
/// Factory for instantiating parquet reader
Expand Down Expand Up @@ -425,6 +427,7 @@ struct PreparedParquetOpen {
baseline_metrics: BaselineMetrics,
file_pruner: Option<FilePruner>,
metadata_size_hint: Option<usize>,
prefetch_size: Option<usize>,
metrics: ExecutionPlanMetricsSet,
parquet_file_reader_factory: Arc<dyn ParquetFileReaderFactory>,
async_file_reader: Box<dyn AsyncFileReader>,
Expand Down Expand Up @@ -828,6 +831,7 @@ impl ParquetMorselizer {
baseline_metrics,
file_pruner,
metadata_size_hint,
prefetch_size: self.prefetch_size,
metrics: self.metrics.clone(),
parquet_file_reader_factory: Arc::clone(&self.parquet_file_reader_factory),
async_file_reader,
Expand Down Expand Up @@ -1482,6 +1486,9 @@ impl RowGroupsPrunedParquetOpen {
active_reader: None,
rg_plan,
reader: prepared.async_file_reader,
parquet_metadata: Arc::clone(&file_metadata),
prefetch_size: prepared.prefetch_size,
prefetched_row_groups: HashSet::new(),
decoder_projection,
arrow_reader_metrics,
predicate_cache_inner_records,
Expand Down Expand Up @@ -2014,6 +2021,7 @@ mod test {
predicate: self.predicate,
table_schema,
metadata_size_hint: self.metadata_size_hint,
prefetch_size: None,
metrics: self.metrics,
parquet_file_reader_factory: self
.parquet_file_reader_factory
Expand Down
228 changes: 225 additions & 3 deletions datafusion/datasource-parquet/src/push_decoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
//! [`PushDecoderStreamState::into_stream`] for consumption.

use std::collections::VecDeque;
use std::ops::Range;
use std::sync::Arc;

use arrow::array::RecordBatch;
Expand All @@ -53,7 +54,7 @@ use parquet::arrow::async_reader::AsyncFileReader;
use parquet::arrow::push_decoder::{ParquetPushDecoder, ParquetPushDecoderBuilder};
use parquet::file::metadata::ParquetMetaData;

use datafusion_common::{DataFusionError, Result};
use datafusion_common::{DataFusionError, HashSet, Result};
use datafusion_physical_expr::expressions::DynamicFilterTracking;
use datafusion_physical_expr_common::physical_expr::PhysicalExpr;
use datafusion_physical_plan::metrics::{BaselineMetrics, Count, Gauge};
Expand Down Expand Up @@ -244,6 +245,14 @@ pub(crate) struct PushDecoderStreamState {
pub(crate) active_reader: Option<ParquetRecordBatchReader>,
pub(crate) rg_plan: VecDeque<RgPlanEntry>,
pub(crate) reader: Box<dyn AsyncFileReader>,
/// Parquet metadata used to identify projected column chunks belonging to
/// subsequent row groups.
pub(crate) parquet_metadata: Arc<ParquetMetaData>,
/// Maximum bytes that may be staged in the push decoder. Required ranges
/// are always fetched, even when they exceed this budget.
pub(crate) prefetch_size: Option<usize>,
/// Row groups whose projected ranges were already fetched speculatively.
pub(crate) prefetched_row_groups: HashSet<usize>,
/// Per-file projection: the mask installed on every decoder and the
/// per-batch transform applied by [`Self::project_batch`].
pub(crate) decoder_projection: DecoderProjection,
Expand Down Expand Up @@ -373,6 +382,20 @@ impl PushDecoderStreamState {
let decoder = self.decoder.as_mut().expect("decoder present");
match decoder.try_next_reader() {
Ok(DecodeResult::NeedsData(ranges)) => {
let buffered_bytes = self
.decoder
.as_ref()
.expect("decoder present")
.buffered_bytes();
let ranges = prefetch_row_group_ranges(
ranges,
buffered_bytes,
self.prefetch_size,
&self.rg_plan,
self.decoder_projection.projection_mask(),
&self.parquet_metadata,
&mut self.prefetched_row_groups,
);
let data = self
.reader
.get_byte_ranges(ranges.clone())
Expand Down Expand Up @@ -423,6 +446,68 @@ impl PushDecoderStreamState {
}
}

/// Append projected column chunks from subsequent row groups to a decoder
/// request while staying within `prefetch_size`.
///
/// The first entry in `rg_plan` is the row group responsible for `ranges`.
/// Complete projected ranges for later row groups are added in scan order so
/// the push decoder can stage them for future calls to `try_next_reader`.
fn prefetch_row_group_ranges(
mut ranges: Vec<Range<u64>>,
buffered_bytes: u64,
prefetch_size: Option<usize>,
rg_plan: &VecDeque<RgPlanEntry>,
projection: &ProjectionMask,
metadata: &ParquetMetaData,
prefetched_row_groups: &mut HashSet<usize>,
) -> Vec<Range<u64>> {
let Some(prefetch_size) = prefetch_size.filter(|size| *size > 0) else {
return ranges;
};

let requested_bytes = ranges
.iter()
.map(|range| range.end - range.start)
.sum::<u64>();
let mut staged_bytes = buffered_bytes.saturating_add(requested_bytes);
let budget = prefetch_size as u64;
if staged_bytes >= budget {
return ranges;
}

for entry in rg_plan.iter().skip(1) {
if prefetched_row_groups.contains(&entry.rg_index) {
continue;
}

let row_group = metadata.row_group(entry.rg_index);

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.

this would be a nice API to add upstream in the parquet crate probably

let row_group_ranges = row_group
.columns()
.iter()
.enumerate()
.filter(|(column_idx, _)| projection.leaf_included(*column_idx))
.map(|(_, column)| {
let (start, len) = column.byte_range();
start..start + len
})
.collect::<Vec<_>>();
let row_group_bytes = row_group_ranges
.iter()
.map(|range| range.end - range.start)
.sum::<u64>();

if staged_bytes.saturating_add(row_group_bytes) > budget {
break;
}

ranges.extend(row_group_ranges);
staged_bytes += row_group_bytes;
prefetched_row_groups.insert(entry.rg_index);
}

ranges
}

#[cfg(test)]
mod tests {
use super::*;
Expand All @@ -444,6 +529,11 @@ mod tests {
/// column statistics are disjoint: RG0 → 0..1000, RG1 → 1000..2000,
/// RG2 → 2000..3000. Returns (metadata, schema).
fn build_three_rg_file() -> (Arc<ParquetMetaData>, SchemaRef) {
let (_, metadata, schema) = build_three_rg_file_data();
(metadata, schema)
}

fn build_three_rg_file_data() -> (Bytes, Arc<ParquetMetaData>, SchemaRef) {
let schema = Arc::new(Schema::new(vec![Field::new("v", DataType::Int64, false)]));
let mut buf = Vec::new();
let props = WriterProperties::builder()
Expand Down Expand Up @@ -474,12 +564,144 @@ mod tests {
reason = "we want a single range covering the whole file"
)]
let ranges = vec![0..len];
md.push_ranges(ranges, vec![file]).unwrap();
md.push_ranges(ranges, vec![file.clone()]).unwrap();
let DecodeResult::Data(meta) = md.try_decode().unwrap() else {
panic!("decoding metadata");
};
assert_eq!(meta.num_row_groups(), 3, "test fixture must have 3 RGs");
(Arc::new(meta), schema)
(file, Arc::new(meta), schema)
}

fn column_range(metadata: &ParquetMetaData, row_group: usize) -> Range<u64> {
let (start, len) = metadata.row_group(row_group).column(0).byte_range();
start..start + len
}

#[test]
fn prefetch_packs_complete_row_groups_within_budget() {
let (metadata, _) = build_three_rg_file();
let rg0 = column_range(&metadata, 0);
let rg1 = column_range(&metadata, 1);
let rg_plan = (0..3).map(|rg_index| RgPlanEntry { rg_index }).collect();
let budget = ((rg0.end - rg0.start) + (rg1.end - rg1.start)) as usize;
let mut prefetched = HashSet::new();

let ranges = prefetch_row_group_ranges(
vec![rg0.clone()],
0,
Some(budget),
&rg_plan,
&ProjectionMask::all(),
&metadata,
&mut prefetched,
);

// RG2 does not fit in the budget, so it is left out entirely.
assert_eq!(ranges, vec![rg0, rg1]);
assert_eq!(prefetched, HashSet::from([1]));
}

#[test]
fn prefetched_bytes_are_staged_for_the_next_push_decoder_reader() {
let (file, metadata, _) = build_three_rg_file_data();
let mut decoder =
ParquetPushDecoderBuilder::try_new_decoder(Arc::clone(&metadata))
.unwrap()
.build()
.unwrap();
let requested = match decoder.try_next_reader().unwrap() {
DecodeResult::NeedsData(ranges) => ranges,
other => panic!("expected initial byte request, got {other:?}"),
};
let rg_plan = (0..3).map(|rg_index| RgPlanEntry { rg_index }).collect();
let rg1 = column_range(&metadata, 1);
let budget = requested
.iter()
.map(|range| range.end - range.start)
.sum::<u64>()
+ rg1.end
- rg1.start;
let mut prefetched = HashSet::new();
let ranges = prefetch_row_group_ranges(
requested,
decoder.buffered_bytes(),
Some(budget as usize),
&rg_plan,
&ProjectionMask::all(),
&metadata,
&mut prefetched,
);
let data = ranges
.iter()
.map(|range| file.slice(range.start as usize..range.end as usize))
.collect();
decoder.push_ranges(ranges, data).unwrap();

let DecodeResult::Data(first_reader) = decoder.try_next_reader().unwrap() else {
panic!("first row group should be ready");
};
assert_eq!(
first_reader
.map(|batch| batch.unwrap().num_rows())
.sum::<usize>(),
1000
);

let DecodeResult::Data(second_reader) = decoder.try_next_reader().unwrap() else {
panic!("prefetched second row group should not require more I/O");
};
assert_eq!(
second_reader
.map(|batch| batch.unwrap().num_rows())
.sum::<usize>(),
1000
);
}

#[test]
fn prefetch_accounts_for_already_buffered_bytes() {
let (metadata, _) = build_three_rg_file();
let rg0 = column_range(&metadata, 0);
let rg1 = column_range(&metadata, 1);
let rg_plan = (0..3).map(|rg_index| RgPlanEntry { rg_index }).collect();
let requested = rg0.end - rg0.start;
let next = rg1.end - rg1.start;
let budget = (requested + next) as usize;
let mut prefetched = HashSet::new();

let ranges = prefetch_row_group_ranges(
vec![rg0.clone()],
1,
Some(budget),
&rg_plan,
&ProjectionMask::all(),
&metadata,
&mut prefetched,
);

assert_eq!(ranges, vec![rg0]);
assert!(prefetched.is_empty());
}

#[test]
fn prefetch_disabled_leaves_request_unchanged() {
let (metadata, _) = build_three_rg_file();
let rg0 = column_range(&metadata, 0);
let rg_plan = (0..3).map(|rg_index| RgPlanEntry { rg_index }).collect();
let mut prefetched = HashSet::new();

let ranges = prefetch_row_group_ranges(
vec![rg0.clone()],
0,
None,
&rg_plan,
&ProjectionMask::all(),
&metadata,
&mut prefetched,
);

assert_eq!(ranges, vec![rg0]);
assert!(prefetched.is_empty());
}

/// Create a fresh `(creation_errors, evaluation_errors)` counter pair
Expand Down
1 change: 1 addition & 0 deletions datafusion/datasource-parquet/src/source.rs
Original file line number Diff line number Diff line change
Expand Up @@ -631,6 +631,7 @@ impl FileSource for ParquetSource {
predicate: self.predicate.clone(),
table_schema: self.table_schema.clone(),
metadata_size_hint: self.metadata_size_hint,
prefetch_size: self.table_parquet_options.global.prefetch_size,
metrics: self.metrics().clone(),
parquet_file_reader_factory,
pushdown_filters: self.pushdown_filters(),
Expand Down
6 changes: 5 additions & 1 deletion datafusion/proto-common/proto/datafusion_common.proto
Original file line number Diff line number Diff line change
Expand Up @@ -631,6 +631,10 @@ message ParquetOptions {
uint64 max_row_group_bytes = 37;
}

oneof prefetch_size_opt {
uint64 prefetch_size = 38;
}

ParquetCdcOptions content_defined_chunking = 35;

// Optional timezone applied to INT96-coerced timestamps when `coerce_int96`
Expand Down Expand Up @@ -715,4 +719,4 @@ enum MetricCategory {
message ExplainAnalyzeCategoriesNode {
bool all = 1;
repeated MetricCategory only = 2;
}
}
Loading
Loading