From b21ac66c1d7bc9634eda4ad1d993dbae4c919941 Mon Sep 17 00:00:00 2001 From: Joe Isaacs Date: Thu, 8 Oct 2026 20:15:19 +0000 Subject: [PATCH] feat(layout): add Eval, RowIdx, and Take exec nodes - `EvalNode` applies its bound expression to each array its child produces as it streams through. The child runs over the same selection, so the expression never sees an unselected row. - `RowIdxNode` emits the global row index of every selected row at start, offset by the graph's row offset. - `TakeNode` runs the codes over the selection and the values over their whole domain. It is `Ready::Closed(values)`: it does not run until the values have closed, joins them once, and then wraps each codes array as a dictionary over them as it arrives. Signed-off-by: Joe Isaacs Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01NWYmfHTFtaEdhdyd3u5hv3 --- vortex-layout/src/plan/exec/eval.rs | 49 +++++++++++++++ vortex-layout/src/plan/exec/mod.rs | 6 ++ vortex-layout/src/plan/exec/row_idx.rs | 49 +++++++++++++++ vortex-layout/src/plan/exec/take.rs | 80 +++++++++++++++++++++++++ vortex-layout/src/plan/exec/tests.rs | 65 ++++++++++++++++++++ vortex-layout/src/plan/plans/eval.rs | 18 ++++++ vortex-layout/src/plan/plans/row_idx.rs | 18 ++++++ vortex-layout/src/plan/plans/take.rs | 18 ++++++ 8 files changed, 303 insertions(+) create mode 100644 vortex-layout/src/plan/exec/eval.rs create mode 100644 vortex-layout/src/plan/exec/row_idx.rs create mode 100644 vortex-layout/src/plan/exec/take.rs diff --git a/vortex-layout/src/plan/exec/eval.rs b/vortex-layout/src/plan/exec/eval.rs new file mode 100644 index 00000000000..43bba4792af --- /dev/null +++ b/vortex-layout/src/plan/exec/eval.rs @@ -0,0 +1,49 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use vortex_error::VortexResult; + +use crate::plan::EvalPlan; +use crate::plan::exec::ExecNode; +use crate::plan::exec::NodeState; +use crate::plan::exec::StepCx; +use crate::plan::exec::selection::Selection; + +const CHILD: usize = 0; + +/// Applies an expression to each array its child produces. +/// +/// The child runs over the same rows and selection, so every array holds only selected rows and +/// the expression never sees a row the selection removed. +pub(crate) struct EvalNode { + plan: EvalPlan, + selection: Selection, +} + +impl EvalNode { + pub(crate) fn new(plan: EvalPlan, selection: Selection) -> Self { + Self { plan, selection } + } +} + +impl ExecNode for EvalNode { + fn start(&mut self, cx: &mut StepCx<'_>) -> VortexResult { + cx.spawn( + CHILD, + self.plan.child_plan()?, + self.selection.rows().clone(), + self.selection.mask().clone(), + ); + Ok(NodeState::Wait) + } + + fn compute(&mut self, cx: &mut StepCx<'_>) -> VortexResult { + for array in cx.input(CHILD).take_all() { + cx.emit(array.apply_bound(self.plan.expression())?); + } + if cx.input(CHILD).finished() { + return Ok(NodeState::Done); + } + Ok(NodeState::Wait) + } +} diff --git a/vortex-layout/src/plan/exec/mod.rs b/vortex-layout/src/plan/exec/mod.rs index 127f5898b8d..21ad122bdd9 100644 --- a/vortex-layout/src/plan/exec/mod.rs +++ b/vortex-layout/src/plan/exec/mod.rs @@ -37,11 +37,14 @@ //! [`Ready::AllClosed`], are barriers between them, as a join build is in a query engine. mod concat; +mod eval; mod filter; mod pack; +mod row_idx; mod segment_scan; mod selection; pub mod synthetic; +mod take; use std::collections::VecDeque; use std::mem; @@ -674,10 +677,13 @@ impl From> for Effects { } pub(crate) use concat::ConcatNode; +pub(crate) use eval::EvalNode; pub(crate) use filter::FilterNode; pub(crate) use pack::PackNode; +pub(crate) use row_idx::RowIdxNode; pub(crate) use segment_scan::SegmentScanNode; pub(crate) use selection::Selection; +pub(crate) use take::TakeNode; #[cfg(test)] mod scheduling_tests; diff --git a/vortex-layout/src/plan/exec/row_idx.rs b/vortex-layout/src/plan/exec/row_idx.rs new file mode 100644 index 00000000000..3989e56211d --- /dev/null +++ b/vortex-layout/src/plan/exec/row_idx.rs @@ -0,0 +1,49 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use vortex_array::IntoArray; +use vortex_buffer::Buffer; +use vortex_error::VortexResult; + +use crate::plan::exec::ExecNode; +use crate::plan::exec::NodeState; +use crate::plan::exec::StepCx; +use crate::plan::exec::selection::Selection; + +/// Produces the global row index of every selected row. +/// +/// Row-index plans sit in the graph's root row domain, so a row's global index is the graph's +/// row offset plus its plan row. The node has no input and emits once, at start. +pub(crate) struct RowIdxNode { + selection: Selection, + /// The global row index of the graph's first plan row. + row_offset: u64, +} + +impl RowIdxNode { + pub(crate) fn new(selection: Selection, row_offset: u64) -> Self { + Self { + selection, + row_offset, + } + } +} + +impl ExecNode for RowIdxNode { + fn start(&mut self, cx: &mut StepCx<'_>) -> VortexResult { + let rows = self.selection.rows(); + let indices = Buffer::from_iter(rows.start + self.row_offset..rows.end + self.row_offset) + .into_array(); + let array = if self.selection.mask().all_true() { + indices + } else { + indices.filter(self.selection.mask().clone())? + }; + cx.emit(array); + Ok(NodeState::Done) + } + + fn compute(&mut self, _cx: &mut StepCx<'_>) -> VortexResult { + Ok(NodeState::Done) + } +} diff --git a/vortex-layout/src/plan/exec/take.rs b/vortex-layout/src/plan/exec/take.rs new file mode 100644 index 00000000000..899c0e9970d --- /dev/null +++ b/vortex-layout/src/plan/exec/take.rs @@ -0,0 +1,80 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::arrays::DictArray; +use vortex_error::VortexResult; +use vortex_mask::Mask; + +use crate::plan::TakePlan; +use crate::plan::exec::ExecNode; +use crate::plan::exec::NodeState; +use crate::plan::exec::Ready; +use crate::plan::exec::StepCx; +use crate::plan::exec::selection::Selection; +use crate::plan::exec::selection::join; + +const CODES: usize = 0; +const VALUES: usize = 1; + +/// Looks up each selected row's code in the full set of values. +/// +/// The codes run over the selection; the values run over their whole domain. The node does not +/// run until the values have closed, so the first compute joins them once, and every compute +/// wraps the codes arrays that have arrived as dictionaries over the joined values. +pub(crate) struct TakeNode { + plan: TakePlan, + selection: Selection, + values: Option, +} + +impl TakeNode { + pub(crate) fn new(plan: TakePlan, selection: Selection) -> Self { + Self { + plan, + selection, + values: None, + } + } +} + +impl ExecNode for TakeNode { + fn ready(&self) -> Ready { + Ready::Closed(&[VALUES]) + } + + fn start(&mut self, cx: &mut StepCx<'_>) -> VortexResult { + if self.selection.mask().all_false() { + return Ok(NodeState::Done); + } + cx.spawn( + CODES, + self.plan.codes()?, + self.selection.rows().clone(), + self.selection.mask().clone(), + ); + let values = self.plan.values()?; + let len = usize::try_from(values.row_count())?; + cx.spawn(VALUES, values, 0..len as u64, Mask::new_true(len)); + Ok(NodeState::Wait) + } + + fn compute(&mut self, cx: &mut StepCx<'_>) -> VortexResult { + let values = match &self.values { + Some(values) => values.clone(), + None => { + let values = join(self.plan.values()?.dtype(), cx.input(VALUES).take_all())?; + self.values = Some(values.clone()); + values + } + }; + for codes in cx.input(CODES).take_all() { + cx.emit(DictArray::try_new(codes, values.clone())?.into_array()); + } + if cx.input(CODES).finished() { + return Ok(NodeState::Done); + } + Ok(NodeState::Wait) + } +} diff --git a/vortex-layout/src/plan/exec/tests.rs b/vortex-layout/src/plan/exec/tests.rs index 390a6e9734c..3dd4530dbb2 100644 --- a/vortex-layout/src/plan/exec/tests.rs +++ b/vortex-layout/src/plan/exec/tests.rs @@ -17,6 +17,10 @@ use vortex_array::arrays::StructArray; use vortex_array::arrays::VarBinViewArray; use vortex_array::assert_arrays_eq; use vortex_array::buffer::BufferHandle; +use vortex_array::expr::get_item; +use vortex_array::expr::gt; +use vortex_array::expr::lit; +use vortex_array::expr::root; use vortex_array::serde::SerializeOptions; use vortex_buffer::Alignment; use vortex_buffer::ByteBufferMut; @@ -29,12 +33,16 @@ use super::*; use crate::LayoutRef; use crate::OwnedLayoutChildren; use crate::layouts::chunked::ChunkedLayout; +use crate::layouts::dict::DictLayout; use crate::layouts::flat::FlatLayout; use crate::layouts::struct_::StructLayout; +use crate::plan::EvalPlan; use crate::plan::Filter; use crate::plan::SegmentScan; +use crate::plan::Take; use crate::plan::exec::selection::join; use crate::plan::lower; +use crate::plan::optimize; use crate::test::SESSION; const ROWS: u64 = 20; @@ -513,6 +521,43 @@ fn bare_scan_is_dense_and_filter_keeps_the_selection(#[case] sel: Sel) -> Vortex Ok(()) } +/// A take reads its values over their whole domain and its codes over the selection, including +/// when a predicate has been pushed onto the values, and emits nothing before the values are +/// whole. +#[rstest] +#[case::values(false)] +#[case::predicate(true)] +fn take_waits_for_whole_values(#[case] predicate: bool) -> VortexResult<()> { + let mut store = Store::default(); + let values = VarBinViewArray::from_iter_str(["a", "b", "c"]).into_array(); + let codes = PrimitiveArray::from_iter((0..ROWS).map(|v| (v % 3) as u8)).into_array(); + let layout = DictLayout::new(store.flat(&values)?, store.flat(&codes)?).into_layout(); + let mut plan = lower(&layout)?; + let mut expected = values.take(codes)?; + if predicate { + let expression = gt(root(), lit("a")) + .bind(plan.dtype())? + .optimize_recursive()?; + expected = expected.apply_bound(&expression)?; + plan = optimize(EvalPlan::try_new(expression, plan)?.into_plan())?; + } + assert!(plan.is::()); + + for rows in [0..10, 10..ROWS] { + let mask = Sel::EveryOther.mask(10); + // Codes (segment 1) land first; nothing comes out until the values (segment 0) do. + let run = run(&store, &plan, rows.clone(), mask.clone(), scripted(&[1, 0]))?; + assert_eq!(reads(&run.events), 2); + assert_eq!( + run.events.iter().position(|e| matches!(e, Event::Piece(_))), + Some(3), + "the only array must follow both deliveries" + ); + assert_view(&expected, &rows, &mask, run.arrays)?; + } + Ok(()) +} + /// A graph sharing a decode cache with one that already ran over the same plan reads nothing and /// returns the same rows, even under a different selection. #[test] @@ -546,3 +591,23 @@ fn shared_decode_cache_skips_reads() -> VortexResult<()> { assert_view(&expected, &rows, &mask, second.arrays)?; Ok(()) } + +#[test] +fn eval_applies_expression_to_selected_rows() -> VortexResult<()> { + let mut store = Store::default(); + let (plan, expected) = two_columns(&mut store)?; + let expression = gt(get_item("a", root()), lit(4_i32)).bind(plan.dtype())?; + let expected = expected.apply_bound(&expression)?; + let plan = EvalPlan::try_new(expression, plan)?.into_plan(); + + let rows = 2..18; + let mask = Sel::EveryOther.mask(16); + let run = run( + &store, + &plan, + rows.clone(), + mask.clone(), + delivery(Delivery::Lifo), + )?; + assert_view(&expected, &rows, &mask, run.arrays) +} diff --git a/vortex-layout/src/plan/plans/eval.rs b/vortex-layout/src/plan/plans/eval.rs index f1edb967e41..afd9931c1ae 100644 --- a/vortex-layout/src/plan/plans/eval.rs +++ b/vortex-layout/src/plan/plans/eval.rs @@ -3,11 +3,13 @@ use std::borrow::Cow; use std::fmt; +use std::ops::Range; use vortex_array::EmptyMetadata; use vortex_array::expr::BoundExpression; use vortex_error::VortexResult; use vortex_error::vortex_bail; +use vortex_mask::Mask; use vortex_session::registry::CachedId; use crate::plan::Plan; @@ -17,6 +19,10 @@ use crate::plan::PlanParts; use crate::plan::PlanRef; use crate::plan::PlanVTable; use crate::plan::check_child_count; +use crate::plan::exec::EvalNode; +use crate::plan::exec::ExecContext; +use crate::plan::exec::ExecNode; +use crate::plan::exec::Selection; use crate::plan::optimizer::PlanReduceRule; /// Applies an expression to the output of its child. @@ -113,6 +119,18 @@ impl PlanVTable for Eval { Cow::Owned(format!("child[{index}]")) } } + + fn exec( + plan: &Plan, + rows: Range, + mask: Mask, + _ctx: &ExecContext, + ) -> VortexResult> { + Ok(Box::new(EvalNode::new( + plan.clone(), + Selection::try_new(rows, mask)?, + ))) + } } fn validate_expression_child(expression: &BoundExpression, child: &PlanRef) -> VortexResult<()> { diff --git a/vortex-layout/src/plan/plans/row_idx.rs b/vortex-layout/src/plan/plans/row_idx.rs index 6b1edbbc433..9ee8e851c97 100644 --- a/vortex-layout/src/plan/plans/row_idx.rs +++ b/vortex-layout/src/plan/plans/row_idx.rs @@ -3,6 +3,7 @@ use std::fmt::Display; use std::fmt::Formatter; +use std::ops::Range; use vortex_array::EmptyMetadata; use vortex_array::dtype::DType; @@ -19,6 +20,7 @@ use vortex_array::scalar_fn::fns::pack::Pack as PackFn; use vortex_error::VortexResult; use vortex_error::vortex_ensure_eq; use vortex_error::vortex_err; +use vortex_mask::Mask; use vortex_session::registry::CachedId; use crate::layouts::row_idx::RowIdx as RowIdxFn; @@ -31,6 +33,10 @@ use crate::plan::PlanParts; use crate::plan::PlanRef; use crate::plan::PlanVTable; use crate::plan::check_child_count; +use crate::plan::exec::ExecContext; +use crate::plan::exec::ExecNode; +use crate::plan::exec::RowIdxNode; +use crate::plan::exec::Selection; use crate::plan::plans::pack::rewrite_partition_root; const ROW_IDX_PARTITION_NAME: &str = "row_idx"; @@ -86,6 +92,18 @@ impl PlanVTable for RowIdx { ) -> VortexResult<()> { check_child_count("RowIdx", children, 0) } + + fn exec( + _plan: &Plan, + rows: Range, + mask: Mask, + ctx: &ExecContext, + ) -> VortexResult> { + Ok(Box::new(RowIdxNode::new( + Selection::try_new(rows, mask)?, + ctx.row_offset(), + ))) + } } /// Plans an expression over a data source and its global row-index domain. diff --git a/vortex-layout/src/plan/plans/take.rs b/vortex-layout/src/plan/plans/take.rs index 95919090f3c..c3ddae8f670 100644 --- a/vortex-layout/src/plan/plans/take.rs +++ b/vortex-layout/src/plan/plans/take.rs @@ -2,12 +2,14 @@ // SPDX-FileCopyrightText: Copyright the Vortex contributors use std::borrow::Cow; +use std::ops::Range; use vortex_array::EmptyMetadata; use vortex_array::dtype::DType; use vortex_array::expr::ExactBoundExpr; use vortex_array::expr::label_bound_tree; use vortex_error::VortexResult; +use vortex_mask::Mask; use vortex_session::registry::CachedId; use crate::plan::Eval; @@ -19,6 +21,10 @@ use crate::plan::PlanParts; use crate::plan::PlanRef; use crate::plan::PlanVTable; use crate::plan::check_child_count; +use crate::plan::exec::ExecContext; +use crate::plan::exec::ExecNode; +use crate::plan::exec::Selection; +use crate::plan::exec::TakeNode; use crate::plan::optimizer::PlanParentReduceRule; const CODES: usize = 0; @@ -117,6 +123,18 @@ impl PlanVTable for Take { _ => Cow::Owned(format!("child[{index}]")), } } + + fn exec( + plan: &Plan, + rows: Range, + mask: Mask, + _ctx: &ExecContext, + ) -> VortexResult> { + Ok(Box::new(TakeNode::new( + plan.clone(), + Selection::try_new(rows, mask)?, + ))) + } } /// Pushes a strict, infallible boolean expression onto the dictionary values of a [`Take`].