diff --git a/Cargo.lock b/Cargo.lock index db3052e5132..0a35c84a451 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7624,7 +7624,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "009a43e88f0f7b3b3c2c9de5384a0ac2cef8f33996b16f9e6738e24661618775" dependencies = [ "async-trait", - "asyncband 0.7.2", + "asyncband 0.7.3", "bytes", "chrono", "futures", @@ -7726,7 +7726,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ddff5c648dcb3922c2bec71e0d26e881ddf9c367cb7c9a59ed7f2d39a4f38218" dependencies = [ "anyhow", - "asyncband 0.7.2", + "asyncband 0.7.3", "base64 0.23.1", "bytes", "futures", @@ -7765,7 +7765,7 @@ version = "0.59.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b454e68e5f6b331bd079a2ee02ebe05078fc0ddaae6822767ef54ab3c2ac0aa8" dependencies = [ - "asyncband 0.7.2", + "asyncband 0.7.3", "futures", "http", "opendal-core 0.59.3", diff --git a/vortex-datafusion/src/persistent/opener.rs b/vortex-datafusion/src/persistent/opener.rs index 9858aedc809..a9dbe082277 100644 --- a/vortex-datafusion/src/persistent/opener.rs +++ b/vortex-datafusion/src/persistent/opener.rs @@ -46,6 +46,7 @@ use vortex::file::OpenOptionsSessionExt; use vortex::io::InstrumentedReadAt; use vortex::layout::LayoutReader; use vortex::layout::scan::scan_builder::ScanBuilder; +use vortex::layout::scan::v2; use vortex::metrics::Label; use vortex::metrics::MetricsRegistry; use vortex::session::VortexSession; @@ -430,7 +431,7 @@ impl FileOpener for VortexOpener { } let stream_target_field = Field::new_struct("", stream_schema.fields().clone(), false); - let stream = scan_builder + let scan_builder = scan_builder .with_metrics_registry(metrics_registry) .with_ordered(has_output_ordering) .map(move |chunk| { @@ -442,8 +443,15 @@ impl FileOpener for VortexOpener { &mut ctx, )?; Ok(RecordBatch::from(arrow.as_struct().clone())) - }) - .into_stream() + }); + let batches = if v2::enabled() { + v2::ScanBuilder::from_default(scan_builder) + .into_stream() + .map(|s| s.boxed()) + } else { + scan_builder.into_stream().map(|s| s.boxed()) + }; + let stream = batches .map_err(|e| exec_datafusion_err!("Failed to create Vortex stream: {e}"))? .map_err(move |e: VortexError| { DataFusionError::External(Box::new(e.with_context(format!( diff --git a/vortex-duckdb/src/file_reader.rs b/vortex-duckdb/src/file_reader.rs index f4c9bb67995..34d5509b106 100644 --- a/vortex-duckdb/src/file_reader.rs +++ b/vortex-duckdb/src/file_reader.rs @@ -26,6 +26,7 @@ use vortex::io::runtime::BlockingRuntime as _; use vortex::io::std_file::StdFileSystem; use vortex::layout::LayoutReaderRef; use vortex::layout::scan::scan_builder::ScanBuilder; +use vortex::layout::scan::v2; use vortex::mask::Mask; use vortex::session::SessionExt as _; @@ -178,8 +179,13 @@ pub fn reader_initialize(file: &mut OpenFileReader, global: &GlobalState) -> Vor .with_projection(global.projection.clone()) .with_some_filter(filter.filter.clone()) .with_selection(filter.row_selection.clone()); - let scan = builder.prepare()?; - let mut splits = scan.execute(filter.row_range.clone())?; + let mut splits = if v2::enabled() { + v2::ScanBuilder::from_default(builder) + .prepare()? + .execute(filter.row_range.clone())? + } else { + builder.prepare()?.execute(filter.row_range.clone())? + }; // threads take last element of file.splits so we need to reverse splits.reverse(); diff --git a/vortex-layout/src/scan/mod.rs b/vortex-layout/src/scan/mod.rs index 98fd1918a42..ab003641eb5 100644 --- a/vortex-layout/src/scan/mod.rs +++ b/vortex-layout/src/scan/mod.rs @@ -12,6 +12,7 @@ mod splits; mod tasks; #[cfg(test)] mod test; +pub mod v2; /// A heuristic for an ideal split size. /// diff --git a/vortex-layout/src/scan/scan_builder.rs b/vortex-layout/src/scan/scan_builder.rs index 3f7e1cd61aa..ec640d7002a 100644 --- a/vortex-layout/src/scan/scan_builder.rs +++ b/vortex-layout/src/scan/scan_builder.rs @@ -297,6 +297,28 @@ impl ScanBuilder { } } + /// Moves this builder's configuration into the alternative scan builder. + // TODO(joe): Remove once the V2 migration is complete. + pub(crate) fn into_parts(self) -> ScanParts { + ScanParts { + session: self.session, + layout_reader: self.layout_reader, + projection: self.projection, + filter: self.filter, + ordered: self.ordered, + row_range: self.row_range, + selection: self.selection, + split_by: self.split_by, + natural_splits: self.natural_splits, + concurrency: self.concurrency, + map_fn: self.map_fn, + metrics_registry: self.metrics_registry, + file_stats: self.file_stats, + limit: self.limit, + row_offset: self.row_offset, + } + } + /// Optimize expressions, compute split ranges, and return an executable repeated scan. pub fn prepare(self) -> VortexResult> { let dtype = self.dtype()?; @@ -390,6 +412,26 @@ impl ScanBuilder { } } +/// Configuration transferred to the alternative scan builder. +// TODO(joe): Remove once the V2 migration is complete. +pub(crate) struct ScanParts { + pub(crate) session: VortexSession, + pub(crate) layout_reader: LayoutReaderRef, + pub(crate) projection: BoundExpression, + pub(crate) filter: Option, + pub(crate) ordered: bool, + pub(crate) row_range: Option>, + pub(crate) selection: Selection, + pub(crate) split_by: SplitBy, + pub(crate) natural_splits: Option>, + pub(crate) concurrency: usize, + pub(crate) map_fn: Arc VortexResult + Send + Sync>, + pub(crate) metrics_registry: Option>, + pub(crate) file_stats: Option>, + pub(crate) limit: Option, + pub(crate) row_offset: u64, +} + enum LazyScanState { Builder(Option>>), Preparing(PreparingScan), diff --git a/vortex-layout/src/scan/v2/mod.rs b/vortex-layout/src/scan/v2/mod.rs new file mode 100644 index 00000000000..f0465f4332b --- /dev/null +++ b/vortex-layout/src/scan/v2/mod.rs @@ -0,0 +1,33 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! An alternative scan builder and execution path for incremental scan development. +//! +//! DuckDB and DataFusion select this path when `VORTEX_SCAN_V2=1`. It initially copies the +//! default builder, preparation, and streaming behavior and shares the existing split tasks. +//! Subsequent changes can replace its execution without changing the default scan path. +//! +//! [`ScanBuilder::from_default`] transfers an already configured builder into this path. +//! Callers can also construct [`ScanBuilder`] directly without setting the environment variable. + +mod repeated_scan; +mod scan_builder; + +use std::env; +use std::sync::LazyLock; + +pub use repeated_scan::RepeatedScanV2; +pub use scan_builder::ScanBuilder; + +/// The environment variable that selects this executor when set to `1`. +pub const ENV_VAR: &str = "VORTEX_SCAN_V2"; + +static ENABLED: LazyLock = LazyLock::new(|| env::var(ENV_VAR).is_ok_and(|v| v == "1")); + +/// Whether [`ENV_VAR`] selects this executor. Read once per process. +pub fn enabled() -> bool { + *ENABLED +} + +#[cfg(test)] +mod tests; diff --git a/vortex-layout/src/scan/v2/repeated_scan.rs b/vortex-layout/src/scan/v2/repeated_scan.rs new file mode 100644 index 00000000000..15ee2d99a02 --- /dev/null +++ b/vortex-layout/src/scan/v2/repeated_scan.rs @@ -0,0 +1,231 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use std::cmp; +use std::iter; +use std::ops::Range; +use std::sync::Arc; + +use futures::Stream; +use futures::StreamExt; +use futures::future::BoxFuture; +use itertools::Either; +use itertools::Itertools; +use vortex_array::ArrayRef; +use vortex_array::dtype::DType; +use vortex_array::expr::BoundExpression; +use vortex_array::iter::ArrayIterator; +use vortex_array::iter::ArrayIteratorAdapter; +use vortex_array::stream::ArrayStream; +use vortex_array::stream::ArrayStreamAdapter; +use vortex_error::VortexExpect; +use vortex_error::VortexResult; +use vortex_io::runtime::BlockingRuntime; +use vortex_io::session::RuntimeSessionExt; +use vortex_scan::selection::Selection; +use vortex_session::VortexSession; +use vortex_utils::parallelism::get_available_parallelism; + +use crate::LayoutReaderRef; +use crate::scan::filter::FilterExpr; +use crate::scan::splits::Splits; +use crate::scan::tasks::TaskContext; +use crate::scan::tasks::split_exec; + +/// A projected subset (by indices, range, and filter) of rows from a Vortex data source. +/// +/// Executes multiple row ranges of the same data source, optionally concurrently. +pub struct RepeatedScanV2 { + session: VortexSession, + layout_reader: LayoutReaderRef, + projection: BoundExpression, + filter: Option, + ordered: bool, + /// Optionally read a subset of the rows in the file. + row_range: Option>, + /// The selection mask to apply to the selected row range. + selection: Selection, + /// The natural splits of the file. + splits: Splits, + /// The number of splits to make progress on concurrently **per-thread**. + concurrency: usize, + /// Function to apply to each [`ArrayRef`] within the spawned split tasks. + map_fn: Arc VortexResult + Send + Sync>, + /// Maximal number of rows to read (after filtering) + limit: Option, + /// The dtype of the projected arrays. + dtype: DType, +} + +impl RepeatedScanV2 { + /// The dtype of the projected arrays. + pub fn dtype(&self) -> &DType { + &self.dtype + } + + /// Execute a row range as a blocking array iterator. + pub fn execute_array_iter( + &self, + row_range: Option>, + runtime: &B, + ) -> VortexResult { + let dtype = self.dtype.clone(); + let stream = self.execute_stream(row_range)?; + let iter = runtime.block_on_stream(stream); + Ok(ArrayIteratorAdapter::new(dtype, iter)) + } + + /// Execute a row range as an array stream. + pub fn execute_array_stream( + &self, + row_range: Option>, + ) -> VortexResult { + let dtype = self.dtype.clone(); + let stream = self.execute_stream(row_range)?; + Ok(ArrayStreamAdapter::new(dtype, stream)) + } +} + +impl RepeatedScanV2 { + /// Constructor just to allow `scan_builder` to create a `RepeatedScanV2`. + #[expect( + clippy::too_many_arguments, + reason = "all arguments are needed for scan construction" + )] + pub(super) fn new( + session: VortexSession, + layout_reader: LayoutReaderRef, + projection: BoundExpression, + filter: Option, + ordered: bool, + row_range: Option>, + selection: Selection, + splits: Splits, + concurrency: usize, + map_fn: Arc VortexResult + Send + Sync>, + limit: Option, + dtype: DType, + ) -> Self { + Self { + session, + layout_reader, + projection, + filter, + ordered, + row_range, + selection, + splits, + concurrency, + map_fn, + limit, + dtype, + } + } + + /// Return one task per selected split of the requested row range. + pub fn execute( + &self, + row_range: Option>, + ) -> VortexResult>>>> { + let selection_range: Option> = match &self.selection { + Selection::IncludeByIndex(buf) if !buf.is_empty() => { + Some(buf[0]..buf[buf.len() - 1] + 1) + } + Selection::IncludeRoaring(roaring) if !roaring.is_empty() => { + Some(roaring.min().vortex_expect("empty")..roaring.max().vortex_expect("empty") + 1) + } + _ => None, + }; + let row_range = intersect_ranges(self.row_range.as_ref(), row_range); + let row_range = intersect_ranges(row_range.as_ref(), selection_range); + + let ranges = match &self.splits { + Splits::Natural(vec) => { + debug_assert!(vec.is_sorted()); + let splits_iter = match row_range { + None => Either::Left(vec.iter().copied()), + Some(range) => { + if range.is_empty() { + return Ok(Vec::new()); + } + let lo = vec.partition_point(|&x| x <= range.start); + let hi = vec.partition_point(|&x| x < range.end); + Either::Right( + iter::once(range.start) + .chain(vec[lo..hi].iter().copied()) + .chain(iter::once(range.end)), + ) + } + }; + + Either::Left(splits_iter.tuple_windows().map(|(start, end)| start..end)) + } + Splits::Ranges(ranges) => Either::Right(match row_range { + None => Either::Left(ranges.iter().cloned()), + Some(range) => { + if range.is_empty() { + return Ok(Vec::new()); + } + Either::Right(ranges.iter().filter_map(move |r| { + let start = cmp::max(r.start, range.start); + let end = cmp::min(r.end, range.end); + (start < end).then_some(start..end) + })) + } + }), + }; + + let mut limit = self.limit; + let mut tasks = Vec::new(); + let ctx = Arc::new(TaskContext { + filter: self.filter.clone().map(|f| Arc::new(FilterExpr::new(f))), + reader: Arc::clone(&self.layout_reader), + projection: self.projection.clone(), + mapper: Arc::clone(&self.map_fn), + }); + + for range in ranges { + let row_mask = self.selection.row_mask(&range); + if row_mask.mask().all_false() { + continue; + } + + tasks.push(split_exec(Arc::clone(&ctx), row_mask, limit.as_mut())?); + if limit.is_some_and(|l| l == 0) { + break; + } + } + + Ok(tasks) + } + + /// Execute a row range as a stream, honoring ordering and concurrency. + pub fn execute_stream( + &self, + row_range: Option>, + ) -> VortexResult> + Send + 'static + use> { + let num_workers = get_available_parallelism().unwrap_or(1); + let concurrency = self.concurrency * num_workers; + let handle = self.session.handle(); + + let stream = + futures::stream::iter(self.execute(row_range)?).map(move |task| handle.spawn(task)); + + let stream = if self.ordered { + stream.buffered(concurrency).boxed() + } else { + stream.buffer_unordered(concurrency).boxed() + }; + + Ok(stream.filter_map(|chunk| async move { chunk.transpose() })) + } +} + +fn intersect_ranges(left: Option<&Range>, right: Option>) -> Option> { + match (left, right) { + (None, None) => None, + (None, Some(r)) => Some(r), + (Some(l), None) => Some(l.clone()), + (Some(l), Some(r)) => Some(cmp::max(l.start, r.start)..cmp::min(l.end, r.end)), + } +} diff --git a/vortex-layout/src/scan/v2/scan_builder.rs b/vortex-layout/src/scan/v2/scan_builder.rs new file mode 100644 index 00000000000..daeaf191e9e --- /dev/null +++ b/vortex-layout/src/scan/v2/scan_builder.rs @@ -0,0 +1,487 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use std::ops::Range; +use std::pin::Pin; +use std::sync::Arc; +use std::task::Context; +use std::task::Poll; +use std::task::ready; + +use futures::Stream; +use futures::StreamExt; +use futures::future::BoxFuture; +use futures::stream::BoxStream; +use vortex_array::ArrayRef; +use vortex_array::dtype::DType; +use vortex_array::expr::BoundExpression; +use vortex_array::iter::ArrayIterator; +use vortex_array::iter::ArrayIteratorAdapter; +use vortex_array::stats::StatsSet; +use vortex_array::stream::ArrayStream; +use vortex_array::stream::ArrayStreamAdapter; +use vortex_error::VortexExpect; +use vortex_error::VortexResult; +use vortex_error::vortex_bail; +use vortex_io::runtime::BlockingRuntime; +use vortex_io::runtime::Handle; +use vortex_io::runtime::Task; +use vortex_io::session::RuntimeSessionExt; +use vortex_metrics::MetricsRegistry; +use vortex_scan::selection::Selection; +use vortex_scan::strict_sorted_buffer::StrictSortedBuffer; +use vortex_session::VortexSession; +use vortex_utils::parallelism::get_available_parallelism; + +use crate::LayoutReader; +use crate::LayoutReaderRef; +use crate::layouts::row_idx::RowIdx; +use crate::layouts::row_idx::RowIdxLayoutReader; +use crate::scan::scan_builder; +use crate::scan::scan_builder::referenced_field_masks; +use crate::scan::split_by::SplitBy; +use crate::scan::splits::Splits; +use crate::scan::splits::attempt_split_ranges; +use crate::scan::v2::RepeatedScanV2; + +/// Alternative builder for scanning a [`LayoutReader`] into arrays, streams, or iterators. +/// +/// This copy initially uses the same layout-reader execution as the default builder, so the new +/// scan implementation can be introduced incrementally behind [`enabled`](super::enabled). +/// +/// A scan has three independent row restriction mechanisms: +/// +/// - [`with_row_range`](Self::with_row_range) selects a contiguous range before scanning. +/// - [`with_selection`](Self::with_selection) applies a [`Selection`] inside that range. +/// - [`with_filter`](Self::with_filter) evaluates an expression predicate during execution. +/// +/// Projection and filter expressions must be bound against the reader dtype. Work is divided by +/// the configured [`SplitBy`] strategy or by explicit selection ranges. +pub struct ScanBuilder { + session: VortexSession, + layout_reader: LayoutReaderRef, + projection: BoundExpression, + filter: Option, + /// Whether the scan needs to return splits in the order they appear in the file. + ordered: bool, + /// Optionally read a subset of the rows in the file. + row_range: Option>, + /// The selection mask to apply to the selected row range. + selection: Selection, + /// How to split the file for concurrent processing. + split_by: SplitBy, + /// Precomputed full-file natural split boundaries; when set, [`prepare`](Self::prepare) + /// uses them instead of walking the layout. + natural_splits: Option>, + /// The number of splits to make progress on concurrently **per-thread**. + concurrency: usize, + /// Function to apply to each [`ArrayRef`] within the spawned split tasks. + map_fn: Arc VortexResult + Send + Sync>, + metrics_registry: Option>, + /// Should we try to prune the file (using stats) on open. + file_stats: Option>, + /// Maximal number of rows to read (after filtering) + limit: Option, + /// The row-offset assigned to the first row of the file. Used by the `row_idx` expression, + /// but not by the scan [`Selection`] which remains relative. + row_offset: u64, +} + +impl ScanBuilder { + /// Create a scan builder over `layout_reader` using `session` for runtime and execution state. + pub fn new(session: VortexSession, layout_reader: Arc) -> Self { + let projection = BoundExpression::new_root(layout_reader.dtype().clone()); + Self { + session, + layout_reader, + projection, + filter: None, + ordered: true, + row_range: None, + selection: Default::default(), + split_by: SplitBy::default(), + natural_splits: None, + // We default to four tasks per worker thread, which allows for some I/O lookahead + // without too much impact on work-stealing. + concurrency: 4, + map_fn: Arc::new(Ok), + metrics_registry: None, + file_stats: None, + limit: None, + row_offset: 0, + } + } + + /// Returns an [`ArrayStream`] with tasks spawned onto the session's runtime handle. + /// + /// See [`ScanBuilder::into_stream`] for more details. + pub fn into_array_stream(self) -> VortexResult { + let dtype = self.dtype()?; + let stream = self.into_stream()?; + Ok(ArrayStreamAdapter::new(dtype, stream)) + } + + /// Returns an [`ArrayIterator`] using the given blocking runtime. + pub fn into_array_iter( + self, + runtime: &B, + ) -> VortexResult { + let stream = self.into_array_stream()?; + let dtype = stream.dtype().clone(); + Ok(ArrayIteratorAdapter::new( + dtype, + runtime.block_on_stream(stream), + )) + } +} + +impl ScanBuilder { + /// Moves every option from a default builder into this executor. + pub fn from_default(builder: scan_builder::ScanBuilder) -> Self { + let parts = builder.into_parts(); + Self { + session: parts.session, + layout_reader: parts.layout_reader, + projection: parts.projection, + filter: parts.filter, + ordered: parts.ordered, + row_range: parts.row_range, + selection: parts.selection, + split_by: parts.split_by, + natural_splits: parts.natural_splits, + concurrency: parts.concurrency, + map_fn: parts.map_fn, + metrics_registry: parts.metrics_registry, + file_stats: parts.file_stats, + limit: parts.limit, + row_offset: parts.row_offset, + } + } + + /// Add a filter expression bound against the reader dtype. + pub fn with_filter(mut self, filter: BoundExpression) -> Self { + self.filter = Some(filter); + self + } + + /// Add or clear a filter expression bound against the reader dtype. + pub fn with_some_filter(mut self, filter: Option) -> Self { + self.filter = filter; + self + } + + /// Set a projection expression bound against the reader dtype. + pub fn with_projection(mut self, projection: BoundExpression) -> Self { + self.projection = projection; + self + } + + /// Returns whether output chunks are yielded in file order. + pub fn ordered(&self) -> bool { + self.ordered + } + + /// Configure whether output chunks must be yielded in file order. + pub fn with_ordered(mut self, ordered: bool) -> Self { + self.ordered = ordered; + self + } + + /// Restrict scanning to a contiguous row range. + pub fn with_row_range(mut self, row_range: Range) -> Self { + self.row_range = Some(row_range); + self + } + + /// Apply a row selection to the selected row range. + pub fn with_selection(mut self, selection: Selection) -> Self { + self.selection = selection; + self + } + + /// Select rows by strictly sorted absolute indices relative to the scan input. + pub fn with_row_indices(mut self, row_indices: StrictSortedBuffer) -> Self { + self.selection = Selection::IncludeByIndex(row_indices); + self + } + + /// Set the root row offset used by row-index expressions. + pub fn with_row_offset(mut self, row_offset: u64) -> Self { + self.row_offset = row_offset; + self + } + + /// Configure how natural scan work is split for concurrency. + pub fn with_split_by(mut self, split_by: SplitBy) -> Self { + self.split_by = split_by; + self + } + + /// Supply precomputed full-file natural split boundaries (see + /// [`full_file_splits`](Self::full_file_splits)) so [`prepare`](Self::prepare) reuses them + /// instead of walking the layout. Callers translating external partitions into row ranges + /// can compute the boundaries once per file and share them across partitions. + /// + /// Takes precedence over [`with_split_by`](Self::with_split_by); boundaries outside the + /// scan's row range are clamped during execution. Boundaries must be strictly increasing. + pub fn with_natural_splits(mut self, boundaries: Arc<[u64]>) -> Self { + debug_assert!( + boundaries.windows(2).all(|w| w[0] < w[1]), + "natural split boundaries must be strictly increasing" + ); + self.natural_splits = Some(boundaries); + self + } + + /// Compute the full-file natural split boundaries for the fields referenced by this scan's + /// projection and filter, ignoring any configured row range. + /// + /// These are the boundaries [`prepare`](Self::prepare) derives for a whole-file scan; hand + /// them back via [`with_natural_splits`](Self::with_natural_splits) to skip the layout walk + /// in `prepare`. + pub fn full_file_splits(&self) -> VortexResult> { + let field_mask = referenced_field_masks(&self.projection, self.filter.as_ref())?; + self.split_by.splits( + self.layout_reader.as_ref(), + &(0..self.layout_reader.row_count()), + &field_mask, + ) + } + + /// Returns the per-worker row-split concurrency. + pub fn concurrency(&self) -> usize { + self.concurrency + } + + /// The number of row splits to make progress on concurrently per-thread, must + /// be greater than 0. + pub fn with_concurrency(mut self, concurrency: usize) -> Self { + assert!(concurrency > 0); + self.concurrency = concurrency; + self + } + + /// Add or clear the metrics registry used by scan execution. + pub fn with_some_metrics_registry(mut self, metrics: Option>) -> Self { + self.metrics_registry = metrics; + self + } + + /// Set the metrics registry used by scan execution. + pub fn with_metrics_registry(mut self, metrics: Arc) -> Self { + self.metrics_registry = Some(metrics); + self + } + + /// Add or clear the maximum number of rows returned after filtering. + pub fn with_some_limit(mut self, limit: Option) -> Self { + self.limit = limit; + self + } + + /// Set the maximum number of rows returned after filtering. + pub fn with_limit(mut self, limit: u64) -> Self { + self.limit = Some(limit); + self + } + + /// The [`DType`] returned by the scan, after applying the projection. + pub fn dtype(&self) -> VortexResult { + Ok(self.projection.dtype().clone()) + } + + /// The session used by the scan. + pub fn session(&self) -> &VortexSession { + &self.session + } + + /// Map each split of the scan. The function will be run on the spawned task. + pub fn map( + self, + map_fn: impl Fn(A) -> VortexResult + 'static + Send + Sync, + ) -> ScanBuilder { + let old_map_fn = self.map_fn; + ScanBuilder { + session: self.session, + layout_reader: self.layout_reader, + projection: self.projection, + filter: self.filter, + ordered: self.ordered, + row_range: self.row_range, + selection: self.selection, + split_by: self.split_by, + natural_splits: self.natural_splits, + concurrency: self.concurrency, + metrics_registry: self.metrics_registry, + file_stats: self.file_stats, + limit: self.limit, + row_offset: self.row_offset, + map_fn: Arc::new(move |a| old_map_fn(a).and_then(&map_fn)), + } + } + + /// Optimize expressions, compute split ranges, and return an executable repeated scan. + pub fn prepare(self) -> VortexResult> { + let dtype = self.dtype()?; + + if self.filter.is_some() && self.limit.is_some() { + vortex_bail!("Vortex doesn't support scans with both a filter and a limit") + } + + let mut layout_reader = self.layout_reader; + + // Row-index expressions need the scan offset applied by the reader. + let mut found_row_idx = self.projection.contains::()?; + if !found_row_idx && let Some(filter) = self.filter.as_ref() { + found_row_idx = filter.contains::()?; + } + if found_row_idx { + layout_reader = Arc::new(RowIdxLayoutReader::new( + self.row_offset, + layout_reader, + self.session.clone(), + )); + } + + let bound_projection = self.projection; + let bound_filter = self.filter; + + let splits = + if let Some(ranges) = attempt_split_ranges(&self.selection, self.row_range.as_ref()) { + Splits::Ranges(ranges) + } else if let Some(boundaries) = self.natural_splits { + // Caller-supplied full-file boundaries; execution clamps them to the row range. + Splits::Natural(boundaries) + } else { + let field_mask = referenced_field_masks(&bound_projection, bound_filter.as_ref())?; + let split_range = self + .row_range + .clone() + .unwrap_or_else(|| 0..layout_reader.row_count()); + Splits::Natural( + self.split_by + .splits(layout_reader.as_ref(), &split_range, &field_mask)? + .into(), + ) + }; + + Ok(RepeatedScanV2::new( + self.session.clone(), + layout_reader, + bound_projection, + bound_filter, + self.ordered, + self.row_range, + self.selection, + splits, + self.concurrency, + self.map_fn, + self.limit, + dtype, + )) + } + + /// Constructs a task per row split of the scan, returned as a vector of futures. + pub fn build(self) -> VortexResult>>>> { + if self.limit.is_some_and(|l| l == 0) { + return Ok(vec![]); + } + + self.prepare()?.execute(None) + } + + /// Returns a [`Stream`] with tasks spawned onto the session's runtime handle. + pub fn into_stream( + self, + ) -> VortexResult> + Send + 'static + use> { + Ok(LazyScanStream::new(self)) + } + + /// Returns an [`Iterator`] using the session's runtime. + pub fn into_iter( + self, + runtime: &B, + ) -> VortexResult> + 'static> { + let stream = self.into_stream()?; + Ok(runtime.block_on_stream(stream)) + } +} + +enum LazyScanState { + Builder(Option>>), + Preparing(PreparingScan), + Stream(BoxStream<'static, VortexResult>), + Error(Option), +} + +type PreparedScanTasks = Vec>>>; + +struct PreparingScan { + ordered: bool, + concurrency: usize, + handle: Handle, + task: Task>>, +} + +struct LazyScanStream { + state: LazyScanState, +} + +impl LazyScanStream { + fn new(builder: ScanBuilder) -> Self { + Self { + state: LazyScanState::Builder(Some(Box::new(builder))), + } + } +} + +impl Unpin for LazyScanStream {} + +impl Stream for LazyScanStream { + type Item = VortexResult; + + fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + loop { + match &mut self.state { + LazyScanState::Builder(builder) => { + let builder = builder.take().vortex_expect("polled after completion"); + let ordered = builder.ordered; + let num_workers = get_available_parallelism().unwrap_or(1); + let concurrency = builder.concurrency * num_workers; + let handle = builder.session.handle(); + let task = handle + .spawn_cpu(move || builder.prepare().and_then(|scan| scan.execute(None))); + self.state = LazyScanState::Preparing(PreparingScan { + ordered, + concurrency, + handle, + task, + }); + } + LazyScanState::Preparing(preparing) => { + match ready!(Pin::new(&mut preparing.task).poll(cx)) { + Ok(tasks) => { + let ordered = preparing.ordered; + let concurrency = preparing.concurrency; + let handle = preparing.handle.clone(); + let stream = + futures::stream::iter(tasks).map(move |task| handle.spawn(task)); + let stream = if ordered { + stream.buffered(concurrency).boxed() + } else { + stream.buffer_unordered(concurrency).boxed() + }; + let stream = stream + .filter_map(|chunk| async move { chunk.transpose() }) + .boxed(); + self.state = LazyScanState::Stream(stream); + } + Err(err) => self.state = LazyScanState::Error(Some(err)), + } + } + LazyScanState::Stream(stream) => return stream.as_mut().poll_next(cx), + LazyScanState::Error(err) => return Poll::Ready(err.take().map(Err)), + } + } + } +} diff --git a/vortex-layout/src/scan/v2/tests.rs b/vortex-layout/src/scan/v2/tests.rs new file mode 100644 index 00000000000..16245c1ee6b --- /dev/null +++ b/vortex-layout/src/scan/v2/tests.rs @@ -0,0 +1,229 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use std::ops::Range; +use std::sync::Arc; + +use futures::TryStreamExt; +use futures::stream; +use rstest::rstest; +use vortex_array::ArrayContext; +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::VortexSessionExecute; +use vortex_array::arrays::ChunkedArray; +use vortex_array::assert_arrays_eq; +use vortex_array::dtype::DType; +use vortex_array::dtype::Nullability::NonNullable; +use vortex_array::dtype::PType; +use vortex_array::expr::gt; +use vortex_array::expr::lit; +use vortex_array::expr::root; +use vortex_buffer::Buffer; +use vortex_error::VortexResult; +use vortex_io::runtime::single::block_on; +use vortex_io::session::RuntimeSessionExt; +use vortex_scan::strict_sorted_buffer::StrictSortedBuffer; +use vortex_session::VortexSession; + +use crate::LayoutRef; +use crate::LayoutStrategy; +use crate::layouts::chunked::writer::ChunkedLayoutStrategy; +use crate::layouts::flat::writer::FlatLayoutStrategy; +use crate::scan::scan_builder::ScanBuilder; +use crate::scan::v2; +use crate::segments::SegmentSource; +use crate::segments::TestSegments; +use crate::sequence::SequenceId; +use crate::sequence::SequentialStreamAdapter; +use crate::sequence::SequentialStreamExt as _; +use crate::test::new_session; + +const CHUNK_ROWS: i32 = 1000; +const DTYPE: DType = DType::Primitive(PType::I32, NonNullable); + +/// Four chunks of consecutive integers, `0..4000`. +async fn write_layout( + session: &VortexSession, +) -> VortexResult<(Arc, LayoutRef)> { + let segments = Arc::new(TestSegments::default()); + let (mut sequence_id, eof) = SequenceId::root().split(); + let chunks = (0..4) + .map(|chunk| { + let values = Buffer::from_iter(chunk * CHUNK_ROWS..(chunk + 1) * CHUNK_ROWS); + Ok((sequence_id.advance(), values.into_array())) + }) + .collect::>(); + let layout = ChunkedLayoutStrategy::new(FlatLayoutStrategy::default()) + .write_stream( + ArrayContext::empty().into(), + Arc::::clone(&segments), + SequentialStreamAdapter::new(DTYPE, stream::iter(chunks)).sendable(), + eof, + session, + ) + .await?; + Ok((segments, layout)) +} + +#[derive(Clone)] +struct Case { + filter: bool, + row_range: Option>, + limit: Option, +} + +fn builder( + session: &VortexSession, + segments: &Arc, + layout: &LayoutRef, + case: &Case, +) -> VortexResult> { + let reader = layout.new_reader( + "".into(), + Arc::clone(segments), + session, + &Default::default(), + )?; + let mut builder = ScanBuilder::new(session.clone(), reader); + if case.filter { + builder = builder.with_filter(gt(root(), lit(1500_i32)).bind(&DTYPE)?); + } + if let Some(row_range) = case.row_range.clone() { + builder = builder.with_row_range(row_range); + } + if let Some(limit) = case.limit { + builder = builder.with_limit(limit); + } + Ok(builder) +} + +async fn await_tasks( + tasks: Vec>>>, +) -> VortexResult { + let mut chunks = Vec::new(); + for task in tasks { + if let Some(chunk) = task.await? { + chunks.push(chunk); + } + } + Ok(ChunkedArray::try_new(chunks, DTYPE)?.into_array()) +} + +fn case(filter: bool, row_range: Option>, limit: Option) -> Case { + Case { + filter, + row_range, + limit, + } +} + +#[rstest] +#[case::everything(case(false, None, None))] +#[case::filter(case(true, None, None))] +#[case::row_range(case(false, Some(500..2500), None))] +#[case::filter_and_row_range(case(true, Some(500..2500), None))] +#[case::limit(case(false, None, Some(1200)))] +#[case::limit_and_row_range(case(false, Some(900..3100), Some(1500)))] +fn stream_matches_default(#[case] case: Case) -> VortexResult<()> { + block_on(|handle| async move { + let session = new_session().with_handle(handle); + let (segments, layout) = write_layout(&session).await?; + + let expected = builder(&session, &segments, &layout, &case)? + .into_stream()? + .try_collect::>() + .await?; + let actual = v2::ScanBuilder::from_default(builder(&session, &segments, &layout, &case)?) + .into_stream()? + .try_collect::>() + .await?; + + assert_arrays_eq!( + ChunkedArray::try_new(actual, DTYPE)?, + ChunkedArray::try_new(expected, DTYPE)?, + &mut session.create_execution_ctx() + ); + Ok(()) + }) +} + +#[rstest] +#[case::whole_scan(case(true, None, None), None)] +#[case::execute_range(case(true, None, None), Some(700..3300))] +#[case::both_ranges(case(false, Some(500..2500), None), Some(2000..4000))] +#[case::empty_range(case(true, None, None), Some(1200..1200))] +fn execute_matches_default( + #[case] case: Case, + #[case] execute_range: Option>, +) -> VortexResult<()> { + block_on(|handle| async move { + let session = new_session().with_handle(handle); + let (segments, layout) = write_layout(&session).await?; + + let default = builder(&session, &segments, &layout, &case)?.prepare()?; + let replacement = + v2::ScanBuilder::from_default(builder(&session, &segments, &layout, &case)?) + .prepare()?; + assert_eq!(replacement.dtype(), default.dtype()); + + let expected = await_tasks(default.execute(execute_range.clone())?).await?; + let actual = await_tasks(replacement.execute(execute_range)?).await?; + assert_arrays_eq!(actual, expected, &mut session.create_execution_ctx()); + Ok(()) + }) +} + +#[test] +fn filter_with_limit_is_rejected() -> VortexResult<()> { + block_on(|handle| async move { + let session = new_session().with_handle(handle); + let (segments, layout) = write_layout(&session).await?; + let case = case(true, None, Some(10)); + assert!( + v2::ScanBuilder::from_default(builder(&session, &segments, &layout, &case)?) + .prepare() + .is_err() + ); + Ok(()) + }) +} + +#[test] +fn conversion_preserves_selection_and_mapper() -> VortexResult<()> { + block_on(|handle| async move { + let session = new_session().with_handle(handle); + let (segments, layout) = write_layout(&session).await?; + let indices = StrictSortedBuffer::try_new(Buffer::from_iter([500_u64, 1500, 2500]))?; + let builder = builder(&session, &segments, &layout, &case(false, None, None))? + .with_row_indices(indices) + .with_ordered(false) + .with_concurrency(2) + .map(|array| Ok(array.len())); + let builder = v2::ScanBuilder::from_default(builder); + assert!(!builder.ordered()); + assert_eq!(builder.concurrency(), 2); + + let lengths = builder.into_stream()?.try_collect::>().await?; + assert_eq!(lengths.into_iter().sum::(), 3); + Ok(()) + }) +} + +#[test] +fn conversion_preserves_natural_splits() -> VortexResult<()> { + block_on(|handle| async move { + let session = new_session().with_handle(handle); + let (segments, layout) = write_layout(&session).await?; + let case = case(false, Some(500..2500), None); + let builder = builder(&session, &segments, &layout, &case)? + .with_natural_splits(vec![0, 2000, 4000].into()) + .map(|array| Ok(array.len())); + let lengths = v2::ScanBuilder::from_default(builder) + .into_stream()? + .try_collect::>() + .await?; + assert_eq!(lengths, [1500, 500]); + Ok(()) + }) +}