diff --git a/vortex-layout/src/scan/mod.rs b/vortex-layout/src/scan/mod.rs index 6b2c07a5b83..98fd1918a42 100644 --- a/vortex-layout/src/scan/mod.rs +++ b/vortex-layout/src/scan/mod.rs @@ -5,12 +5,11 @@ pub mod arrow; mod filter; pub mod layout; pub mod multi; -mod plan; -pub mod plan_v2; pub mod repeated_scan; pub mod scan_builder; pub mod split_by; mod splits; +mod tasks; #[cfg(test)] mod test; diff --git a/vortex-layout/src/scan/plan_v2.rs b/vortex-layout/src/scan/plan_v2.rs deleted file mode 100644 index 0af9745402e..00000000000 --- a/vortex-layout/src/scan/plan_v2.rs +++ /dev/null @@ -1,455 +0,0 @@ -// SPDX-License-Identifier: Apache-2.0 -// SPDX-FileCopyrightText: Copyright the Vortex contributors - -// See https://github.com/vortex-data/vortex/issues/9062 - -use std::ops::BitAnd; -use std::ops::Range; -use std::sync::Arc; - -use bit_vec::BitVec; -use futures::FutureExt; -use futures::future::BoxFuture; -use vortex_array::ArrayRef; -use vortex_array::MaskFuture; -use vortex_array::dtype::DType; -use vortex_array::expr::Expression; -use vortex_array::expr::root; -use vortex_array::expr::transform::replace; -use vortex_error::VortexResult; -use vortex_error::vortex_bail; -use vortex_mask::Mask; -use vortex_scan::row_mask::RowMask; - -use crate::ArrayFuture; -use crate::LayoutReaderRef; -use crate::scan::filter::FilterExpr; - -pub(crate) struct PlanV2 { - projection: ScanPlanRef, - predicates: Vec, - filter: Option, -} - -impl PlanV2 { - pub(crate) fn new( - projection: ScanPlanRef, - predicates: Vec, - filter: Option, - ) -> Self { - Self { - projection, - predicates, - filter, - } - } - - pub(crate) fn task_context( - &self, - mapper: Arc VortexResult + Send + Sync>, - ) -> Arc> { - Arc::new(TaskContext { - filter: self - .filter - .clone() - .map(|filter| Arc::new(FilterExpr::new(filter))), - predicates: self.predicates.clone(), - projection: Arc::clone(&self.projection), - mapper, - }) - } -} - -/// Shared handle to a heap-allocated V2 physical scan plan. -pub type ScanPlanRef = Arc; - -/// A heap-allocated physical scan plan. -/// -/// A source plan represents an instantiated layout. [`apply_expr`](Self::apply_expr) derives -/// another plan whose root value is the applied expression, and [`optimize`](Self::optimize) -/// rewrites that derived plan before execution. Execution therefore selects an already-bound plan -/// and supplies only its row range and mask. -pub trait ScanPlan: 'static + Send + Sync { - /// Apply `expr` to this plan's root value and return the resulting plan. - fn apply_expr(self: Arc, expr: Expression) -> VortexResult; - - /// Optimize this plan and return the resulting plan. - fn optimize(self: Arc) -> VortexResult; - - /// Returns the name of the underlying layout reader for debugging. - fn name(&self) -> &Arc; - - /// Returns the dtype produced by this plan. - fn dtype(&self) -> &DType; - - /// Returns the number of rows in this plan's row domain. - fn row_count(&self) -> u64; - - /// Returns a mask where all false values are proven false for this plan. - fn pruning_evaluation(&self, row_range: &Range, mask: Mask) -> VortexResult; - - /// Evaluates this boolean plan and intersects it with `mask`. - fn filter_evaluation( - &self, - row_range: &Range, - mask: MaskFuture, - ) -> VortexResult; - - /// Evaluates this plan over the selected rows. - fn projection_evaluation( - &self, - row_range: &Range, - mask: MaskFuture, - ) -> VortexResult; -} - -/// Compatibility V2 source and expression plan backed by a layout reader. -/// -/// Applying an expression and optimizing it produce new heap-allocated plans. Execution delegates -/// the resulting expression to the established reader implementation. Layout-specific source plans -/// can replace this compatibility node without changing the split execution loop. -pub struct LayoutReaderScanPlanV2 { - reader: LayoutReaderRef, - expr: Expression, - dtype: DType, -} - -impl LayoutReaderScanPlanV2 { - /// Create a V2 source plan for `reader`. - pub fn new(reader: LayoutReaderRef) -> Self { - let dtype = reader.dtype().clone(); - Self { - reader, - expr: root(), - dtype, - } - } - - fn try_new(reader: LayoutReaderRef, expr: Expression) -> VortexResult { - let dtype = expr.return_dtype(reader.dtype())?; - Ok(Self { - reader, - expr, - dtype, - }) - } -} - -impl ScanPlan for LayoutReaderScanPlanV2 { - fn apply_expr(self: Arc, expr: Expression) -> VortexResult { - let expr = replace(expr, &root(), self.expr.clone()); - Ok(Arc::new(Self::try_new(Arc::clone(&self.reader), expr)?)) - } - - fn optimize(self: Arc) -> VortexResult { - let expr = self.expr.optimize_recursive(self.reader.dtype())?; - Ok(Arc::new(Self::try_new(Arc::clone(&self.reader), expr)?)) - } - - fn name(&self) -> &Arc { - self.reader.name() - } - - fn dtype(&self) -> &DType { - &self.dtype - } - - fn row_count(&self) -> u64 { - self.reader.row_count() - } - - fn pruning_evaluation(&self, row_range: &Range, mask: Mask) -> VortexResult { - self.reader.pruning_evaluation(row_range, &self.expr, mask) - } - - fn filter_evaluation( - &self, - row_range: &Range, - mask: MaskFuture, - ) -> VortexResult { - self.reader.filter_evaluation(row_range, &self.expr, mask) - } - - fn projection_evaluation( - &self, - row_range: &Range, - mask: MaskFuture, - ) -> VortexResult { - self.reader - .projection_evaluation(row_range, &self.expr, mask) - } -} - -/// Environment variable selecting the scan planning implementation. -pub const SCAN_IMPL_ENV: &str = "VORTEX_SCAN_IMPL"; - -/// Returns whether V2 heap-allocated planning is enabled for this process. -/// -/// The existing `plan` path remains the default on this extraction branch. Set -/// `VORTEX_SCAN_IMPL=planv2` to exercise the V2 path with the same execution implementation. -pub fn plan_v2_enabled() -> VortexResult { - match std::env::var(SCAN_IMPL_ENV) { - Ok(value) => parse_scan_impl(&value), - Err(std::env::VarError::NotPresent) => Ok(false), - Err(std::env::VarError::NotUnicode(value)) => { - vortex_bail!("{SCAN_IMPL_ENV} must be valid unicode, got {value:?}") - } - } -} - -fn parse_scan_impl(value: &str) -> VortexResult { - match value { - "" | "plan" | "v1" | "legacy" | "layout-reader" => Ok(false), - "planv2" | "plan-v2" | "v2" | "planned" | "scan-plan" => Ok(true), - other => vortex_bail!( - "{SCAN_IMPL_ENV} must be one of plan, v1, legacy, layout-reader, planv2, plan-v2, v2, planned, or scan-plan, got {other:?}" - ), - } -} - -/// Execute one split using a V2 physical scan plan. -/// -/// The execution order intentionally mirrors [`crate::scan::plan::split_exec`]. Expressions were -/// consumed during planning, so execution selects a predicate or projection plan without passing -/// an expression. -pub(crate) fn split_exec( - ctx: Arc>, - read_mask: RowMask, - limit: Option<&mut u64>, -) -> VortexResult>>> { - let row_range = read_mask.row_range(); - let row_mask = read_mask.mask().clone(); - - let filter_mask = match ctx.filter.as_ref() { - None => { - let row_mask = match limit { - Some(l) if *l == 0 => Mask::new_false(row_mask.len()), - Some(l) => { - let true_count = row_mask.true_count(); - let mask_limit = usize::try_from(*l) - .map(|l| l.min(true_count)) - .unwrap_or(true_count); - let row_mask = row_mask.limit(mask_limit); - *l -= mask_limit as u64; - row_mask - } - None => row_mask, - }; - - MaskFuture::ready(row_mask) - } - Some(filter) => { - if filter.conjuncts().len() != ctx.predicates.len() { - vortex_bail!( - "physical predicate count {} does not match conjunct count {}", - ctx.predicates.len(), - filter.conjuncts().len() - ); - } - - let ctx = Arc::clone(&ctx); - let filter = Arc::clone(filter); - let row_range = row_range.clone(); - - MaskFuture::new(row_mask.len(), async move { - let mut mask = row_mask; - let mut dynamic_versions = vec![None; filter.conjuncts().len()]; - - for (idx, predicate) in ctx.predicates.iter().enumerate() { - if mask.all_false() { - return Ok(mask); - } - - dynamic_versions[idx] = filter.dynamic_updates(idx).map(|du| du.version()); - let conjunct_mask = predicate - .pruning_evaluation(&row_range, mask.clone())? - .await?; - mask = mask.bitand(&conjunct_mask); - } - - let mut remaining = BitVec::from_elem(filter.conjuncts().len(), true); - while let Some(idx) = filter.next_conjunct(&remaining) { - remaining.set(idx, false); - if mask.all_false() { - return Ok(mask); - } - - let current_version = filter.dynamic_updates(idx).map(|du| du.version()); - if let Some(version) = current_version - && dynamic_versions[idx].is_none_or(|old| old < version) - { - dynamic_versions[idx] = Some(version); - let conjunct_mask = ctx.predicates[idx] - .pruning_evaluation(&row_range, mask.clone())? - .await?; - mask = mask.bitand(&conjunct_mask); - } - if mask.all_false() { - return Ok(mask); - } - - let conjunct_mask = ctx.predicates[idx] - .filter_evaluation(&row_range, MaskFuture::ready(mask))? - .await?; - filter.report_selectivity(idx, conjunct_mask.density()); - mask = conjunct_mask; - } - - Ok(mask) - }) - } - }; - - let projection_future = ctx - .projection - .projection_evaluation(&row_range, filter_mask.clone())?; - - let mapper = Arc::clone(&ctx.mapper); - let array_fut = async move { - let mask = filter_mask.await?; - if mask.all_false() { - return Ok(None); - } - - let array = projection_future.await?; - mapper(array).map(Some) - }; - - Ok(array_fut.boxed()) -} - -/// Information needed to execute one split from a V2 physical scan plan. -pub(crate) struct TaskContext { - filter: Option>, - predicates: Vec, - projection: ScanPlanRef, - mapper: Arc VortexResult + Send + Sync>, -} - -#[cfg(test)] -mod tests { - use std::any::Any; - - use vortex_array::dtype::FieldMask; - use vortex_array::dtype::FieldName; - use vortex_array::dtype::Nullability; - use vortex_array::dtype::PType; - use vortex_array::dtype::StructFields; - use vortex_array::expr::eq; - use vortex_array::expr::get_item; - use vortex_array::expr::lit; - - use super::*; - use crate::LayoutReader; - use crate::RowSplits; - use crate::SplitRange; - - #[test] - fn scan_impl_accepts_v1_and_v2_values() -> VortexResult<()> { - for value in ["", "plan", "v1", "legacy", "layout-reader"] { - assert!(!parse_scan_impl(value)?); - } - for value in ["planv2", "plan-v2", "v2", "planned", "scan-plan"] { - assert!(parse_scan_impl(value)?); - } - Ok(()) - } - - #[test] - fn scan_impl_rejects_unknown_value() { - assert!(parse_scan_impl("unknown").is_err()); - } - - struct TestLayoutReader { - name: Arc, - dtype: DType, - } - - impl TestLayoutReader { - fn new() -> Self { - Self { - name: Arc::from("test"), - dtype: DType::Struct( - StructFields::from_iter([( - FieldName::from("a"), - DType::Primitive(PType::I32, Nullability::NonNullable), - )]), - Nullability::NonNullable, - ), - } - } - } - - impl LayoutReader for TestLayoutReader { - fn name(&self) -> &Arc { - &self.name - } - - fn as_any(&self) -> &dyn Any { - self - } - - fn dtype(&self) -> &DType { - &self.dtype - } - - fn row_count(&self) -> u64 { - 1 - } - - fn register_splits( - &self, - _field_mask: &[FieldMask], - _split_range: &SplitRange, - _splits: &mut RowSplits, - ) -> VortexResult<()> { - unimplemented!("not needed for scan-plan construction") - } - - fn pruning_evaluation( - &self, - _row_range: &Range, - _expr: &Expression, - _mask: Mask, - ) -> VortexResult { - unimplemented!("not needed for scan-plan construction") - } - - fn filter_evaluation( - &self, - _row_range: &Range, - _expr: &Expression, - _mask: MaskFuture, - ) -> VortexResult { - unimplemented!("not needed for scan-plan construction") - } - - fn projection_evaluation( - &self, - _row_range: &Range, - _expr: &Expression, - _mask: MaskFuture, - ) -> VortexResult { - unimplemented!("not needed for scan-plan construction") - } - } - - #[test] - fn scan_plan_v2_applies_expressions_to_the_current_root() -> VortexResult<()> { - let reader: LayoutReaderRef = Arc::new(TestLayoutReader::new()); - let source: ScanPlanRef = Arc::new(LayoutReaderScanPlanV2::new(reader)); - - let field = Arc::clone(&source) - .apply_expr(get_item("a", root()))? - .optimize()?; - assert_eq!( - field.dtype(), - &DType::Primitive(PType::I32, Nullability::NonNullable) - ); - - let predicate = field.apply_expr(eq(root(), lit(1_i32)))?.optimize()?; - assert_eq!(predicate.dtype(), &DType::Bool(Nullability::NonNullable)); - - Ok(()) - } -} diff --git a/vortex-layout/src/scan/repeated_scan.rs b/vortex-layout/src/scan/repeated_scan.rs index 117fa276168..681f33639bc 100644 --- a/vortex-layout/src/scan/repeated_scan.rs +++ b/vortex-layout/src/scan/repeated_scan.rs @@ -26,10 +26,10 @@ use vortex_session::VortexSession; use vortex_utils::parallelism::get_available_parallelism; use crate::LayoutReaderRef; -use crate::scan::plan; -use crate::scan::plan_v2; -use crate::scan::plan_v2::ScanPlanRef; +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. /// @@ -37,7 +37,9 @@ use crate::scan::splits::Splits; /// data source. pub struct RepeatedScan { session: VortexSession, - execution: ExecutionPlan, + layout_reader: LayoutReaderRef, + projection: Expression, + filter: Option, ordered: bool, /// Optionally read a subset of the rows in the file. row_range: Option>, @@ -55,16 +57,6 @@ pub struct RepeatedScan { dtype: DType, } -enum ExecutionPlan { - Plan(plan::Plan), - PlanV2(plan_v2::PlanV2), -} - -enum ExecutionTaskContext { - Plan(Arc>), - PlanV2(Arc>), -} - impl RepeatedScan { pub fn dtype(&self) -> &DType { &self.dtype @@ -97,7 +89,7 @@ impl RepeatedScan { clippy::too_many_arguments, reason = "all arguments are needed for scan construction" )] - pub fn new_plan( + pub fn new( session: VortexSession, layout_reader: LayoutReaderRef, projection: Expression, @@ -113,40 +105,9 @@ impl RepeatedScan { ) -> Self { Self { session, - execution: ExecutionPlan::Plan(plan::Plan::new(layout_reader, projection, filter)), - ordered, - row_range, - selection, - splits, - concurrency, - map_fn, - limit, - dtype, - } - } - - /// Construct a repeated scan from a prepared heap-allocated physical plan. - #[expect( - clippy::too_many_arguments, - reason = "all arguments are needed for scan construction" - )] - pub fn new_plan_v2( - session: VortexSession, - projection: ScanPlanRef, - predicates: Vec, - 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, - execution: ExecutionPlan::PlanV2(plan_v2::PlanV2::new(projection, predicates, filter)), + layout_reader, + projection, + filter, ordered, row_range, selection, @@ -212,14 +173,12 @@ impl RepeatedScan { let mut limit = self.limit; let mut tasks = Vec::new(); - let ctx = match &self.execution { - ExecutionPlan::Plan(plan) => { - ExecutionTaskContext::Plan(plan.task_context(Arc::clone(&self.map_fn))) - } - ExecutionPlan::PlanV2(plan_v2) => { - ExecutionTaskContext::PlanV2(plan_v2.task_context(Arc::clone(&self.map_fn))) - } - }; + 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); @@ -227,14 +186,7 @@ impl RepeatedScan { continue; } - tasks.push(match &ctx { - ExecutionTaskContext::Plan(ctx) => { - plan::split_exec(Arc::clone(ctx), row_mask, limit.as_mut())? - } - ExecutionTaskContext::PlanV2(ctx) => { - plan_v2::split_exec(Arc::clone(ctx), row_mask, limit.as_mut())? - } - }); + tasks.push(split_exec(Arc::clone(&ctx), row_mask, limit.as_mut())?); if limit.is_some_and(|l| l == 0) { break; } diff --git a/vortex-layout/src/scan/scan_builder.rs b/vortex-layout/src/scan/scan_builder.rs index 7354861a806..03f8c49649d 100644 --- a/vortex-layout/src/scan/scan_builder.rs +++ b/vortex-layout/src/scan/scan_builder.rs @@ -18,7 +18,6 @@ use vortex_array::dtype::DType; use vortex_array::dtype::FieldMask; use vortex_array::expr::Expression; use vortex_array::expr::analysis::referenced_field_paths; -use vortex_array::expr::forms::conjuncts; use vortex_array::expr::root; use vortex_array::iter::ArrayIterator; use vortex_array::iter::ArrayIteratorAdapter; @@ -41,9 +40,6 @@ use vortex_utils::parallelism::get_available_parallelism; use crate::LayoutReader; use crate::LayoutReaderRef; use crate::layouts::row_idx::RowIdxLayoutReader; -use crate::scan::plan_v2::LayoutReaderScanPlanV2; -use crate::scan::plan_v2::ScanPlanRef; -use crate::scan::plan_v2::plan_v2_enabled; use crate::scan::repeated_scan::RepeatedScan; use crate::scan::split_by::SplitBy; use crate::scan::splits::Splits; @@ -313,34 +309,7 @@ impl ScanBuilder { )?) }; - if plan_v2_enabled()? { - let source: ScanPlanRef = - Arc::new(LayoutReaderScanPlanV2::new(Arc::clone(&layout_reader))); - let projection_plan = Arc::clone(&source).apply_expr(projection)?.optimize()?; - let predicate_plans = filter - .as_ref() - .map(conjuncts) - .unwrap_or_default() - .into_iter() - .map(|expr| Arc::clone(&source).apply_expr(expr)?.optimize()) - .collect::>>()?; - return Ok(RepeatedScan::new_plan_v2( - self.session.clone(), - projection_plan, - predicate_plans, - filter, - self.ordered, - self.row_range, - self.selection, - splits, - self.concurrency, - self.map_fn, - self.limit, - dtype, - )); - } - - Ok(RepeatedScan::new_plan( + Ok(RepeatedScan::new( self.session.clone(), layout_reader, projection, diff --git a/vortex-layout/src/scan/plan.rs b/vortex-layout/src/scan/tasks.rs similarity index 85% rename from vortex-layout/src/scan/plan.rs rename to vortex-layout/src/scan/tasks.rs index 439cbee1be3..a86546e15ef 100644 --- a/vortex-layout/src/scan/plan.rs +++ b/vortex-layout/src/scan/tasks.rs @@ -16,46 +16,11 @@ use vortex_error::VortexResult; use vortex_mask::Mask; use vortex_scan::row_mask::RowMask; -use crate::LayoutReaderRef; +use crate::LayoutReader; use crate::scan::filter::FilterExpr; pub type TaskFuture = BoxFuture<'static, VortexResult>; -pub(crate) struct Plan { - layout_reader: LayoutReaderRef, - projection: Expression, - filter: Option, -} - -impl Plan { - pub(crate) fn new( - layout_reader: LayoutReaderRef, - projection: Expression, - filter: Option, - ) -> Self { - Self { - layout_reader, - projection, - filter, - } - } - - pub(crate) fn task_context( - &self, - mapper: Arc VortexResult + Send + Sync>, - ) -> Arc> { - Arc::new(TaskContext { - filter: self - .filter - .clone() - .map(|filter| Arc::new(FilterExpr::new(filter))), - reader: Arc::clone(&self.layout_reader), - projection: self.projection.clone(), - mapper, - }) - } -} - /// Logic for executing a single split reading task. /// N.B. read_mask should be evaluated against all_false() before calling this /// method to avoid creating an empty TaskFuture. @@ -188,9 +153,13 @@ pub fn split_exec( /// Information needed to execute a single split task. /// /// Row selection is evaluated before creating a split task so it's not included -pub(crate) struct TaskContext { - filter: Option>, - reader: LayoutReaderRef, - projection: Expression, - mapper: Arc VortexResult + Send + Sync>, +pub struct TaskContext { + /// The shared filter expression. + pub filter: Option>, + /// The layout reader. + pub reader: Arc, + /// The projection expression to apply to gather the scanned rows. + pub projection: Expression, + /// Function that maps into an A. + pub mapper: Arc VortexResult + Send + Sync>, }