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
6 changes: 3 additions & 3 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

14 changes: 11 additions & 3 deletions vortex-datafusion/src/persistent/opener.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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| {
Expand All @@ -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!(
Expand Down
10 changes: 8 additions & 2 deletions vortex-duckdb/src/file_reader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 _;

Expand Down Expand Up @@ -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();
Expand Down
1 change: 1 addition & 0 deletions vortex-layout/src/scan/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ mod splits;
mod tasks;
#[cfg(test)]
mod test;
pub mod v2;

/// A heuristic for an ideal split size.
///
Expand Down
42 changes: 42 additions & 0 deletions vortex-layout/src/scan/scan_builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -297,6 +297,28 @@ impl<A: 'static + Send> ScanBuilder<A> {
}
}

/// 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<A> {
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<RepeatedScan<A>> {
let dtype = self.dtype()?;
Expand Down Expand Up @@ -390,6 +412,26 @@ impl<A: 'static + Send> ScanBuilder<A> {
}
}

/// Configuration transferred to the alternative scan builder.
// TODO(joe): Remove once the V2 migration is complete.
pub(crate) struct ScanParts<A> {
pub(crate) session: VortexSession,
pub(crate) layout_reader: LayoutReaderRef,
pub(crate) projection: BoundExpression,
pub(crate) filter: Option<BoundExpression>,
pub(crate) ordered: bool,
pub(crate) row_range: Option<Range<u64>>,
pub(crate) selection: Selection,
pub(crate) split_by: SplitBy,
pub(crate) natural_splits: Option<Arc<[u64]>>,
pub(crate) concurrency: usize,
pub(crate) map_fn: Arc<dyn Fn(ArrayRef) -> VortexResult<A> + Send + Sync>,
pub(crate) metrics_registry: Option<Arc<dyn MetricsRegistry>>,
pub(crate) file_stats: Option<Arc<[StatsSet]>>,
pub(crate) limit: Option<u64>,
pub(crate) row_offset: u64,
}

enum LazyScanState<A: 'static + Send> {
Builder(Option<Box<ScanBuilder<A>>>),
Preparing(PreparingScan<A>),
Expand Down
33 changes: 33 additions & 0 deletions vortex-layout/src/scan/v2/mod.rs
Original file line number Diff line number Diff line change
@@ -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<bool> = 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;
Loading
Loading