diff --git a/datafusion/datasource-parquet/src/opener/mod.rs b/datafusion/datasource-parquet/src/opener/mod.rs index f6736152fe796..241fa9e058f35 100644 --- a/datafusion/datasource-parquet/src/opener/mod.rs +++ b/datafusion/datasource-parquet/src/opener/mod.rs @@ -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; @@ -397,6 +398,7 @@ enum ParquetOpenState { BuildStream(Box), /// Terminal state: the final opened stream is ready to return. Ready(BoxStream<'static, Result>), + LoadInitialData(BoxFuture<'static, Result>>>), /// Terminal state: reading complete Done, } @@ -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) @@ -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"); } @@ -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> = vec![Box::new(ParquetStreamMorsel::new(stream))]; @@ -1364,7 +1373,7 @@ impl BloomFiltersLoadedParquetOpen { impl RowGroupsPrunedParquetOpen { /// Build the final parquet stream once all pruning work is complete. - fn build_stream(self) -> Result>> { + fn build_stream(self) -> Result { let RowGroupsPrunedParquetOpen { prepared, mut row_groups, @@ -1700,7 +1709,8 @@ impl RowGroupsPrunedParquetOpen { .file_metrics .row_groups_pruned_dynamic_filter .clone(); - let stream = PushDecoderStreamState { + let read_ahead = prepared.extensions.get_arc::(); + let state = PushDecoderStreamState { decoder: Some(decoder), active_reader: None, rg_plan, @@ -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, @@ -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())?)) } } } diff --git a/datafusion/datasource-parquet/src/push_decoder.rs b/datafusion/datasource-parquet/src/push_decoder.rs index a63a3922904a0..6a1a4acff8389 100644 --- a/datafusion/datasource-parquet/src/push_decoder.rs +++ b/datafusion/datasource-parquet/src/push_decoder.rs @@ -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; @@ -414,6 +415,7 @@ pub(crate) struct PushDecoderStreamState { pub(crate) pending_prefetch: Option, pub(crate) prefetch_metrics: crate::metrics::PrefetchMetrics, pub(crate) prefetch_reservation: Option, + pub(crate) initial_read_ahead: Option>, /// Per-file projection: the mask installed on every decoder and the /// per-batch transform applied by [`Self::project_batch`]. pub(crate) decoder_projection: DecoderProjection, @@ -540,6 +542,59 @@ impl RowFilterContext { } impl PushDecoderStreamState { + pub(crate) async fn prepare_initial_data( + mut self, + reservation: Arc, + ) -> Result { + 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 @@ -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() { @@ -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 = 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::() + .unwrap() + .values(), + ); + } + values + }) + .await + .unwrap(); + if let Some(limit) = limit { + assert_eq!(values.len(), limit); + } else { + let mut expected: Vec = + (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; diff --git a/datafusion/datasource-parquet/src/source.rs b/datafusion/datasource-parquet/src/source.rs index 1add32d85c7ca..f42a6584db276 100644 --- a/datafusion/datasource-parquet/src/source.rs +++ b/datafusion/datasource-parquet/src/source.rs @@ -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; @@ -321,6 +322,7 @@ pub struct ParquetSource { /// in the opener. sort_order_for_reorder: Option, row_group_prefetch: Option, + scan_read_ahead: Option>, } impl ParquetSource { @@ -348,6 +350,7 @@ impl ParquetSource { reverse_row_groups: false, sort_order_for_reorder: None, row_group_prefetch: None, + scan_read_ahead: None, } } @@ -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, + ) -> 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 @@ -600,6 +619,13 @@ impl From for Arc { } impl FileSource for ParquetSource { + fn read_ahead_budget(&self) -> Option> { + // 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, diff --git a/datafusion/datasource/src/file.rs b/datafusion/datasource/src/file.rs index f1a94f2e12363..27f6593546434 100644 --- a/datafusion/datasource/src/file.rs +++ b/datafusion/datasource/src/file.rs @@ -92,6 +92,13 @@ pub trait FileSource: Any + Send + Sync { Ok(Box::new(FileOpenerMorselizer::new(opener))) } + /// Optional execution-local budget for shared file-range read-ahead. + fn read_ahead_budget( + &self, + ) -> Option> { + None + } + /// Returns the table schema for the overall table (including partition columns, if any) /// /// This method returns the unprojected schema: the full schema of the data diff --git a/datafusion/datasource/src/file_stream/mod.rs b/datafusion/datasource/src/file_stream/mod.rs index 619d2cf2ef479..7580ac8339a58 100644 --- a/datafusion/datasource/src/file_stream/mod.rs +++ b/datafusion/datasource/src/file_stream/mod.rs @@ -23,6 +23,7 @@ mod builder; mod metrics; +pub mod read_ahead; mod scan_state; pub(crate) mod work_source; diff --git a/datafusion/datasource/src/file_stream/read_ahead.rs b/datafusion/datasource/src/file_stream/read_ahead.rs new file mode 100644 index 0000000000000..24a993b9ef03d --- /dev/null +++ b/datafusion/datasource/src/file_stream/read_ahead.rs @@ -0,0 +1,379 @@ +// 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. + +//! Experimental shared queue for preparing upcoming file-range scan jobs. + +use std::collections::VecDeque; +use std::fmt; +use std::sync::Arc; +use std::task::{Context, Poll, Waker}; + +use arrow::record_batch::RecordBatch; +use datafusion_common::{DataFusionError, Result}; +use datafusion_common_runtime::JoinSet; +use datafusion_execution::memory_pool::{MemoryConsumer, MemoryPool, MemoryReservation}; +use datafusion_physical_plan::metrics::{ + Count, ExecutionPlanMetricsSet, Gauge, MetricBuilder, +}; +use futures::stream::BoxStream; +use futures::task::{ArcWake, waker_ref}; +use futures::{FutureExt, StreamExt}; +use parking_lot::Mutex; + +use crate::PartitionedFile; +use crate::morsel::Morselizer; + +const MIN_JOB_BYTES: usize = 16 * 1024 * 1024; + +/// Execution-local scan backlog budget. This prototype is opt-in and not serialized. +#[derive(Debug)] +pub struct ReadAheadBudget { + max_jobs: usize, + max_bytes: usize, + pool: Arc, + used: Mutex, + peak_bytes: Gauge, + peak_jobs: Gauge, + admitted: Count, + initial_bytes: Count, + denied: Count, + pub(crate) fallbacks: Count, +} + +impl ReadAheadBudget { + /// Construct a backlog bounded by job count and compressed bytes. + pub fn new( + max_jobs: usize, + max_bytes: usize, + pool: Arc, + metrics: &ExecutionPlanMetricsSet, + ) -> Arc { + Arc::new(Self { + max_jobs, + max_bytes, + pool, + used: Mutex::new(0), + peak_bytes: MetricBuilder::new(metrics) + .global_gauge("scan_read_ahead_peak_bytes"), + peak_jobs: MetricBuilder::new(metrics) + .global_gauge("scan_read_ahead_peak_jobs"), + admitted: MetricBuilder::new(metrics).global_counter("scan_read_ahead_jobs"), + initial_bytes: MetricBuilder::new(metrics) + .global_counter("scan_read_ahead_initial_bytes"), + denied: MetricBuilder::new(metrics) + .global_counter("scan_read_ahead_budget_denials"), + fallbacks: MetricBuilder::new(metrics) + .global_counter("scan_read_ahead_demand_fallbacks"), + }) + } + + fn reserve(self: &Arc) -> Option> { + let reservation = Arc::new(ReadAheadReservation { + budget: Arc::clone(self), + reservation: MemoryConsumer::new("Scan read-ahead backlog") + .register(&self.pool), + }); + reservation.try_resize(MIN_JOB_BYTES).then_some(reservation) + } +} + +/// Reservation carried in a file extension through planning and initial data I/O. +/// It is released when the first row group's reader takes ownership of the bytes. +#[derive(Debug)] +pub struct ReadAheadReservation { + budget: Arc, + reservation: MemoryReservation, +} + +impl ReadAheadReservation { + /// Record successfully prefetched initial payload bytes. + pub fn record_read(&self, bytes: usize) { + self.budget.initial_bytes.add(bytes); + } + + /// Account for projected compressed bytes before reading them. Failure leaves + /// the old reservation intact; the caller falls back to ordinary demand I/O. + pub fn try_resize(&self, bytes: usize) -> bool { + let bytes = bytes.max(MIN_JOB_BYTES); + let mut used = self.budget.used.lock(); + let Some(target) = used + .checked_sub(self.reservation.size()) + .and_then(|n| n.checked_add(bytes)) + else { + self.budget.denied.add(1); + return false; + }; + if target > self.budget.max_bytes || self.reservation.try_resize(bytes).is_err() { + self.budget.denied.add(1); + return false; + } + *used = target; + self.budget.peak_bytes.set_max(target); + true + } +} + +impl Drop for ReadAheadReservation { + fn drop(&mut self) { + let mut used = self.budget.used.lock(); + *used -= self.reservation.free(); + } +} + +#[derive(Default)] +struct Waiters(Mutex>); + +impl ArcWake for Waiters { + fn wake_by_ref(arc_self: &Arc) { + let waiters = std::mem::take(&mut *arc_self.0.lock()); + for waker in waiters { + waker.wake(); + } + } +} + +/// Completed tasks in JoinSet form the ready queue. Dropping the queue aborts +/// outstanding preparation and releases their reservations when cancellation runs. +pub(super) struct ReadAheadQueue { + budget: Arc, + jobs: Mutex>>>>, + waiters: Arc, +} + +impl fmt::Debug for ReadAheadQueue { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("ReadAheadQueue") + .field("budget", &self.budget) + .finish_non_exhaustive() + } +} + +impl ReadAheadQueue { + pub(super) fn new(budget: Arc) -> Self { + Self { + budget, + jobs: Mutex::new(JoinSet::new()), + waiters: Arc::default(), + } + } + + pub(super) fn fill( + &self, + files: &Mutex>, + morselizer: &Arc, + ) { + let mut jobs = self.jobs.lock(); + while jobs.len() < self.budget.max_jobs { + if files.lock().is_empty() { + break; + } + let Some(reservation) = self.budget.reserve() else { + break; + }; + let Some(mut file) = files.lock().pop_front() else { + break; + }; + file.extensions.insert_arc(reservation); + let morselizer = Arc::clone(morselizer); + jobs.spawn(async move { prepare_file(morselizer, file).await }); + self.budget.admitted.add(1); + self.budget.peak_jobs.set_max(jobs.len()); + } + } + + pub(super) fn demand_fallback(&self) { + self.budget.fallbacks.add(1); + } + + pub(super) fn poll_ready( + &self, + cx: &mut Context<'_>, + ) -> Poll>>>> { + { + let mut waiters = self.waiters.0.lock(); + if !waiters.iter().any(|w| w.will_wake(cx.waker())) { + waiters.push(cx.waker().clone()); + } + } + let waker = waker_ref(&self.waiters); + let mut context = Context::from_waker(&waker); + let mut jobs = self.jobs.lock(); + std::pin::pin!(jobs.join_next()) + .poll_unpin(&mut context) + .map(|result| { + result.map(|joined| { + joined.unwrap_or_else(|error| { + Err(DataFusionError::External(Box::new(error))) + }) + }) + }) + } +} + +async fn prepare_file( + morselizer: Arc, + file: PartitionedFile, +) -> Result>> { + let mut planners = VecDeque::from([morselizer.plan_file(file)?]); + let mut morsels = Vec::new(); + while let Some(planner) = planners.pop_front() { + if let Some(mut plan) = planner.plan()? { + morsels.extend(plan.take_morsels()); + let mut children = plan.take_ready_planners(); + if let Some(pending) = plan.take_pending_planner() { + children.push(pending.await?); + } + for child in children.into_iter().rev() { + planners.push_front(child); + } + } + tokio::task::yield_now().await; + } + Ok(futures::stream::iter(morsels) + .flat_map(|m| m.into_stream()) + .boxed()) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::morsel::{Morsel, MorselPlan, MorselPlanner}; + use datafusion_execution::memory_pool::GreedyMemoryPool; + use tokio::sync::oneshot; + + #[test] + fn budget_caps_backlog_and_releases_reservations() { + let pool: Arc = Arc::new(GreedyMemoryPool::new(128 << 20)); + let budget = ReadAheadBudget::new( + 2, + 32 << 20, + Arc::clone(&pool), + &ExecutionPlanMetricsSet::new(), + ); + let first = budget.reserve().unwrap(); + let second = budget.reserve().unwrap(); + assert!(budget.reserve().is_none()); + assert!(!first.try_resize(32 << 20)); + assert_eq!(pool.reserved(), 32 << 20); + drop(second); + assert!(first.try_resize(32 << 20)); + drop(first); + assert_eq!(pool.reserved(), 0); + assert_eq!(*budget.used.lock(), 0); + } + + #[derive(Debug)] + struct GatedMorselizer(Mutex>>); + + impl Morselizer for GatedMorselizer { + fn plan_file(&self, file: PartitionedFile) -> Result> { + Ok(Box::new(GatedPlanner { + file, + gate: self.0.lock().pop_front(), + })) + } + } + + #[derive(Debug)] + struct GatedPlanner { + file: PartitionedFile, + gate: Option>, + } + + impl MorselPlanner for GatedPlanner { + fn plan(self: Box) -> Result> { + let Self { file, gate } = *self; + if let Some(gate) = gate { + Ok(Some(MorselPlan::new().with_pending_planner(async move { + gate.await + .map_err(|e| DataFusionError::External(Box::new(e)))?; + Ok(Box::new(Self { file, gate: None }) as Box) + }))) + } else { + Ok(Some( + MorselPlan::new().with_morsels(vec![Box::new(ReadyFile(file))]), + )) + } + } + } + + #[derive(Debug)] + struct ReadyFile(PartitionedFile); + + impl Morsel for ReadyFile { + fn into_stream(self: Box) -> BoxStream<'static, Result> { + futures::stream::once(async move { + let batch = + RecordBatch::new_empty(Arc::new(arrow::datatypes::Schema::empty())); + drop(self.0); + Ok(batch) + }) + .boxed() + } + } + + #[tokio::test] + async fn queue_returns_ready_work_and_cancels_remaining_io() { + let pool: Arc = Arc::new(GreedyMemoryPool::new(128 << 20)); + let budget = ReadAheadBudget::new( + 2, + 32 << 20, + Arc::clone(&pool), + &ExecutionPlanMetricsSet::new(), + ); + let queue = ReadAheadQueue::new(Arc::clone(&budget)); + let (first_tx, first_rx) = oneshot::channel(); + let (second_tx, second_rx) = oneshot::channel(); + let morselizer: Arc = + Arc::new(GatedMorselizer(Mutex::new(VecDeque::from([ + first_rx, second_rx, + ])))); + let files = Mutex::new(VecDeque::from([ + PartitionedFile::new("first", 1), + PartitionedFile::new("second", 1), + ])); + queue.fill(&files, &morselizer); + assert!(files.lock().is_empty()); + assert_eq!(pool.reserved(), 32 << 20); + second_tx.send(()).unwrap(); + let mut ready = tokio::time::timeout( + std::time::Duration::from_secs(5), + futures::future::poll_fn(|cx| queue.poll_ready(cx)), + ) + .await + .unwrap() + .unwrap() + .unwrap(); + assert!(ready.next().await.unwrap().is_ok()); + drop(ready); + assert_eq!(pool.reserved(), 16 << 20); + assert!( + futures::future::poll_fn(|cx| queue.poll_ready(cx)) + .now_or_never() + .is_none() + ); + drop(queue); + tokio::time::timeout(std::time::Duration::from_secs(5), async { + while pool.reserved() != 0 { + tokio::task::yield_now().await; + } + }) + .await + .unwrap(); + assert!(first_tx.send(()).is_err()); + } +} diff --git a/datafusion/datasource/src/file_stream/scan_state.rs b/datafusion/datasource/src/file_stream/scan_state.rs index d1e57604f971c..e6129fd57d2eb 100644 --- a/datafusion/datasource/src/file_stream/scan_state.rs +++ b/datafusion/datasource/src/file_stream/scan_state.rs @@ -17,6 +17,7 @@ use datafusion_common::internal_datafusion_err; use std::collections::VecDeque; +use std::sync::Arc; use std::task::{Context, Poll}; use crate::morsel::{Morsel, MorselPlanner, Morselizer, PendingMorselPlanner}; @@ -67,7 +68,7 @@ pub(super) struct ScanState { /// Remaining row limit, if any. remain: Option, /// The morselizer used to plan files. - morselizer: Box, + morselizer: Arc, /// Behavior if opening or scanning a file fails. on_error: OnError, /// CPU-ready planners for the current file. @@ -96,7 +97,7 @@ impl ScanState { Self { work_source, remain, - morselizer, + morselizer: morselizer.into(), on_error, ready_planners: Default::default(), ready_morsels: Default::default(), @@ -126,6 +127,12 @@ impl ScanState { let _processing_timer: ScopedTimerGuard<'_> = self.metrics.time_processing.timer(); + if self.reader.is_none() + && let WorkSource::Shared(shared) = &self.work_source + { + shared.fill_read_ahead(&self.morselizer); + } + // Try and resolve outstanding IO first. If it is still pending, check // the current reader or ready morsels before yielding. New planning // work must still wait for this I/O to resolve. @@ -265,6 +272,32 @@ impl ScanState { }; } + if let WorkSource::Shared(shared) = &self.work_source + && let Some(queue) = shared.read_ahead() + { + match queue.poll_ready(cx) { + Poll::Ready(Some(Ok(stream))) => { + self.metrics.files_opened.add(1); + self.metrics.time_scanning_total.start(); + self.metrics.time_scanning_until_data.start(); + self.reader = Some(stream); + return ScanAndReturn::Continue; + } + Poll::Ready(Some(Err(error))) => { + self.metrics.file_open_errors.add(1); + return match self.on_error { + OnError::Skip => { + self.metrics.files_processed.add(1); + ScanAndReturn::Continue + } + OnError::Fail => ScanAndReturn::Error(error), + }; + } + Poll::Pending => return ScanAndReturn::Return(Poll::Pending), + Poll::Ready(None) => {} // No backlog: ordinary demand planning guarantees progress. + } + } + // No outstanding work remains, so begin planning the next unopened file. let Some(part_file) = self.work_source.pop_front() else { return ScanAndReturn::Done(None); diff --git a/datafusion/datasource/src/file_stream/work_source.rs b/datafusion/datasource/src/file_stream/work_source.rs index c00048453b304..a0154050d68f8 100644 --- a/datafusion/datasource/src/file_stream/work_source.rs +++ b/datafusion/datasource/src/file_stream/work_source.rs @@ -71,19 +71,10 @@ pub(crate) struct SharedWorkSource { #[derive(Debug, Default)] pub(super) struct SharedWorkSourceInner { files: Mutex>, + read_ahead: Option, } impl SharedWorkSource { - /// Create a shared work source containing the provided unopened files. - pub(crate) fn new(files: impl IntoIterator) -> Self { - let files = files.into_iter().collect(); - Self { - inner: Arc::new(SharedWorkSourceInner { - files: Mutex::new(files), - }), - } - } - /// Create a shared work source for the unopened files in `config`. /// /// Files are reordered by the file source (e.g. by statistics for TopK) @@ -97,13 +88,40 @@ impl SharedWorkSource { .cloned() .collect(); let files = config.file_source.reorder_files(files); - Self::new(files) + Self { + inner: Arc::new(SharedWorkSourceInner { + files: Mutex::new(files.into()), + read_ahead: config + .file_source + .read_ahead_budget() + .map(super::read_ahead::ReadAheadQueue::new), + }), + } + } + + pub(super) fn read_ahead(&self) -> Option<&super::read_ahead::ReadAheadQueue> { + self.inner.read_ahead.as_ref() + } + + pub(super) fn fill_read_ahead( + &self, + morselizer: &Arc, + ) { + if let Some(queue) = self.read_ahead() { + queue.fill(&self.inner.files, morselizer); + } } /// Pop the next file from the shared work queue. /// /// Returns `None` if the queue is empty fn pop_front(&self) -> Option { - self.inner.files.lock().pop_front() + let file = self.inner.files.lock().pop_front(); + if file.is_some() + && let Some(queue) = self.read_ahead() + { + queue.demand_fallback(); + } + file } }