Skip to content
Draft
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
67 changes: 45 additions & 22 deletions datafusion/datasource-parquet/src/opener/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ use crate::{
};
use arrow::array::RecordBatch;
use arrow::datatypes::DataType;
use datafusion_datasource::file_stream::read_ahead::ReadAheadReservation;
use datafusion_datasource::morsel::{Morsel, MorselPlan, MorselPlanner, Morselizer};
use datafusion_physical_expr::projection::ProjectionExprs;
use datafusion_physical_expr_adapter::replace_columns_with_literals;
Expand Down Expand Up @@ -397,6 +398,7 @@ enum ParquetOpenState {
BuildStream(Box<RowGroupsPrunedParquetOpen>),
/// Terminal state: the final opened stream is ready to return.
Ready(BoxStream<'static, Result<RecordBatch>>),
LoadInitialData(BoxFuture<'static, Result<BoxStream<'static, Result<RecordBatch>>>>),
/// Terminal state: reading complete
Done,
}
Expand All @@ -416,6 +418,7 @@ impl fmt::Debug for ParquetOpenState {
ParquetOpenState::PruneWithBloomFilters(_) => "PruneWithBloomFilters",
ParquetOpenState::BuildStream(_) => "BuildStream",
ParquetOpenState::Ready(_) => "Ready",
ParquetOpenState::LoadInitialData(_) => "LoadInitialData",
ParquetOpenState::Done => "Done",
};
f.write_str(state)
Expand Down Expand Up @@ -603,10 +606,11 @@ impl ParquetOpenState {
ParquetOpenState::PruneWithBloomFilters(loaded) => Ok(
ParquetOpenState::BuildStream(Box::new(loaded.prune_bloom_filters())),
),
ParquetOpenState::BuildStream(prepared) => {
Ok(ParquetOpenState::Ready(prepared.build_stream()?))
}
ParquetOpenState::BuildStream(prepared) => prepared.build_stream(),
ParquetOpenState::Ready(stream) => Ok(ParquetOpenState::Ready(stream)),
ParquetOpenState::LoadInitialData(future) => {
Ok(ParquetOpenState::LoadInitialData(future))
}
ParquetOpenState::Done => {
panic!("ParquetOpenFuture polled after completion");
}
Expand Down Expand Up @@ -722,6 +726,11 @@ impl MorselPlanner for ParquetMorselPlanner {
)))
})))
}
ParquetOpenState::LoadInitialData(future) => {
Ok(Some(Self::schedule_io(async move {
Ok(ParquetOpenState::Ready(future.await?))
})))
}
ParquetOpenState::Ready(stream) => {
let morsels: Vec<Box<dyn Morsel>> =
vec![Box::new(ParquetStreamMorsel::new(stream))];
Expand Down Expand Up @@ -1364,7 +1373,7 @@ impl BloomFiltersLoadedParquetOpen {

impl RowGroupsPrunedParquetOpen {
/// Build the final parquet stream once all pruning work is complete.
fn build_stream(self) -> Result<BoxStream<'static, Result<RecordBatch>>> {
fn build_stream(self) -> Result<ParquetOpenState> {
let RowGroupsPrunedParquetOpen {
prepared,
mut row_groups,
Expand Down Expand Up @@ -1700,7 +1709,8 @@ impl RowGroupsPrunedParquetOpen {
.file_metrics
.row_groups_pruned_dynamic_filter
.clone();
let stream = PushDecoderStreamState {
let read_ahead = prepared.extensions.get_arc::<ReadAheadReservation>();
let state = PushDecoderStreamState {
decoder: Some(decoder),
active_reader: None,
rg_plan,
Expand All @@ -1716,6 +1726,7 @@ impl RowGroupsPrunedParquetOpen {
prepared.partition_index,
),
prefetch_reservation: None,
initial_read_ahead: None,
decoder_projection,
arrow_reader_metrics,
predicate_cache_inner_records,
Expand All @@ -1727,24 +1738,36 @@ impl RowGroupsPrunedParquetOpen {
filter_installed,
row_filter_skipped_fully_matched,
byte_progress,
}
.into_stream();

// Wrap the stream so a dynamic filter can stop the file scan early, but
// only when the pruner is still watching a filter that can change
// mid-scan. For a static (or already-complete) predicate the up-front
// `prune_file` check already captured everything that can be pruned, so
// per-batch re-checking would only add overhead.
match prepared.file_pruner {
Some(file_pruner) if file_pruner.is_watching() => {
Ok(EarlyStoppingStream::new(
stream,
file_pruner,
files_ranges_pruned_statistics,
)
.boxed())
};

let wrap = move |stream| {
// Wrap the stream so a dynamic filter can stop the file scan early, but
// only when the pruner is still watching a filter that can change
// mid-scan. For a static (or already-complete) predicate the up-front
// `prune_file` check already captured everything that can be pruned, so
// per-batch re-checking would only add overhead.
match prepared.file_pruner {
Some(file_pruner) if file_pruner.is_watching() => {
Ok(EarlyStoppingStream::new(
stream,
file_pruner,
files_ranges_pruned_statistics,
)
.boxed())
}
_ => Ok(stream),
}
_ => Ok(stream),
};
if let Some(reservation) = read_ahead {
Ok(ParquetOpenState::LoadInitialData(
async move {
let state = state.prepare_initial_data(reservation).await?;
wrap(state.into_stream())
}
.boxed(),
))
} else {
Ok(ParquetOpenState::Ready(wrap(state.into_stream())?))
}
}
}
Expand Down
164 changes: 164 additions & 0 deletions datafusion/datasource-parquet/src/push_decoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@

use bytes::Bytes;
use datafusion_common_runtime::SpawnedTask;
use datafusion_datasource::file_stream::read_ahead::ReadAheadReservation;
use datafusion_execution::memory_pool::{MemoryConsumer, MemoryPool, MemoryReservation};
use std::collections::VecDeque;
use std::ops::Range;
Expand Down Expand Up @@ -414,6 +415,7 @@ pub(crate) struct PushDecoderStreamState {
pub(crate) pending_prefetch: Option<PrefetchedRowGroup>,
pub(crate) prefetch_metrics: crate::metrics::PrefetchMetrics,
pub(crate) prefetch_reservation: Option<MemoryReservation>,
pub(crate) initial_read_ahead: Option<Arc<ReadAheadReservation>>,
/// 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 @@ -540,6 +542,59 @@ impl RowFilterContext {
}

impl PushDecoderStreamState {
pub(crate) async fn prepare_initial_data(
mut self,
reservation: Arc<ReadAheadReservation>,
) -> Result<Self> {
self.initial_read_ahead = Some(Arc::clone(&reservation));
let Some(entry) = self.rg_plan.front() else {
return Ok(self);
};
let group = entry.rg_index;
let Some(ranges) = column_ranges(
&self.parquet_metadata,
entry,
if entry.fully_matched {
self.decoder_projection.projection_mask()
} else {
&self.fetch_projection
},
) else {
return Ok(self);
};
let ranges = merge_ranges(ranges);
let bytes = ranges.iter().try_fold(0usize, |sum, r| {
sum.checked_add(usize::try_from(r.end - r.start).ok()?)
});
let Some(bytes) = bytes else {
return Ok(self);
};
if bytes == 0 || !reservation.try_resize(bytes) {
return Ok(self);
}
let result = self
.reader
.lock()
.await
.get_byte_ranges(ranges.clone())
.await;
match result {
Ok(data) => {
self.decoder
.as_mut()
.expect("decoder present")
.push_ranges(ranges, data)?;
reservation.record_read(bytes);
self.upfront_row_group = Some(group);
self.initial_read_ahead = Some(reservation);
}
Err(error) => {
debug!("Initial scan read-ahead failed, retrying on demand: {error}")
}
}
Ok(self)
}

/// Drive the state machine to completion as a [`futures::Stream`] of record batches.
///
/// The returned stream is fused and boxed so the caller can wrap it (for
Expand Down Expand Up @@ -788,6 +843,7 @@ impl PushDecoderStreamState {
self.byte_progress.credit(entry.bytes);
}
self.active_reader = Some(reader);
self.initial_read_ahead = None;
// The extracted reader now owns required bytes. Release any
// unused speculation (e.g. pages removed by a row filter).
if !self.progressive_io || self.prefetch_reservation.is_some() {
Expand Down Expand Up @@ -1447,6 +1503,114 @@ mod tests {
}
}

#[tokio::test]
async fn shared_read_ahead_preserves_rows_falls_back_and_cancels() {
use datafusion_datasource::file_scan_config::FileScanConfigBuilder;
use datafusion_datasource::source::DataSourceExec;
use datafusion_datasource::{PartitionedFile, file_groups::FileGroup};
use datafusion_execution::memory_pool::GreedyMemoryPool;
use datafusion_execution::{
TaskContext, config::SessionConfig, object_store::ObjectStoreUrl,
};
use datafusion_physical_plan::ExecutionPlan;

for (budget, limit) in [
(0, None),
(1, None),
(32 << 20, None),
(32 << 20, Some(123)),
] {
let pool: Arc<dyn MemoryPool> = Arc::new(GreedyMemoryPool::new(512 << 20));
let (data, metadata, schema) = build_three_rg_file_data();
let files = (0..4)
.map(|i| {
PartitionedFile::new(format!("queue{i}.parquet"), data.len() as u64)
})
.collect();
let control = Arc::new(ReadControl {
block_second: limit.is_some(),
..Default::default()
});
// RG1 is fully rejected inside arrow-rs, without a reader being returned.
let predicate = Arc::new(BinaryExpr::new(
Arc::new(BinaryExpr::new(
Arc::new(Column::new("v", 0)),
Operator::Modulo,
lit(2000i64),
)),
Operator::Lt,
lit(1000i64),
));
let source = crate::source::ParquetSource::new(schema)
.with_progressive_io(false)
.with_row_group_prefetch(1 << 20, Arc::clone(&pool))
.with_scan_read_ahead(4, budget, Arc::clone(&pool))
.with_enable_page_index(false)
.with_pushdown_filters(true)
.with_predicate(predicate)
.with_parquet_file_reader_factory(Arc::new(TestReader {
data,
metadata,
control,
}));
let config = FileScanConfigBuilder::new(
ObjectStoreUrl::local_filesystem(),
Arc::new(source),
)
.with_file_group(FileGroup::new(files))
.with_limit(limit)
.build();
let exec = DataSourceExec::new(Arc::new(config));
let task = TaskContext::default()
.with_session_config(SessionConfig::new().with_batch_size(100));
let mut stream = exec.execute(0, Arc::new(task)).unwrap();
let values = tokio::time::timeout(std::time::Duration::from_secs(5), async {
let mut values = Vec::new();
while let Some(batch) = stream.next().await {
values.extend_from_slice(
batch
.unwrap()
.column(0)
.as_any()
.downcast_ref::<Int64Array>()
.unwrap()
.values(),
);
}
values
})
.await
.unwrap();
if let Some(limit) = limit {
assert_eq!(values.len(), limit);
} else {
let mut expected: Vec<i64> =
(0..4).flat_map(|_| (0..1000).chain(2000..3000)).collect();
let mut values = values;
values.sort_unstable();
expected.sort_unstable();
assert_eq!(values, expected);
}
let metrics = exec.metrics().unwrap();
let metric = |name| {
metrics
.sum_by_name(name)
.map(|value| value.as_usize())
.unwrap_or(0)
};
if budget >= 32 << 20 {
assert!(metric("scan_read_ahead_jobs") > 0);
assert!(metric("scan_read_ahead_initial_bytes") > 0);
assert!(metric("scan_read_ahead_peak_bytes") <= budget);
} else {
assert_eq!(metric("scan_read_ahead_jobs"), 0);
}
drop(stream);
drop(exec);
assert_pool_released(&pool).await;
}
}

#[tokio::test]
async fn upfront_prefetch_starts_before_evaluating_row_filters() {
use datafusion_common::config::ConfigOptions;
Expand Down
26 changes: 26 additions & 0 deletions datafusion/datasource-parquet/src/source.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
// under the License.

//! ParquetSource implementation for reading parquet files
use datafusion_datasource::file_stream::read_ahead::ReadAheadBudget;
use std::fmt::Debug;
use std::fmt::Formatter;
use std::sync::Arc;
Expand Down Expand Up @@ -321,6 +322,7 @@ pub struct ParquetSource {
/// in the opener.
sort_order_for_reorder: Option<LexOrdering>,
row_group_prefetch: Option<RowGroupPrefetchOptions>,
scan_read_ahead: Option<Arc<ReadAheadBudget>>,
}

impl ParquetSource {
Expand Down Expand Up @@ -348,6 +350,7 @@ impl ParquetSource {
reverse_row_groups: false,
sort_order_for_reorder: None,
row_group_prefetch: None,
scan_read_ahead: None,
}
}

Expand Down Expand Up @@ -376,6 +379,22 @@ impl ParquetSource {
self
}

/// Experimental shared file-range queue, enabled only for reorderable sibling
/// streams. Prepares initial row-group bytes before workers claim each job.
/// A fixed cap bounds backlog bytes across the scan.
/// This execution-local option is not serialized and does not enable pushdown.
pub fn with_scan_read_ahead(
mut self,
max_jobs: usize,
max_bytes: usize,
memory_pool: Arc<dyn MemoryPool>,
) -> Self {
self.scan_read_ahead = (max_jobs > 0 && max_bytes > 0).then(|| {
ReadAheadBudget::new(max_jobs, max_bytes, memory_pool, &self.metrics)
});
self
}

/// Fetch pages progressively as decoding and row filtering require
/// them (the default). When false, the first demand read for each row group
/// fetches the output and predicate pages selected at file open together. This reduces
Expand Down Expand Up @@ -600,6 +619,13 @@ impl From<ParquetSource> for Arc<dyn FileSource> {
}

impl FileSource for ParquetSource {
fn read_ahead_budget(&self) -> Option<Arc<ReadAheadBudget>> {
// Preparing all projected columns early is the upfront-I/O policy.
(!self.table_parquet_options.global.progressive_io)
.then(|| self.scan_read_ahead.clone())
.flatten()
}

fn create_file_opener(
&self,
_object_store: Arc<dyn ObjectStore>,
Expand Down
Loading