From cc4c7705dd8cc69ac2ea6b789b5af797fb99dc9e Mon Sep 17 00:00:00 2001 From: Joe Isaacs Date: Thu, 8 Oct 2026 14:50:23 +0000 Subject: [PATCH 1/4] feat(layout): add a Filter plan over dense segment scans Add `Filter`, a plan operator with no predicate that keeps only the rows of the selection it is executed with. Its child produces every row of the same row domain, and values outside the selection are unspecified. Flat layouts now lower to a `Filter` over their `SegmentScan`, so a bare scan can stay dense and return every row of its range while the filter above it keeps the selected ones. Plan display snapshots are updated for the extra level. Signed-off-by: Joe Isaacs Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01NWYmfHTFtaEdhdyd3u5hv3 --- vortex-layout/src/plan/lower.rs | 6 +- vortex-layout/src/plan/mod.rs | 3 + vortex-layout/src/plan/plans/filter.rs | 89 ++++++++++++ vortex-layout/src/plan/plans/mod.rs | 4 + vortex-layout/src/plan/tests.rs | 190 ++++++++++++++++--------- 5 files changed, 223 insertions(+), 69 deletions(-) create mode 100644 vortex-layout/src/plan/plans/filter.rs diff --git a/vortex-layout/src/plan/lower.rs b/vortex-layout/src/plan/lower.rs index e8911734549..d04e3e128d4 100644 --- a/vortex-layout/src/plan/lower.rs +++ b/vortex-layout/src/plan/lower.rs @@ -25,6 +25,7 @@ use crate::layouts::list::VALIDITY_CHILD_INDEX; use crate::layouts::struct_::Struct; use crate::layouts::struct_::StructLayout; use crate::plan::ConcatPlan; +use crate::plan::FilterPlan; use crate::plan::ListPackPlan; use crate::plan::PackPlan; use crate::plan::PlanChildren; @@ -36,9 +37,12 @@ use crate::plan::TakePlan; /// /// The root operator is built immediately. Its child container owns a hidden clone of the source /// layout and lowers each child independently on first access. +/// +/// A flat layout lowers to a [`Filter`](crate::plan::Filter) over its segment scan, so the scan +/// returns only the rows it is executed with. pub fn lower(layout: &LayoutRef) -> VortexResult { if let Some(layout) = layout.as_opt::() { - return Ok(lower_flat(layout).into_plan()); + return Ok(FilterPlan::new(lower_flat(layout).into_plan()).into_plan()); } if let Some(layout) = layout.as_opt::() { return Ok(lower_chunked(layout)?.into_plan()); diff --git a/vortex-layout/src/plan/mod.rs b/vortex-layout/src/plan/mod.rs index b4782481b49..7e45533ed84 100644 --- a/vortex-layout/src/plan/mod.rs +++ b/vortex-layout/src/plan/mod.rs @@ -31,6 +31,9 @@ pub use plans::ConcatPlan; pub use plans::Eval; pub use plans::EvalData; pub use plans::EvalPlan; +pub use plans::Filter; +pub use plans::FilterData; +pub use plans::FilterPlan; pub use plans::ListPack; pub use plans::ListPackData; pub use plans::ListPackPlan; diff --git a/vortex-layout/src/plan/plans/filter.rs b/vortex-layout/src/plan/plans/filter.rs new file mode 100644 index 00000000000..da5907f3e80 --- /dev/null +++ b/vortex-layout/src/plan/plans/filter.rs @@ -0,0 +1,89 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use std::borrow::Cow; + +use vortex_array::EmptyMetadata; +use vortex_error::VortexResult; +use vortex_error::vortex_bail; +use vortex_error::vortex_err; +use vortex_session::registry::CachedId; + +use crate::plan::Plan; +use crate::plan::PlanChildren; +use crate::plan::PlanId; +use crate::plan::PlanParts; +use crate::plan::PlanRef; +use crate::plan::PlanVTable; +use crate::plan::check_child_count; + +/// Keeps only the selected rows of its child. +/// +/// A filter holds no predicate. The rows it keeps are the selection it is executed with; its child +/// produces every row of the same row domain, and values outside the selection are unspecified. +/// The filter's dtype and row domain are its child's. +#[derive(Clone, Debug)] +pub struct Filter; + +/// Operator-specific data for a [`Filter`] plan. +#[derive(Clone, Debug)] +pub struct FilterData; + +/// A plan that keeps only the selected rows of its child. +pub type FilterPlan = Plan; + +impl FilterPlan { + /// Creates a filter over `child`. + pub fn new(child: PlanRef) -> Self { + PlanParts { + vtable: Filter, + dtype: child.dtype().clone(), + row_count: child.row_count(), + children: vec![child].into(), + data: FilterData, + } + .into_typed() + } + + /// Returns the plan whose rows are filtered. + pub fn child_plan(&self) -> VortexResult { + self.child_required(0) + } +} + +impl PlanVTable for Filter { + type PlanData = FilterData; + type Metadata = EmptyMetadata; + + fn id(&self) -> PlanId { + static ID: CachedId = CachedId::new("vortex.plan.filter"); + *ID + } + + fn metadata(_plan: &Plan) -> Option { + Some(EmptyMetadata) + } + + fn with_children( + plan: &Plan, + children: &PlanChildren, + _data: &mut Self::PlanData, + ) -> VortexResult<()> { + check_child_count("Filter", children, 1)?; + let child = children + .get(0)? + .ok_or_else(|| vortex_err!("Filter child is absent"))?; + if child.dtype() != plan.dtype() || child.row_count() != plan.row_count() { + vortex_bail!("Filter child does not match the filter's dtype and row count"); + } + Ok(()) + } + + fn child_name(_plan: &Plan, index: usize) -> Cow<'_, str> { + if index == 0 { + Cow::Borrowed("child") + } else { + Cow::Owned(format!("child[{index}]")) + } + } +} diff --git a/vortex-layout/src/plan/plans/mod.rs b/vortex-layout/src/plan/plans/mod.rs index 7b6a1893d0b..b20731b8d34 100644 --- a/vortex-layout/src/plan/plans/mod.rs +++ b/vortex-layout/src/plan/plans/mod.rs @@ -3,6 +3,7 @@ mod concat; pub(crate) mod eval; +mod filter; mod list_pack; mod pack; mod row_idx; @@ -17,6 +18,9 @@ pub use eval::Eval; pub use eval::EvalData; pub(crate) use eval::EvalIdentityRule; pub use eval::EvalPlan; +pub use filter::Filter; +pub use filter::FilterData; +pub use filter::FilterPlan; pub use list_pack::ListPack; pub use list_pack::ListPackData; pub use list_pack::ListPackPlan; diff --git a/vortex-layout/src/plan/tests.rs b/vortex-layout/src/plan/tests.rs index 783eb455327..fccc06ebb2f 100644 --- a/vortex-layout/src/plan/tests.rs +++ b/vortex-layout/src/plan/tests.rs @@ -94,12 +94,14 @@ fn unsupported_layout_has_no_plan() -> VortexResult<()> { } #[test] -fn flat_plan_has_no_children() -> VortexResult<()> { +fn flat_plan_is_a_filtered_segment_scan_without_children() -> VortexResult<()> { let plan = make_plan(flat(3, primitive(PType::I32, Nullability::NonNullable), 0))?; - assert!(plan.is::()); - assert_eq!(plan.child_count(), 0); - assert!(plan.child(0)?.is_none()); + assert!(plan.is::()); + let scan = child_of(&plan, 0)?; + assert!(scan.is::()); + assert_eq!(scan.child_count(), 0); + assert!(scan.child(0)?.is_none()); Ok(()) } @@ -323,7 +325,7 @@ fn optimize_drops_identity_expressions() -> VortexResult<()> { let expression = root().bind(child.dtype())?; let plan: PlanRef = EvalPlan::try_new(expression, child)?.into_plan(); - assert!(optimize(plan)?.is::()); + assert!(optimize(plan)?.is::()); Ok(()) } @@ -346,8 +348,8 @@ fn optimize_rewrites_nested_children() -> VortexResult<()> { let optimized = optimize(wrapped)?; assert!(optimized.is::()); - assert!(child_of(&optimized, 0)?.is::()); - assert!(child_of(&optimized, 1)?.is::()); + assert!(child_of(&optimized, 0)?.is::()); + assert!(child_of(&optimized, 1)?.is::()); Ok(()) } @@ -371,11 +373,13 @@ fn plan_display_matches_array_tree_display_shape() -> VortexResult<()> { let plan: PlanRef = plan.into_plan(); assert_eq!(plan.to_string(), "vortex.plan.eval(i32, rows=3) expr=$.a"); - insta::assert_snapshot!(plan.display_tree(), @r" + insta::assert_snapshot!(plan.display_tree(), @" root: vortex.plan.eval(i32, rows=3) expr=$.a child: vortex.plan.pack({a=i32, b=i32}, rows=3) - a: vortex.plan.segment_scan(i32, rows=3) - b: vortex.plan.segment_scan(i32, rows=3) + a: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) + b: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); struct DepthExtractor; @@ -391,11 +395,13 @@ fn plan_display_matches_array_tree_display_shape() -> VortexResult<()> { } } - insta::assert_snapshot!(plan.tree_display_builder().with(DepthExtractor), @r" + insta::assert_snapshot!(plan.tree_display_builder().with(DepthExtractor), @" root: depth=0 child: depth=1 a: depth=2 + child: depth=3 b: depth=2 + child: depth=3 "); let nullable_fields = StructFields::from_iter([ @@ -413,11 +419,14 @@ fn plan_display_matches_array_tree_display_shape() -> VortexResult<()> { ) .into_layout(); let nullable = make_plan(nullable_layout)?; - insta::assert_snapshot!(nullable.tree_display_builder(), @r" + insta::assert_snapshot!(nullable.tree_display_builder(), @" root: a: + child: b: + child: validity: + child: "); Ok(()) } @@ -433,10 +442,12 @@ fn chunked_plan_display_names_chunks() -> VortexResult<()> { .into_layout(); let plan = make_plan(layout)?; - insta::assert_snapshot!(plan.display_tree(), @r" + insta::assert_snapshot!(plan.display_tree(), @" root: vortex.plan.concat(i32, rows=3) - chunks[0]: vortex.plan.segment_scan(i32, rows=2) - chunks[1]: vortex.plan.segment_scan(i32, rows=1) + chunks[0]: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) + chunks[1]: vortex.plan.filter(i32, rows=1) + child: vortex.plan.segment_scan(i32, rows=1) "); Ok(()) } @@ -450,10 +461,12 @@ fn dict_plan_display_names_logical_children() -> VortexResult<()> { .into_layout(); let plan = make_plan(layout)?; - insta::assert_snapshot!(plan.display_tree(), @r" + insta::assert_snapshot!(plan.display_tree(), @" root: vortex.plan.take(i32, rows=3) - codes: vortex.plan.segment_scan(u8, rows=3) - values: vortex.plan.segment_scan(i32, rows=2) + codes: vortex.plan.filter(u8, rows=3) + child: vortex.plan.segment_scan(u8, rows=3) + values: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) "); Ok(()) } @@ -471,10 +484,12 @@ fn list_plan_display_handles_optional_validity() -> VortexResult<()> { .into_layout(); let non_nullable = make_plan(non_nullable_layout)?; - insta::assert_snapshot!(non_nullable.display_tree(), @r" + insta::assert_snapshot!(non_nullable.display_tree(), @" root: vortex.plan.list_pack(list(i32), rows=2) - elements: vortex.plan.segment_scan(i32, rows=4) - offsets: vortex.plan.segment_scan(u32, rows=3) + elements: vortex.plan.filter(i32, rows=4) + child: vortex.plan.segment_scan(i32, rows=4) + offsets: vortex.plan.filter(u32, rows=3) + child: vortex.plan.segment_scan(u32, rows=3) "); let nullable_layout = ListLayout::new( @@ -486,11 +501,14 @@ fn list_plan_display_handles_optional_validity() -> VortexResult<()> { .into_layout(); let nullable = make_plan(nullable_layout)?; - insta::assert_snapshot!(nullable.display_tree(), @r" + insta::assert_snapshot!(nullable.display_tree(), @" root: vortex.plan.list_pack(list(i32)?, rows=2) - elements: vortex.plan.segment_scan(i32, rows=4) - offsets: vortex.plan.segment_scan(u32, rows=3) - validity: vortex.plan.segment_scan(bool, rows=2) + elements: vortex.plan.filter(i32, rows=4) + child: vortex.plan.segment_scan(i32, rows=4) + offsets: vortex.plan.filter(u32, rows=3) + child: vortex.plan.segment_scan(u32, rows=3) + validity: vortex.plan.filter(bool, rows=2) + child: vortex.plan.segment_scan(bool, rows=2) "); Ok(()) } @@ -546,11 +564,14 @@ fn expression_partitions_across_row_idx_and_struct() -> VortexResult<()> { child: vortex.plan.eval({child_0=bool, child_1=bool}, rows=3) expr=pack(child_0: $.a, child_1: $.b) child: vortex.plan.pack({a=bool, b=bool}, rows=3) a: vortex.plan.take(bool, rows=3) - codes: vortex.plan.segment_scan(u8, rows=3) + codes: vortex.plan.filter(u8, rows=3) + child: vortex.plan.segment_scan(u8, rows=3) values: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) b: vortex.plan.eval(bool, rows=3) expr=($ > 7i32) - child: vortex.plan.segment_scan(i32, rows=3) + child: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); Ok(()) } @@ -600,16 +621,18 @@ fn row_idx_and_data_expression_pushes_data_into_chunks() -> VortexResult<()> { let expression = and(gt(row_idx(), lit(10_u64)), gt(root(), lit(5_i32))); let plan = optimize(make_row_idx_plan(expression, child)?)?; - insta::assert_snapshot!(plan.display_tree(), @r" + insta::assert_snapshot!(plan.display_tree(), @" root: vortex.plan.eval(bool, rows=3) expr=($.row_idx and $.child) child: vortex.plan.pack({row_idx=bool, child=bool}, rows=3) row_idx: vortex.plan.eval(bool, rows=3) expr=($ > 10u64) child: vortex.plan.row_idx(u64, rows=3) child: vortex.plan.concat(bool, rows=3) chunks[0]: vortex.plan.eval(bool, rows=1) expr=($ > 5i32) - child: vortex.plan.segment_scan(i32, rows=1) + child: vortex.plan.filter(i32, rows=1) + child: vortex.plan.segment_scan(i32, rows=1) chunks[1]: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) "); Ok(()) } @@ -635,17 +658,22 @@ fn expression_pushes_through_struct_field_and_dictionary_values() -> VortexResul root: vortex.plan.eval(bool, rows=3) expr=($.a > 5i32) child: vortex.plan.pack({a=i32, b=i32}, rows=3) a: vortex.plan.take(i32, rows=3) - codes: vortex.plan.segment_scan(u8, rows=3) - values: vortex.plan.segment_scan(i32, rows=2) - b: vortex.plan.segment_scan(i32, rows=3) + codes: vortex.plan.filter(u8, rows=3) + child: vortex.plan.segment_scan(u8, rows=3) + values: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) + b: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); let optimized = optimize(plan)?; insta::assert_snapshot!(optimized.display_tree(), @" root: vortex.plan.take(bool, rows=3) - codes: vortex.plan.segment_scan(u8, rows=3) + codes: vortex.plan.filter(u8, rows=3) + child: vortex.plan.segment_scan(u8, rows=3) values: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) "); Ok(()) } @@ -680,21 +708,28 @@ fn expression_pushes_through_struct_field_with_heterogeneous_chunks() -> VortexR child: vortex.plan.pack({a=i32, b=i32}, rows=5) a: vortex.plan.concat(i32, rows=5) chunks[0]: vortex.plan.take(i32, rows=3) - codes: vortex.plan.segment_scan(u8, rows=3) - values: vortex.plan.segment_scan(i32, rows=2) - chunks[1]: vortex.plan.segment_scan(i32, rows=2) - b: vortex.plan.segment_scan(i32, rows=5) + codes: vortex.plan.filter(u8, rows=3) + child: vortex.plan.segment_scan(u8, rows=3) + values: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) + chunks[1]: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) + b: vortex.plan.filter(i32, rows=5) + child: vortex.plan.segment_scan(i32, rows=5) "); let optimized = optimize(plan)?; insta::assert_snapshot!(optimized.display_tree(), @" root: vortex.plan.concat(bool, rows=5) chunks[0]: vortex.plan.take(bool, rows=3) - codes: vortex.plan.segment_scan(u8, rows=3) + codes: vortex.plan.filter(u8, rows=3) + child: vortex.plan.segment_scan(u8, rows=3) values: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) chunks[1]: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) "); Ok(()) } @@ -728,9 +763,10 @@ fn expression_pushes_through_nested_struct_fields_in_one_pass() -> VortexResult< let plan = make_eval(expression, make_plan(layout)?)?.into_plan(); let optimized = optimize(plan)?; - insta::assert_snapshot!(optimized.display_tree(), @r" + insta::assert_snapshot!(optimized.display_tree(), @" root: vortex.plan.eval(bool, rows=3) expr=($ > 5i32) - child: vortex.plan.segment_scan(i32, rows=3) + child: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); Ok(()) } @@ -756,17 +792,19 @@ fn expression_pushes_through_single_field_nested_structs() -> VortexResult<()> { let expression = gt(get_item("b", get_item("a", root())), lit(5_i32)); let plan = make_eval(expression, make_plan(layout)?)?.into_plan(); - insta::assert_snapshot!(plan.display_tree(), @r" + insta::assert_snapshot!(plan.display_tree(), @" root: vortex.plan.eval(bool, rows=3) expr=($.a.b > 5i32) child: vortex.plan.pack({a={b=i32}}, rows=3) a: vortex.plan.pack({b=i32}, rows=3) - b: vortex.plan.segment_scan(i32, rows=3) + b: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); let optimized = optimize(plan)?; - insta::assert_snapshot!(optimized.display_tree(), @r" + insta::assert_snapshot!(optimized.display_tree(), @" root: vortex.plan.eval(bool, rows=3) expr=($ > 5i32) - child: vortex.plan.segment_scan(i32, rows=3) + child: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); Ok(()) } @@ -800,9 +838,10 @@ fn expression_pushes_through_three_single_field_structs() -> VortexResult<()> { let plan = make_eval(expression, make_plan(layout)?)?.into_plan(); let optimized = optimize(plan)?; - insta::assert_snapshot!(optimized.display_tree(), @r" + insta::assert_snapshot!(optimized.display_tree(), @" root: vortex.plan.eval(bool, rows=3) expr=($ > 5i32) - child: vortex.plan.segment_scan(i32, rows=3) + child: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); Ok(()) } @@ -847,15 +886,18 @@ fn compound_expression_pushes_through_nested_struct_and_dictionary() -> VortexRe let plan = make_eval(expression, make_plan(layout)?)?.into_plan(); let optimized = optimize(plan)?; - insta::assert_snapshot!(optimized.display_tree(), @r" + insta::assert_snapshot!(optimized.display_tree(), @" root: vortex.plan.eval(bool, rows=3) expr=($.b and $.c) child: vortex.plan.pack({b=bool, c=bool}, rows=3) b: vortex.plan.take(bool, rows=3) - codes: vortex.plan.segment_scan(u8, rows=3) + codes: vortex.plan.filter(u8, rows=3) + child: vortex.plan.segment_scan(u8, rows=3) values: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) c: vortex.plan.eval(bool, rows=3) expr=($ > 7i32) - child: vortex.plan.segment_scan(i32, rows=3) + child: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); Ok(()) } @@ -894,11 +936,13 @@ fn expression_pushes_through_dictionary_of_struct_values() -> VortexResult<()> { let plan = make_eval(expression, make_plan(layout)?)?.into_plan(); let optimized = optimize(plan)?; - insta::assert_snapshot!(optimized.display_tree(), @r" + insta::assert_snapshot!(optimized.display_tree(), @" root: vortex.plan.take(bool, rows=3) - codes: vortex.plan.segment_scan(u8, rows=3) + codes: vortex.plan.filter(u8, rows=3) + child: vortex.plan.segment_scan(u8, rows=3) values: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) "); Ok(()) } @@ -938,10 +982,14 @@ fn multi_field_struct_expression_pushes_into_each_field() -> VortexResult<()> { root: vortex.plan.eval(bool, rows=3) expr=(($.a > 5i32) and ($.b > 7i32)) child: vortex.plan.pack({a=i32, b=i32, c=i32}, rows=3) a: vortex.plan.take(i32, rows=3) - codes: vortex.plan.segment_scan(u8, rows=3) - values: vortex.plan.segment_scan(i32, rows=2) - b: vortex.plan.segment_scan(i32, rows=3) - c: vortex.plan.segment_scan(i32, rows=3) + codes: vortex.plan.filter(u8, rows=3) + child: vortex.plan.segment_scan(u8, rows=3) + values: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) + b: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) + c: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); let optimized = optimize(plan)?; @@ -949,11 +997,14 @@ fn multi_field_struct_expression_pushes_into_each_field() -> VortexResult<()> { root: vortex.plan.eval(bool, rows=3) expr=($.a and $.b) child: vortex.plan.pack({a=bool, b=bool}, rows=3) a: vortex.plan.take(bool, rows=3) - codes: vortex.plan.segment_scan(u8, rows=3) + codes: vortex.plan.filter(u8, rows=3) + child: vortex.plan.segment_scan(u8, rows=3) values: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) b: vortex.plan.eval(bool, rows=3) expr=($ > 7i32) - child: vortex.plan.segment_scan(i32, rows=3) + child: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); let reoptimized = optimize(optimized.clone())?; assert!(PlanRef::ptr_eq(&optimized, &reoptimized)); @@ -1022,9 +1073,12 @@ fn multi_field_struct_expression_keeps_cross_field_refinement() -> VortexResult< root: vortex.plan.eval(bool, rows=3) expr=(($.a + $.b) > 10i32) child: vortex.plan.pack({a=i32, b=i32}, rows=3) a: vortex.plan.take(i32, rows=3) - codes: vortex.plan.segment_scan(u8, rows=3) - values: vortex.plan.segment_scan(i32, rows=2) - b: vortex.plan.segment_scan(i32, rows=3) + codes: vortex.plan.filter(u8, rows=3) + child: vortex.plan.segment_scan(u8, rows=3) + values: vortex.plan.filter(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) + b: vortex.plan.filter(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); assert_eq!( optimized.display_tree().to_string(), From c31fbcf6c1703f84290fe56362f6e31a7805db08 Mon Sep 17 00:00:00 2001 From: Joe Isaacs Date: Thu, 8 Oct 2026 20:14:36 +0000 Subject: [PATCH 2/4] feat(layout): add SegmentScan, Filter, Concat, and Pack exec nodes Implement `PlanVTable::exec` for the operators a projection of a struct of chunked columns needs: - `SegmentScanNode` requests its segment at start, decodes the delivered bytes, and emits the node's rows as one dense array. A `Filter` over a scan fuses into the same node with a filter mask, so the array holds exactly the selected rows. Decoded segments are shared through the graph's `DecodeCache`. - `FilterNode` keeps the selected rows of any other child's dense arrays as they stream through, tracking its position with a cursor. - `ConcatNode` spawns only the chunks overlapping its rows, one port each in chunk order, and emits the current chunk's rows as they arrive. Chunks whose reads land early wait in their ports. - `PackNode` zips its fields as their rows arrive: whenever every port has rows, it takes as many as the shortest holds from each and emits one struct. Fields chunked alike stream one struct per chunk. The layout tests drive lowered plans over in-memory segments with scripted IO ordering and check what reaches the root. The scheduling tests add hand-built sources under the real operators, including random trees over random views with reads completing in random order, checked against a model of the same tree. 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/concat.rs | 97 ++++ vortex-layout/src/plan/exec/filter.rs | 98 ++++ vortex-layout/src/plan/exec/mod.rs | 13 + vortex-layout/src/plan/exec/pack.rs | 111 ++++ .../src/plan/exec/scheduling_tests.rs | 324 ++++++++++- vortex-layout/src/plan/exec/segment_scan.rs | 118 ++++ vortex-layout/src/plan/exec/selection.rs | 55 ++ vortex-layout/src/plan/exec/synthetic.rs | 2 +- vortex-layout/src/plan/exec/tests.rs | 548 ++++++++++++++++++ vortex-layout/src/plan/plans/concat.rs | 18 + vortex-layout/src/plan/plans/filter.rs | 33 ++ vortex-layout/src/plan/plans/pack.rs | 18 + vortex-layout/src/plan/plans/segment_scan.rs | 22 + 13 files changed, 1454 insertions(+), 3 deletions(-) create mode 100644 vortex-layout/src/plan/exec/concat.rs create mode 100644 vortex-layout/src/plan/exec/filter.rs create mode 100644 vortex-layout/src/plan/exec/pack.rs create mode 100644 vortex-layout/src/plan/exec/segment_scan.rs create mode 100644 vortex-layout/src/plan/exec/selection.rs create mode 100644 vortex-layout/src/plan/exec/tests.rs diff --git a/vortex-layout/src/plan/exec/concat.rs b/vortex-layout/src/plan/exec/concat.rs new file mode 100644 index 00000000000..e6dbd4ec00d --- /dev/null +++ b/vortex-layout/src/plan/exec/concat.rs @@ -0,0 +1,97 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use std::ops::Range; + +use vortex_error::VortexResult; + +use crate::plan::ConcatPlan; +use crate::plan::exec::ExecNode; +use crate::plan::exec::NodeState; +use crate::plan::exec::StepCx; +use crate::plan::exec::selection::Selection; + +/// Emits its chunks' rows in chunk order. +/// +/// Only chunks overlapping the selection's rows are spawned, one port each in chunk order, and +/// a chunk whose slice of the mask is all false is skipped. Chunks are read independently and +/// may finish in any order; the node emits whatever the current chunk has produced, moves on +/// once that chunk has finished, and leaves what later chunks produced early waiting in their +/// ports. Nothing above it sees the disorder. +pub(crate) struct ConcatNode { + plan: ConcatPlan, + selection: Selection, + /// The port whose chunk is being emitted. Ports are numbered in chunk order. + current: usize, + ports: usize, +} + +impl ConcatNode { + pub(crate) fn new(plan: ConcatPlan, selection: Selection) -> Self { + Self { + plan, + selection, + current: 0, + ports: 0, + } + } + + fn chunk_rows(&self, index: usize) -> Range { + let offsets = self.plan.row_offsets(); + let end = offsets + .get(index + 1) + .copied() + .unwrap_or_else(|| self.plan.row_count()); + offsets[index]..end + } +} + +impl ExecNode for ConcatNode { + fn start(&mut self, cx: &mut StepCx<'_>) -> VortexResult { + let rows = self.selection.rows().clone(); + let offsets = self.plan.row_offsets(); + // Every split visits a small part of a file; skip the chunks before and after it. + let first = offsets + .partition_point(|&offset| offset <= rows.start) + .saturating_sub(1); + let end = offsets.partition_point(|&offset| offset < rows.end); + for index in first..end { + let chunk = self.chunk_rows(index); + let local = rows.start.max(chunk.start)..rows.end.min(chunk.end); + if local.start >= local.end { + continue; + } + let mask = self.selection.slice(&local)?; + if mask.all_false() { + continue; + } + cx.spawn( + self.ports, + self.plan.child_required(index)?, + local.start - chunk.start..local.end - chunk.start, + mask, + ); + self.ports += 1; + } + if self.ports == 0 { + return Ok(NodeState::Done); + } + Ok(NodeState::Wait) + } + + fn compute(&mut self, cx: &mut StepCx<'_>) -> VortexResult { + while self.current < self.ports { + let input = cx.input(self.current); + let arrays = input.take_all(); + let finished = input.finished(); + for array in arrays { + cx.emit(array); + } + if !finished { + return Ok(NodeState::Wait); + } + self.current += 1; + } + Ok(NodeState::Done) + } +} diff --git a/vortex-layout/src/plan/exec/filter.rs b/vortex-layout/src/plan/exec/filter.rs new file mode 100644 index 00000000000..e2294c074e1 --- /dev/null +++ b/vortex-layout/src/plan/exec/filter.rs @@ -0,0 +1,98 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use vortex_array::ArrayRef; +use vortex_array::Canonical; +use vortex_array::IntoArray; +use vortex_array::VortexSessionExecute; +use vortex_error::VortexResult; +use vortex_mask::Mask; +use vortex_session::VortexSession; + +use crate::plan::FilterPlan; +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; + +/// The selected fraction of an array at or above which a predicate runs over the whole array. +const EXPR_EVAL_THRESHOLD: f64 = 0.2; + +/// Keeps the selected rows of the dense arrays its child produces. +/// +/// The child runs over the same rows with the selection as its care hint, and returns every row +/// of its range in order. A cursor tracks how far along the range the arrays taken so far reach, +/// so each one is filtered by its own slice of the mask as it arrives. +pub(crate) struct FilterNode { + plan: FilterPlan, + selection: Selection, + session: VortexSession, + /// Rows of the child's range consumed so far. + cursor: usize, +} + +impl FilterNode { + pub(crate) fn new(plan: FilterPlan, selection: Selection, session: VortexSession) -> Self { + Self { + plan, + selection, + session, + cursor: 0, + } + } +} + +/// Keeps the rows of `array` that `mask` selects. `predicate` says the array is a predicate's +/// result. +pub(crate) fn keep_selected( + array: ArrayRef, + mask: Mask, + predicate: bool, + session: &VortexSession, +) -> VortexResult { + if mask.all_true() { + return Ok(array); + } + if predicate && mask.density() >= EXPR_EVAL_THRESHOLD { + // A predicate over a mostly selected array runs over every row and its result is + // filtered, as the default scan's flat reader does. Filtering lazily would push the + // filter back through the predicate onto the encoded input, which for some encodings + // costs more than the predicate. + let mut ctx = session.create_execution_ctx(); + return array + .execute::(&mut ctx)? + .into_array() + .filter(mask); + } + array.filter(mask) +} + +impl ExecNode for FilterNode { + 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 { + let predicate = self.plan.dtype().is_boolean(); + for array in cx.input(CHILD).take_all() { + let mask = self + .selection + .mask() + .slice(self.cursor..self.cursor + array.len()); + self.cursor += array.len(); + cx.emit(keep_selected(array, mask, predicate, &self.session)?); + } + 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 5cb10f6e569..127f5898b8d 100644 --- a/vortex-layout/src/plan/exec/mod.rs +++ b/vortex-layout/src/plan/exec/mod.rs @@ -36,6 +36,11 @@ //! merge pipelines, and nodes that wait for a port to close, [`Ready::Closed`] and //! [`Ready::AllClosed`], are barriers between them, as a join build is in a query engine. +mod concat; +mod filter; +mod pack; +mod segment_scan; +mod selection; pub mod synthetic; use std::collections::VecDeque; @@ -668,5 +673,13 @@ impl From> for Effects { } } +pub(crate) use concat::ConcatNode; +pub(crate) use filter::FilterNode; +pub(crate) use pack::PackNode; +pub(crate) use segment_scan::SegmentScanNode; +pub(crate) use selection::Selection; + #[cfg(test)] mod scheduling_tests; +#[cfg(test)] +mod tests; diff --git a/vortex-layout/src/plan/exec/pack.rs b/vortex-layout/src/plan/exec/pack.rs new file mode 100644 index 00000000000..cedec0eb215 --- /dev/null +++ b/vortex-layout/src/plan/exec/pack.rs @@ -0,0 +1,111 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::arrays::StructArray; +use vortex_array::dtype::DType; +use vortex_array::dtype::Nullability; +use vortex_array::validity::Validity; +use vortex_error::VortexResult; +use vortex_error::vortex_err; + +use crate::plan::PackPlan; +use crate::plan::exec::ExecNode; +use crate::plan::exec::NodeState; +use crate::plan::exec::StepCx; +use crate::plan::exec::selection::Selection; +use crate::plan::exec::selection::join; + +/// Zips its fields into struct arrays as their rows arrive. +/// +/// Every field runs over the same rows and selection, so the fields' outputs have the same +/// length and line up row for row. Whenever every port has rows available, the node takes as +/// many rows as the shortest port holds from each of them and emits one struct. Fields chunked +/// alike therefore stream one struct per chunk; fields chunked differently are sliced at the +/// boundaries they share. +pub(crate) struct PackNode { + plan: PackPlan, + selection: Selection, + ports: usize, +} + +impl PackNode { + pub(crate) fn new(plan: PackPlan, selection: Selection) -> Self { + let ports = plan.children().len(); + Self { + plan, + selection, + ports, + } + } + + /// Builds a struct of `len` rows from `arrays`, one per port, the last being the validity + /// when the struct is nullable. + fn assemble(&self, mut arrays: Vec, len: usize) -> VortexResult { + let validity = if self.plan.dtype().is_nullable() { + Validity::Array( + arrays + .pop() + .ok_or_else(|| vortex_err!("Nullable Pack is missing its validity port"))?, + ) + } else { + Validity::NonNullable + }; + StructArray::try_new_with_dtype(arrays, self.plan.fields().clone(), len, validity) + } + + fn field_dtype(&self, port: usize) -> DType { + self.plan + .fields() + .field_by_index(port) + .unwrap_or(DType::Bool(Nullability::NonNullable)) + } +} + +impl ExecNode for PackNode { + fn start(&mut self, cx: &mut StepCx<'_>) -> VortexResult { + if self.ports == 0 { + // A struct with no fields still has rows, so no child can carry them. + let len = self.selection.mask().true_count(); + cx.emit(self.assemble(Vec::new(), len)?.into_array()); + return Ok(NodeState::Done); + } + if self.selection.mask().all_false() { + return Ok(NodeState::Done); + } + for port in 0..self.ports { + cx.spawn( + port, + self.plan.child_required(port)?, + self.selection.rows().clone(), + self.selection.mask().clone(), + ); + } + Ok(NodeState::Wait) + } + + fn compute(&mut self, cx: &mut StepCx<'_>) -> VortexResult { + loop { + let len = cx + .inputs() + .iter() + .map(|input| input.available()) + .min() + .unwrap_or(0); + if len == 0 { + break; + } + let mut arrays = Vec::with_capacity(self.ports); + for port in 0..self.ports { + let dtype = self.field_dtype(port); + arrays.push(join(&dtype, cx.input(port).take(len)?)?); + } + cx.emit(self.assemble(arrays, len)?.into_array()); + } + if cx.all_finished() { + return Ok(NodeState::Done); + } + Ok(NodeState::Wait) + } +} diff --git a/vortex-layout/src/plan/exec/scheduling_tests.rs b/vortex-layout/src/plan/exec/scheduling_tests.rs index 6fb2ee452f7..afe0fe27cc2 100644 --- a/vortex-layout/src/plan/exec/scheduling_tests.rs +++ b/vortex-layout/src/plan/exec/scheduling_tests.rs @@ -3,8 +3,8 @@ //! Tests of the graph's scheduling, driven through hand-built nodes over in-memory segments. //! -//! Every source emits its own row indices, so what reaches the root can be checked against a -//! sequence. +//! Every source emits its own row indices, so whatever the tree does to them can be checked +//! against a model of the same tree built from plain arrays. #![allow(clippy::cast_possible_truncation)] @@ -12,14 +12,21 @@ use std::ops::Range; use std::sync::Arc; use parking_lot::Mutex; +use rand::RngExt; +use rand::SeedableRng; +use rand::rngs::StdRng; use rstest::rstest; use vortex_array::ArrayRef; use vortex_array::IntoArray; use vortex_array::VortexSessionExecute; use vortex_array::arrays::ChunkedArray; use vortex_array::arrays::PrimitiveArray; +use vortex_array::arrays::StructArray; use vortex_array::assert_arrays_eq; use vortex_array::dtype::DType; +use vortex_array::dtype::Nullability; +use vortex_array::dtype::StructFields; +use vortex_array::validity::Validity; use vortex_buffer::Buffer; use vortex_error::VortexResult; use vortex_error::vortex_err; @@ -28,7 +35,11 @@ use vortex_mask::Mask; use super::synthetic::Probe; use super::synthetic::RowSource; use super::synthetic::empty_segment; +use super::synthetic::row_dtype; use super::*; +use crate::plan::ConcatPlan; +use crate::plan::FilterPlan; +use crate::plan::PackPlan; use crate::test::SESSION; #[derive(Debug, PartialEq, Eq)] @@ -94,6 +105,10 @@ fn fifo(_: &[IoRequest]) -> usize { 0 } +fn lifo(inflight: &[IoRequest]) -> usize { + inflight.len() - 1 +} + /// Delivers segments in the given order. fn scripted(order: &[u32]) -> impl FnMut(&[IoRequest]) -> usize + '_ { let mut next = order.iter(); @@ -136,6 +151,39 @@ fn segment(id: u32) -> Option { Some(SegmentId::from(id)) } +fn fields(dtypes: impl IntoIterator) -> StructFields { + StructFields::from_iter( + dtypes + .into_iter() + .enumerate() + .map(|(i, dtype)| (format!("f{i}"), dtype)), + ) +} + +fn struct_of(arrays: Vec, len: usize) -> VortexResult { + let names = fields(arrays.iter().map(|array| array.dtype().clone())) + .names() + .clone(); + Ok(StructArray::try_new(names, arrays, len, Validity::NonNullable)?.into_array()) +} + +fn pack(children: Vec) -> VortexResult { + let row_count = children[0].row_count(); + Ok(PackPlan::try_new( + fields(children.iter().map(|child| child.dtype().clone())), + Nullability::NonNullable, + row_count, + children, + None, + )? + .into_plan()) +} + +fn concat(children: Vec) -> VortexResult { + let dtype = children[0].dtype().clone(); + Ok(ConcatPlan::try_new(dtype, children)?.into_plan()) +} + #[rstest] #[case::one_piece(vec![], false)] #[case::many_pieces(vec![3, 1, 4, 1, 5], false)] @@ -151,6 +199,62 @@ fn source_pieces_come_out_in_order( assert_rows(indices((5..15).filter(|i| *i != 9)), run.arrays) } +/// Chunks whose reads land in reverse order come out in chunk order, and nothing comes out +/// before the first chunk has landed. +#[rstest] +#[case::fifo(&[0, 1, 2, 3])] +#[case::lifo(&[3, 2, 1, 0])] +#[case::first_last(&[1, 2, 3, 0])] +fn concat_reorders_out_of_order_reads(#[case] order: &[u32]) -> VortexResult<()> { + let chunks = (0..4) + .map(|i| RowSource::plan(10, vec![4, 6], segment(i), false, false)) + .collect(); + let plan = concat(chunks)?; + let run = run(&plan, 0..40, Mask::new_true(40), scripted(order))?; + + let first_piece = run + .events + .iter() + .position(|e| matches!(e, Event::Piece(_))) + .expect("something came out"); + let first_chunk_landed = run + .events + .iter() + .position(|e| *e == Event::Delivered(0)) + .expect("chunk 0 landed"); + assert!(first_piece > first_chunk_landed); + // Each chunk's own rows are local to the chunk. + assert_rows(indices((0..4).flat_map(|_| 0..10)), run.arrays) +} + +/// Fields cut at different places zip into structs at the boundaries they share, and the +/// structs come out in row order under any delivery order. +#[rstest] +fn pack_zips_misaligned_fields(#[values(true, false)] reverse: bool) -> VortexResult<()> { + let a = concat(vec![ + RowSource::plan(7, vec![], segment(0), false, false), + RowSource::plan(13, vec![5], segment(1), false, false), + ])?; + let b = RowSource::plan(20, vec![10], segment(2), false, false); + let c = concat(vec![ + RowSource::plan(3, vec![], segment(3), false, false), + RowSource::plan(17, vec![], segment(4), false, false), + ])?; + let plan = pack(vec![a, b, c])?; + let pick: fn(&[IoRequest]) -> usize = if reverse { lifo } else { fifo }; + let run = run(&plan, 0..20, Mask::new_true(20), pick)?; + + let expected = struct_of( + vec![ + indices((0..7).chain(0..13)), + indices(0..20), + indices((0..3).chain(0..17)), + ], + 20, + )?; + assert_rows(expected, run.arrays) +} + /// A node runs only when its readiness rule holds, whatever order its inputs finish in. #[rstest] #[case::any(Ready::Any)] @@ -180,6 +284,18 @@ fn probe_runs_only_when_ready(#[case] ready: Ready) -> VortexResult<()> { Ok(()) } +/// A filter over a dense source keeps the selected rows as the pieces stream through, however +/// the source cuts them. +#[rstest] +#[case::one_piece(vec![])] +#[case::odd_cuts(vec![1, 3, 2, 7])] +fn filter_streams_dense_pieces(#[case] cut: Vec) -> VortexResult<()> { + let plan = FilterPlan::new(RowSource::plan(20, cut, None, true, true)).into_plan(); + let mask = Mask::from_iter((0..14).map(|i| i % 3 == 0)); + let run = run(&plan, 3..17, mask.clone(), fifo)?; + assert_rows(indices(3..17).filter(mask)?, run.arrays) +} + /// A source that yields between pieces is run once per piece, and its pieces reach the root in /// order, one per run. #[test] @@ -190,6 +306,210 @@ fn yielding_source_is_run_once_per_piece() -> VortexResult<()> { assert_rows(indices(0..12), run.arrays) } +/// A tree of operators over sources, mirrored as plain arrays. +#[derive(Debug)] +enum Tree { + /// A source of this many rows, cut into these pieces, read from this segment, yielding. + Source { + rows: u64, + cut: Vec, + segment: Option, + yielding: bool, + }, + Concat(Vec), + Pack(Vec), + /// A filter over a dense source. + Filter(Box), +} + +impl Tree { + fn rows(&self) -> u64 { + match self { + Tree::Source { rows, .. } => *rows, + Tree::Concat(children) => children.iter().map(Tree::rows).sum(), + Tree::Pack(children) => children[0].rows(), + Tree::Filter(child) => child.rows(), + } + } + + fn dtype(&self) -> DType { + match self { + Tree::Pack(children) => DType::Struct( + fields(children.iter().map(Tree::dtype)), + Nullability::NonNullable, + ), + _ => row_dtype(), + } + } + + fn plan(&self, dense: bool) -> VortexResult { + Ok(match self { + Tree::Source { + rows, + cut, + segment, + yielding, + } => RowSource::plan( + *rows, + cut.clone(), + segment.map(SegmentId::from), + *yielding, + dense, + ), + Tree::Concat(children) => concat( + children + .iter() + .map(|child| child.plan(false)) + .collect::>()?, + )?, + Tree::Pack(children) => pack( + children + .iter() + .map(|child| child.plan(false)) + .collect::>()?, + )?, + Tree::Filter(child) => FilterPlan::new(child.plan(true)?).into_plan(), + }) + } + + /// The selected rows of `rows`, as the tree's plan should produce them. + fn expected(&self, rows: Range, mask: &Mask) -> VortexResult { + match self { + Tree::Source { .. } | Tree::Filter(_) => indices(rows).filter(mask.clone()), + Tree::Concat(children) => { + let mut parts = Vec::new(); + let mut offset = 0; + for child in children { + let chunk = offset..offset + child.rows(); + offset = chunk.end; + let local = rows.start.max(chunk.start)..rows.end.min(chunk.end); + if local.start >= local.end { + continue; + } + let mask = mask.slice( + (local.start - rows.start) as usize..(local.end - rows.start) as usize, + ); + parts.push( + child + .expected(local.start - chunk.start..local.end - chunk.start, &mask)?, + ); + } + join(&children[0].dtype(), parts) + } + Tree::Pack(children) => { + let arrays = children + .iter() + .map(|child| child.expected(rows.clone(), mask)) + .collect::>>()?; + struct_of(arrays, mask.true_count()) + } + } + } +} + +/// A random tree. Sources get distinct segments so delivery order can be scripted. +fn random_tree(rng: &mut StdRng, depth: usize, rows: Option, next_segment: &mut u32) -> Tree { + let kind = if depth == 0 { + Kind::Source + } else { + [Kind::Source, Kind::Concat, Kind::Pack, Kind::Filter][rng.random_range(0..4)] + }; + random_tree_of(rng, kind, depth, rows, next_segment) +} + +#[derive(Clone, Copy)] +enum Kind { + Source, + Concat, + Pack, + Filter, +} + +fn random_tree_of( + rng: &mut StdRng, + kind: Kind, + depth: usize, + rows: Option, + next_segment: &mut u32, +) -> Tree { + let rows = rows.unwrap_or_else(|| rng.random_range(1..24)); + match kind { + Kind::Concat => { + // Chunks summing to `rows`. Chunks of one column are leaves: a source, or a filter + // over one, so every chunk has the column's dtype. + let mut children = Vec::new(); + let mut left = rows; + while left > 0 { + let chunk = rng.random_range(1..=left); + let leaf = if rng.random_bool(0.3) { + Kind::Filter + } else { + Kind::Source + }; + children.push(random_tree_of(rng, leaf, 0, Some(chunk), next_segment)); + left -= chunk; + } + Tree::Concat(children) + } + Kind::Pack => Tree::Pack( + (0..rng.random_range(1..4)) + .map(|_| random_tree(rng, depth - 1, Some(rows), next_segment)) + .collect(), + ), + Kind::Filter => Tree::Filter(Box::new(random_tree_of( + rng, + Kind::Source, + 0, + Some(rows), + next_segment, + ))), + Kind::Source => { + let segment = rng.random_bool(0.7).then(|| { + *next_segment += 1; + *next_segment - 1 + }); + let cut = (0..rng.random_range(0..4)) + .map(|_| rng.random_range(1..8)) + .collect(); + Tree::Source { + rows, + cut, + segment, + yielding: rng.random_bool(0.3), + } + } + } +} + +/// Random trees over random views, with reads completing in random order, come out as the +/// model says. +#[test] +fn random_trees_match_their_model() -> VortexResult<()> { + let mut rng = StdRng::seed_from_u64(7); + for _ in 0..300 { + let mut next_segment = 0; + let tree = random_tree(&mut rng, 3, None, &mut next_segment); + let total = tree.rows(); + let start = rng.random_range(0..=total); + let end = rng.random_range(start..=total); + let density = rng.random_range(0.0..=1.0); + let mask = Mask::from_iter((start..end).map(|_| rng.random_bool(density))); + + let plan = tree.plan(false)?; + let mut order = StdRng::seed_from_u64(rng.random()); + let run = run(&plan, start..end, mask.clone(), |inflight| { + order.random_range(0..inflight.len()) + }) + .map_err(|e| vortex_err!("{tree:?} over {start}..{end}: {e}"))?; + + let expected = tree.expected(start..end, &mask)?; + assert_eq!(expected.dtype(), &tree.dtype()); + assert_rows(expected, run.arrays) + .map_err(|e| vortex_err!("{tree:?} over {start}..{end}: {e}"))?; + } + Ok(()) +} + #[test] fn input_take_slices_at_the_boundary() -> VortexResult<()> { let mut input = Input { diff --git a/vortex-layout/src/plan/exec/segment_scan.rs b/vortex-layout/src/plan/exec/segment_scan.rs new file mode 100644 index 00000000000..17c1ff0508d --- /dev/null +++ b/vortex-layout/src/plan/exec/segment_scan.rs @@ -0,0 +1,118 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use vortex_array::ArrayRef; +use vortex_array::buffer::BufferHandle; +use vortex_array::serde::SerializedArray; +use vortex_error::VortexResult; +use vortex_error::vortex_ensure; +use vortex_error::vortex_err; +use vortex_mask::Mask; + +use crate::plan::SegmentScanPlan; +use crate::plan::exec::ExecContext; +use crate::plan::exec::ExecNode; +use crate::plan::exec::NodeState; +use crate::plan::exec::StepCx; +use crate::plan::exec::selection::Selection; + +/// Reads one segment, then emits its rows as a single array. +/// +/// Two masks cover the node's rows: +/// +/// - The selection is always given and says which rows the parent cares about. It is a hint: rows +/// outside it may be skipped or decoded, and the node may ignore it entirely. +/// - The optional filter says which rows to return. With a filter the array holds exactly the +/// filter's rows, in order; without one it holds every row, dense, and values at rows outside +/// the selection are unspecified. A filter only selects rows the selection cares about. +/// +/// `start` publishes the read, unless the filter selects nothing or the segment is already +/// decoded, and the compute that receives the bytes decodes them, slices them to the node's +/// rows, and applies the filter. +pub(crate) struct SegmentScanNode { + plan: SegmentScanPlan, + selection: Selection, + filter: Option, + ctx: ExecContext, +} + +impl SegmentScanNode { + pub(crate) fn try_new( + plan: SegmentScanPlan, + selection: Selection, + filter: Option, + ctx: ExecContext, + ) -> VortexResult { + if let Some(filter) = &filter { + vortex_ensure!( + filter.len() == selection.mask().len(), + "SegmentScan filter of length {} does not cover rows {:?}", + filter.len(), + selection.rows() + ); + } + Ok(Self { + plan, + selection, + filter, + ctx, + }) + } + + /// Decodes the whole segment. + fn decode(&self, segment: BufferHandle) -> VortexResult { + let serialized = match self.plan.array_tree() { + Some(tree) => SerializedArray::from_flatbuffer_and_segment(tree.clone(), segment)?, + None => SerializedArray::try_from(segment)?, + }; + let row_count = usize::try_from(self.plan.row_count())?; + serialized.decode( + self.plan.dtype(), + row_count, + self.plan.array_ctx(), + self.ctx.session(), + ) + } + + /// Slices the whole decoded segment to the node's rows, filtered when the node has a filter. + fn select(&self, array: ArrayRef) -> VortexResult { + let rows = self.selection.rows(); + let mut array = if rows.start == 0 && rows.end == self.plan.row_count() { + array + } else { + array.slice(usize::try_from(rows.start)?..usize::try_from(rows.end)?)? + }; + if let Some(filter) = &self.filter + && !filter.all_true() + { + array = array.filter(filter.clone())?; + } + Ok(array) + } +} + +impl ExecNode for SegmentScanNode { + fn start(&mut self, cx: &mut StepCx<'_>) -> VortexResult { + if self.filter.as_ref().is_some_and(Mask::all_false) { + return Ok(NodeState::Done); + } + if let Some(array) = self.ctx.decoded().get(self.plan.segment_id()) { + cx.emit(self.select(array)?); + return Ok(NodeState::Done); + } + cx.request(self.plan.segment_id()); + Ok(NodeState::Wait) + } + + fn compute(&mut self, cx: &mut StepCx<'_>) -> VortexResult { + let segment = cx + .take_delivery() + .ok_or_else(|| vortex_err!("SegmentScan ran without its delivery"))?; + let array = self.decode(segment)?; + self.ctx + .decoded() + .insert(self.plan.segment_id(), array.clone()); + cx.emit(self.select(array)?); + Ok(NodeState::Done) + } +} diff --git a/vortex-layout/src/plan/exec/selection.rs b/vortex-layout/src/plan/exec/selection.rs new file mode 100644 index 00000000000..ce6f1c34033 --- /dev/null +++ b/vortex-layout/src/plan/exec/selection.rs @@ -0,0 +1,55 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use std::ops::Range; + +use vortex_array::ArrayRef; +use vortex_array::Canonical; +use vortex_array::IntoArray; +use vortex_array::arrays::ChunkedArray; +use vortex_array::dtype::DType; +use vortex_error::VortexResult; +use vortex_error::vortex_ensure; +use vortex_mask::Mask; + +/// A node's selection over `rows` of its plan domain, indexed by plan row. +#[derive(Clone, Debug)] +pub(crate) struct Selection { + rows: Range, + mask: Mask, +} + +impl Selection { + pub(crate) fn try_new(rows: Range, mask: Mask) -> VortexResult { + vortex_ensure!( + rows.start <= rows.end && mask.len() as u64 == rows.end - rows.start, + "Mask of length {} does not cover rows {rows:?}", + mask.len() + ); + Ok(Self { rows, mask }) + } + + pub(crate) fn rows(&self) -> &Range { + &self.rows + } + + pub(crate) fn mask(&self) -> &Mask { + &self.mask + } + + /// The mask over `rows`, which must lie within this selection. + pub(crate) fn slice(&self, rows: &Range) -> VortexResult { + let start = usize::try_from(rows.start - self.rows.start)?; + let end = usize::try_from(rows.end - self.rows.start)?; + Ok(self.mask.slice(start..end)) + } +} + +/// Joins arrays covering consecutive rows into one. An empty list joins to an empty array. +pub(crate) fn join(dtype: &DType, mut arrays: Vec) -> VortexResult { + match arrays.len() { + 0 => Ok(Canonical::empty(dtype).into_array()), + 1 => Ok(arrays.remove(0)), + _ => Ok(ChunkedArray::try_new(arrays, dtype.clone())?.into_array()), + } +} diff --git a/vortex-layout/src/plan/exec/synthetic.rs b/vortex-layout/src/plan/exec/synthetic.rs index 1d6b74e44ed..a32b4350149 100644 --- a/vortex-layout/src/plan/exec/synthetic.rs +++ b/vortex-layout/src/plan/exec/synthetic.rs @@ -128,7 +128,7 @@ pub fn row_dtype() -> DType { /// /// A selected source emits only the rows its mask selects, as every operator does. A dense /// source ignores its mask and emits every row of its range, as a bare segment scan does, so -/// a filter plan can sit over it. The indices are the plan's own rows, so +/// a [`Filter`](crate::plan::Filter) can sit over it. The indices are the plan's own rows, so /// a parent that reorders or filters can be checked against a sequence. Piece lengths count /// emitted rows; the last piece takes whatever remains. pub struct RowSource { diff --git a/vortex-layout/src/plan/exec/tests.rs b/vortex-layout/src/plan/exec/tests.rs new file mode 100644 index 00000000000..390a6e9734c --- /dev/null +++ b/vortex-layout/src/plan/exec/tests.rs @@ -0,0 +1,548 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +// Fixtures use a 20-row domain and name columns after their struct fields. +#![allow(clippy::cast_possible_truncation, clippy::many_single_char_names)] + +use std::ops::Range; + +use rstest::rstest; +use vortex_array::ArrayContext; +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::VortexSessionExecute; +use vortex_array::arrays::BoolArray; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::arrays::StructArray; +use vortex_array::arrays::VarBinViewArray; +use vortex_array::assert_arrays_eq; +use vortex_array::buffer::BufferHandle; +use vortex_array::serde::SerializeOptions; +use vortex_buffer::Alignment; +use vortex_buffer::ByteBufferMut; +use vortex_error::VortexResult; +use vortex_error::vortex_err; +use vortex_mask::Mask; +use vortex_session::registry::ReadContext; + +use super::*; +use crate::LayoutRef; +use crate::OwnedLayoutChildren; +use crate::layouts::chunked::ChunkedLayout; +use crate::layouts::flat::FlatLayout; +use crate::layouts::struct_::StructLayout; +use crate::plan::Filter; +use crate::plan::SegmentScan; +use crate::plan::exec::selection::join; +use crate::plan::lower; +use crate::test::SESSION; + +const ROWS: u64 = 20; + +/// In-memory segments, standing in for the IO service. +#[derive(Default)] +struct Store { + segments: Vec, +} + +impl Store { + fn flat(&mut self, array: &ArrayRef) -> VortexResult { + let ctx = ArrayContext::empty(); + let buffers = array.serialize( + &ctx, + &SESSION, + &SerializeOptions { + offset: 0, + include_padding: true, + }, + )?; + let mut bytes = ByteBufferMut::empty_aligned(Alignment::new(64)); + for buffer in buffers { + bytes.extend_from_slice(buffer.as_ref()); + } + let segment_id = SegmentId::from(u32::try_from(self.segments.len())?); + self.segments.push(BufferHandle::new_host(bytes.freeze())); + Ok(FlatLayout::new( + array.len() as u64, + array.dtype().clone(), + segment_id, + ReadContext::new(ctx.to_ids()), + ) + .into_layout()) + } + + fn chunked(&mut self, array: &ArrayRef, sizes: &[usize]) -> VortexResult { + let mut chunks = Vec::with_capacity(sizes.len()); + let mut start = 0; + for size in sizes { + chunks.push(self.flat(&array.slice(start..start + size)?)?); + start += size; + } + assert_eq!(start, array.len(), "chunk sizes must cover the array"); + Ok(ChunkedLayout::new( + array.len() as u64, + array.dtype().clone(), + OwnedLayoutChildren::layout_children(chunks), + ) + .into_layout()) + } + + fn read(&self, segment_id: SegmentId) -> BufferHandle { + self.segments[*segment_id as usize].clone() + } +} + +/// A struct whose columns are chunked at different row boundaries. +/// +/// ```text +/// a i32 chunks [0,7) [7,12) [12,20) +/// b i64 chunks [0,10) [10,20) +/// c utf8 flat [0,20) +/// d i32? chunks [0,1) [1,3) [3,6) [6,20) +/// e struct x: chunks [0,5) [5,20), y: flat +/// ``` +fn fixture(store: &mut Store) -> VortexResult<(PlanRef, ArrayRef)> { + let a = PrimitiveArray::from_iter(0..ROWS as i32).into_array(); + let b = PrimitiveArray::from_iter((0..ROWS as i64).map(|v| v * 10)).into_array(); + let c = VarBinViewArray::from_iter_str((0..ROWS).map(|v| format!("row-{v}"))).into_array(); + let d = PrimitiveArray::from_option_iter((0..ROWS as i32).map(|v| (v % 3 != 0).then_some(v))) + .into_array(); + let x = PrimitiveArray::from_iter((0..ROWS).map(|v| v as u8)).into_array(); + let y = BoolArray::from_iter((0..ROWS).map(|v| v % 2 == 0)).into_array(); + let e = StructArray::from_fields(&[("x", x.clone()), ("y", y.clone())])?.into_array(); + let expected = StructArray::from_fields(&[ + ("a", a.clone()), + ("b", b.clone()), + ("c", c.clone()), + ("d", d.clone()), + ("e", e.clone()), + ])? + .into_array(); + + let e_layout = StructLayout::new( + ROWS, + e.dtype().clone(), + vec![store.chunked(&x, &[5, 15])?, store.flat(&y)?], + ) + .into_layout(); + let layout = StructLayout::new( + ROWS, + expected.dtype().clone(), + vec![ + store.chunked(&a, &[7, 5, 8])?, + store.chunked(&b, &[10, 10])?, + store.flat(&c)?, + store.chunked(&d, &[1, 2, 3, 14])?, + e_layout, + ], + ) + .into_layout(); + Ok((lower(&layout)?, expected)) +} + +/// Which outstanding request the scripted IO service completes next. +#[derive(Clone, Copy, Debug)] +enum Delivery { + /// Oldest request first. + Fifo, + /// Newest request first. + Lifo, +} + +#[derive(Debug, PartialEq, Eq)] +enum Event { + /// A batch returned by `compute`, as segment ids. + Io(Vec), + /// A result delivered for this segment. + Delivered(u32), + /// An array of this many rows reached the root. + Piece(usize), +} + +struct Run { + arrays: Vec, + events: Vec, +} + +/// Drives a graph the way an owner does, completing reads only while the graph waits. +fn run( + store: &Store, + plan: &PlanRef, + rows: Range, + mask: Mask, + pick: impl FnMut(&[IoRequest]) -> usize, +) -> VortexResult { + run_with(store, plan, rows, mask, pick, DecodeCache::default()) +} + +/// Like [`run`], sharing `decoded` with other graphs. +fn run_with( + store: &Store, + plan: &PlanRef, + rows: Range, + mask: Mask, + mut pick: impl FnMut(&[IoRequest]) -> usize, + decoded: DecodeCache, +) -> VortexResult { + let mut graph = ExecGraph::try_new(SESSION.clone(), plan, rows, mask, 0, decoded)?; + let mut inflight: Vec = Vec::new(); + let mut arrays = Vec::new(); + let mut events = Vec::new(); + loop { + match graph.state() { + ExecState::Done => break, + ExecState::NeedsCompute => match graph.compute()? { + ExecOutput::Piece(array) => { + assert!(!array.is_empty(), "the root never emits an empty array"); + events.push(Event::Piece(array.len())); + arrays.push(array); + } + ExecOutput::NeedsIO(batch) => { + events.push(Event::Io(batch.iter().map(|r| *r.segment_id).collect())); + inflight.extend(batch); + } + ExecOutput::Yield => {} + }, + ExecState::Waiting => { + if inflight.is_empty() { + return Err(vortex_err!("graph waits with no reads in flight")); + } + let request = inflight.remove(pick(&inflight)); + events.push(Event::Delivered(*request.segment_id)); + graph.set_io_result(request.id, store.read(request.segment_id))?; + } + } + } + assert!(inflight.is_empty(), "graph finished with reads in flight"); + Ok(Run { arrays, events }) +} + +fn delivery(order: Delivery) -> impl FnMut(&[IoRequest]) -> usize { + move |inflight| match order { + Delivery::Fifo => 0, + Delivery::Lifo => inflight.len() - 1, + } +} + +/// Delivers segments in the given order. +fn scripted(order: &[u32]) -> impl FnMut(&[IoRequest]) -> usize + '_ { + let mut next = order.iter(); + move |inflight| { + let segment = *next.next().expect("script covers every read"); + inflight + .iter() + .position(|r| *r.segment_id == segment) + .expect("scripted segment is in flight") + } +} + +fn reads(events: &[Event]) -> usize { + events + .iter() + .map(|event| match event { + Event::Io(batch) => batch.len(), + _ => 0, + }) + .sum() +} + +fn pieces(events: &[Event]) -> Vec { + events + .iter() + .filter_map(|event| match event { + Event::Piece(len) => Some(*len), + _ => None, + }) + .collect() +} + +/// A view's selection over its rows. +#[derive(Clone, Copy, Debug)] +enum Sel { + All, + None, + EveryOther, + /// Rows relative to the view start. + Rows(&'static [usize]), +} + +impl Sel { + fn mask(self, len: usize) -> Mask { + match self { + Sel::All => Mask::new_true(len), + Sel::None => Mask::new_false(len), + Sel::EveryOther => Mask::from_iter((0..len).map(|i| i % 2 == 0)), + Sel::Rows(rows) => Mask::from_indices(len, rows.iter().copied()), + } + } +} + +/// Checks that the arrays, in the order they came out, are the selected rows of the view. +fn assert_view( + expected: &ArrayRef, + rows: &Range, + mask: &Mask, + arrays: Vec, +) -> VortexResult<()> { + let actual = join(expected.dtype(), arrays)?; + let expected = expected + .slice(rows.start as usize..rows.end as usize)? + .filter(mask.clone())?; + let mut ctx = SESSION.create_execution_ctx(); + assert_arrays_eq!(actual, expected, &mut ctx); + Ok(()) +} + +#[rstest] +#[case::full(0..20, Sel::All)] +#[case::crosses_every_boundary(3..17, Sel::All)] +#[case::every_other(5..15, Sel::EveryOther)] +#[case::sparse(0..20, Sel::Rows(&[0, 6, 7, 19]))] +#[case::one_row(8..9, Sel::All)] +#[case::inside_one_chunk_per_column(13..16, Sel::All)] +#[case::nothing_selected(2..18, Sel::None)] +#[case::empty_range(5..5, Sel::All)] +fn views_of_one_plan( + #[case] rows: Range, + #[case] sel: Sel, + #[values(Delivery::Fifo, Delivery::Lifo)] order: Delivery, +) -> VortexResult<()> { + let mut store = Store::default(); + let (plan, expected) = fixture(&mut store)?; + let mask = sel.mask((rows.end - rows.start) as usize); + + let run = run(&store, &plan, rows.clone(), mask.clone(), delivery(order))?; + assert_view(&expected, &rows, &mask, run.arrays) +} + +#[test] +fn many_views_share_one_plan() -> VortexResult<()> { + let mut store = Store::default(); + let (plan, expected) = fixture(&mut store)?; + + // Split the row domain into views that cut through every column at different places. + for split in [[0, 4, 11, 20], [0, 9, 13, 20], [0, 1, 2, 20]] { + for window in split.windows(2) { + let rows = window[0]..window[1]; + let mask = Sel::EveryOther.mask((rows.end - rows.start) as usize); + let run = run( + &store, + &plan, + rows.clone(), + mask.clone(), + delivery(Delivery::Lifo), + )?; + assert_view(&expected, &rows, &mask, run.arrays)?; + } + } + Ok(()) +} + +#[test] +fn first_compute_publishes_every_read_in_one_batch() -> VortexResult<()> { + let mut store = Store::default(); + let (plan, _) = fixture(&mut store)?; + + let run = run( + &store, + &plan, + 0..ROWS, + Mask::new_true(ROWS as usize), + delivery(Delivery::Fifo), + )?; + let Some(Event::Io(first)) = run.events.first() else { + return Err(vortex_err!("first event must be a read batch")); + }; + assert_eq!(first.len(), store.segments.len()); + assert_eq!(reads(&run.events), store.segments.len()); + Ok(()) +} + +#[test] +fn unselected_chunks_are_never_read() -> VortexResult<()> { + let mut store = Store::default(); + let (plan, expected) = fixture(&mut store)?; + // Rows 7..12 only: a reads one chunk, b both, c one, d one, e.x one, e.y one. + let mask = Mask::from_indices(ROWS as usize, 7..12); + + let run = run( + &store, + &plan, + 0..ROWS, + mask.clone(), + delivery(Delivery::Fifo), + )?; + assert_eq!(reads(&run.events), 7); + assert_view(&expected, &(0..ROWS), &mask, run.arrays) +} + +#[test] +fn nothing_selected_issues_no_io_and_emits_nothing() -> VortexResult<()> { + let mut store = Store::default(); + let (plan, _) = fixture(&mut store)?; + let mask = Mask::new_false(ROWS as usize); + + let run = run(&store, &plan, 0..ROWS, mask, delivery(Delivery::Fifo))?; + assert_eq!(reads(&run.events), 0); + assert!(run.arrays.is_empty()); + Ok(()) +} + +/// `{a, b}` with `a` chunked `[0,7) [7,12) [12,20)` and `b` flat. +fn two_columns(store: &mut Store) -> VortexResult<(PlanRef, ArrayRef)> { + let a = PrimitiveArray::from_iter(0..ROWS as i32).into_array(); + let b = PrimitiveArray::from_iter((0..ROWS as i64).map(|v| v * 10)).into_array(); + let expected = StructArray::from_fields(&[("a", a.clone()), ("b", b.clone())])?.into_array(); + let layout = StructLayout::new( + ROWS, + expected.dtype().clone(), + vec![store.chunked(&a, &[7, 5, 8])?, store.flat(&b)?], + ) + .into_layout(); + Ok((lower(&layout)?, expected)) +} + +/// Pack emits a struct as soon as every field has rows, however the reads complete, and the +/// structs come out in row order. +/// +/// Segments: a0=0, a1=1, a2=2, b=3. +#[rstest] +// `b` arrives last: nothing can be emitted before it, then everything at once. +#[case::b_last(&[1, 0, 2, 3], &[20])] +// `b` and `a0` first: rows 0..7 go out; `a2` arrives before `a1` and waits in its port. +#[case::a_chunk_last(&[3, 0, 2, 1], &[7, 13])] +// Chunks in order: one struct per chunk of `a`. +#[case::in_order(&[3, 0, 1, 2], &[7, 5, 8])] +fn pack_streams_in_row_order( + #[case] order: &[u32], + #[case] expected_pieces: &[usize], +) -> VortexResult<()> { + let mut store = Store::default(); + let (plan, expected) = two_columns(&mut store)?; + let mask = Mask::new_true(ROWS as usize); + + let run = run(&store, &plan, 0..ROWS, mask.clone(), scripted(order))?; + assert_eq!(pieces(&run.events), expected_pieces); + assert_view(&expected, &(0..ROWS), &mask, run.arrays) +} + +/// Fields chunked alike stream one struct per chunk. +#[test] +fn aligned_fields_stream_one_struct_per_chunk() -> VortexResult<()> { + let mut store = Store::default(); + let a = PrimitiveArray::from_iter(0..ROWS as i32).into_array(); + let b = PrimitiveArray::from_iter((0..ROWS as i64).map(|v| v * 10)).into_array(); + let expected = StructArray::from_fields(&[("a", a.clone()), ("b", b.clone())])?.into_array(); + let layout = StructLayout::new( + ROWS, + expected.dtype().clone(), + vec![ + store.chunked(&a, &[7, 5, 8])?, + store.chunked(&b, &[7, 5, 8])?, + ], + ) + .into_layout(); + let plan = lower(&layout)?; + let mask = Mask::new_true(ROWS as usize); + + let run = run( + &store, + &plan, + 0..ROWS, + mask.clone(), + delivery(Delivery::Fifo), + )?; + assert_eq!(pieces(&run.events), [7, 5, 8]); + assert_view(&expected, &(0..ROWS), &mask, run.arrays) +} + +#[test] +fn state_is_side_effect_free() -> VortexResult<()> { + let mut store = Store::default(); + let (plan, _) = two_columns(&mut store)?; + let graph = ExecGraph::try_new( + SESSION.clone(), + &plan, + 0..ROWS, + Mask::new_true(ROWS as usize), + 0, + DecodeCache::default(), + )?; + for _ in 0..3 { + assert_eq!(graph.state(), ExecState::NeedsCompute); + } + Ok(()) +} + +/// A bare segment scan returns every row of its range whatever it is told to care about, and a +/// filter over the same scan returns only the selected rows. +#[rstest] +#[case::every_other(Sel::EveryOther)] +#[case::sparse(Sel::Rows(&[1, 4]))] +#[case::nothing(Sel::None)] +fn bare_scan_is_dense_and_filter_keeps_the_selection(#[case] sel: Sel) -> VortexResult<()> { + let mut store = Store::default(); + let values = PrimitiveArray::from_iter(0..ROWS as i32).into_array(); + let filtered = lower(&store.flat(&values)?)?; + assert!(filtered.is::()); + let scan = filtered.child_required(0)?; + assert!(scan.is::()); + + let rows = 3..9; + let mask = sel.mask(6); + let mut ctx = SESSION.create_execution_ctx(); + + let dense = run( + &store, + &scan, + rows.clone(), + mask.clone(), + delivery(Delivery::Fifo), + )?; + assert_eq!(dense.arrays.len(), 1); + assert_arrays_eq!(dense.arrays[0], values.slice(3..9)?, &mut ctx); + + let kept = run( + &store, + &filtered, + rows.clone(), + mask.clone(), + delivery(Delivery::Fifo), + )?; + assert_view(&values, &rows, &mask, kept.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] +fn shared_decode_cache_skips_reads() -> VortexResult<()> { + let mut store = Store::default(); + let (plan, expected) = fixture(&mut store)?; + let decoded = DecodeCache::default(); + + let rows = 0..ROWS; + let all = Mask::new_true(ROWS as usize); + let first = run_with( + &store, + &plan, + rows.clone(), + all, + delivery(Delivery::Fifo), + decoded.clone(), + )?; + assert!(reads(&first.events) > 0); + + let mask = Sel::EveryOther.mask(ROWS as usize); + let second = run_with( + &store, + &plan, + rows.clone(), + mask.clone(), + delivery(Delivery::Fifo), + decoded, + )?; + assert_eq!(reads(&second.events), 0); + assert_view(&expected, &rows, &mask, second.arrays)?; + Ok(()) +} diff --git a/vortex-layout/src/plan/plans/concat.rs b/vortex-layout/src/plan/plans/concat.rs index 2e814156309..f8b704c2131 100644 --- a/vortex-layout/src/plan/plans/concat.rs +++ b/vortex-layout/src/plan/plans/concat.rs @@ -2,12 +2,14 @@ // SPDX-FileCopyrightText: Copyright the Vortex contributors use std::borrow::Cow; +use std::ops::Range; use std::sync::Arc; use vortex_array::EmptyMetadata; use vortex_array::dtype::DType; use vortex_error::VortexResult; use vortex_error::vortex_bail; +use vortex_mask::Mask; use vortex_session::registry::CachedId; use crate::plan::Eval; @@ -18,6 +20,10 @@ use crate::plan::PlanId; use crate::plan::PlanParts; use crate::plan::PlanRef; use crate::plan::PlanVTable; +use crate::plan::exec::ConcatNode; +use crate::plan::exec::ExecContext; +use crate::plan::exec::ExecNode; +use crate::plan::exec::Selection; use crate::plan::optimizer::PlanParentReduceRule; /// Concatenates its children row-wise. @@ -142,6 +148,18 @@ impl PlanVTable for Concat { fn child_name(_plan: &Plan, index: usize) -> Cow<'_, str> { Cow::Owned(format!("chunks[{index}]")) } + + fn exec( + plan: &Plan, + rows: Range, + mask: Mask, + _ctx: &ExecContext, + ) -> VortexResult> { + Ok(Box::new(ConcatNode::new( + plan.clone(), + Selection::try_new(rows, mask)?, + ))) + } } /// Pushes an expression into every chunk of a [`Concat`]. diff --git a/vortex-layout/src/plan/plans/filter.rs b/vortex-layout/src/plan/plans/filter.rs index da5907f3e80..e7424b236ba 100644 --- a/vortex-layout/src/plan/plans/filter.rs +++ b/vortex-layout/src/plan/plans/filter.rs @@ -2,11 +2,13 @@ // SPDX-FileCopyrightText: Copyright the Vortex contributors use std::borrow::Cow; +use std::ops::Range; use vortex_array::EmptyMetadata; use vortex_error::VortexResult; use vortex_error::vortex_bail; use vortex_error::vortex_err; +use vortex_mask::Mask; use vortex_session::registry::CachedId; use crate::plan::Plan; @@ -15,7 +17,13 @@ use crate::plan::PlanId; use crate::plan::PlanParts; use crate::plan::PlanRef; use crate::plan::PlanVTable; +use crate::plan::SegmentScan; use crate::plan::check_child_count; +use crate::plan::exec::ExecContext; +use crate::plan::exec::ExecNode; +use crate::plan::exec::FilterNode; +use crate::plan::exec::SegmentScanNode; +use crate::plan::exec::Selection; /// Keeps only the selected rows of its child. /// @@ -86,4 +94,29 @@ impl PlanVTable for Filter { Cow::Owned(format!("child[{index}]")) } } + + /// Runs fused with a segment-scan child, as one node that keeps the selected rows itself; + /// over any other child, as a filter node that filters the child's whole pieces. + fn exec( + plan: &Plan, + rows: Range, + mask: Mask, + ctx: &ExecContext, + ) -> VortexResult> { + let child = plan.child_plan()?; + let Some(scan) = child.as_opt::() else { + return Ok(Box::new(FilterNode::new( + plan.clone(), + Selection::try_new(rows, mask)?, + ctx.session().clone(), + ))); + }; + let filter = Some(mask.clone()); + Ok(Box::new(SegmentScanNode::try_new( + scan.clone(), + Selection::try_new(rows, mask)?, + filter, + ctx.clone(), + )?)) + } } diff --git a/vortex-layout/src/plan/plans/pack.rs b/vortex-layout/src/plan/plans/pack.rs index 6497cc63d45..81311678376 100644 --- a/vortex-layout/src/plan/plans/pack.rs +++ b/vortex-layout/src/plan/plans/pack.rs @@ -2,6 +2,7 @@ // SPDX-FileCopyrightText: Copyright the Vortex contributors use std::borrow::Cow; +use std::ops::Range; use vortex_array::EmptyMetadata; use vortex_array::dtype::DType; @@ -27,6 +28,7 @@ use vortex_error::VortexResult; use vortex_error::vortex_bail; use vortex_error::vortex_ensure; use vortex_error::vortex_err; +use vortex_mask::Mask; use vortex_session::registry::CachedId; use crate::plan::Eval; @@ -37,6 +39,10 @@ use crate::plan::PlanId; use crate::plan::PlanParts; use crate::plan::PlanRef; use crate::plan::PlanVTable; +use crate::plan::exec::ExecContext; +use crate::plan::exec::ExecNode; +use crate::plan::exec::PackNode; +use crate::plan::exec::Selection; use crate::plan::optimizer::PlanParentReduceRule; /// Assembles a struct from one child per field, plus an optional trailing validity child. @@ -204,6 +210,18 @@ impl PlanVTable for Pack { } Cow::Borrowed("validity") } + + fn exec( + plan: &Plan, + rows: Range, + mask: Mask, + _ctx: &ExecContext, + ) -> VortexResult> { + Ok(Box::new(PackNode::new( + plan.clone(), + Selection::try_new(rows, mask)?, + ))) + } } fn validate_field_child( diff --git a/vortex-layout/src/plan/plans/segment_scan.rs b/vortex-layout/src/plan/plans/segment_scan.rs index d2b20df89ad..521d44b589f 100644 --- a/vortex-layout/src/plan/plans/segment_scan.rs +++ b/vortex-layout/src/plan/plans/segment_scan.rs @@ -1,10 +1,13 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors +use std::ops::Range; + use vortex_array::EmptyMetadata; use vortex_array::dtype::DType; use vortex_buffer::ByteBuffer; use vortex_error::VortexResult; +use vortex_mask::Mask; use vortex_session::registry::CachedId; use vortex_session::registry::ReadContext; @@ -14,6 +17,10 @@ use crate::plan::PlanId; use crate::plan::PlanParts; use crate::plan::PlanVTable; use crate::plan::check_child_count; +use crate::plan::exec::ExecContext; +use crate::plan::exec::ExecNode; +use crate::plan::exec::SegmentScanNode; +use crate::plan::exec::Selection; use crate::segments::SegmentId; /// Reads one serialized array segment. @@ -92,4 +99,19 @@ impl PlanVTable for SegmentScan { check_child_count("SegmentScan", children, 0)?; Ok(()) } + + fn exec( + plan: &Plan, + rows: Range, + mask: Mask, + ctx: &ExecContext, + ) -> VortexResult> { + // A bare scan returns every row; a Filter over it runs the same node with a filter. + Ok(Box::new(SegmentScanNode::try_new( + plan.clone(), + Selection::try_new(rows, mask)?, + None, + ctx.clone(), + )?)) + } } From d092cfd1f885e68efd0a673e934a107ad4de7ed2 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 9 Oct 2026 19:16:15 +0000 Subject: [PATCH 3/4] refactor(layout): keep Filter out of lower, and let a segment scan keep its own rows A flat layout lowers to a bare segment scan again, as on develop. Every plan now produces only the rows its split selects, a segment scan included, so the scan's source applies the mask; a Filter over a scan compiles to the same source. Signed-off-by: Claude Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01NWYmfHTFtaEdhdyd3u5hv3 --- vortex-layout/src/plan/lower.rs | 6 +- vortex-layout/src/plan/pipeline/tests.rs | 45 ++--- vortex-layout/src/plan/plans/segment_scan.rs | 5 +- vortex-layout/src/plan/tests.rs | 190 +++++++------------ 4 files changed, 89 insertions(+), 157 deletions(-) diff --git a/vortex-layout/src/plan/lower.rs b/vortex-layout/src/plan/lower.rs index d04e3e128d4..e8911734549 100644 --- a/vortex-layout/src/plan/lower.rs +++ b/vortex-layout/src/plan/lower.rs @@ -25,7 +25,6 @@ use crate::layouts::list::VALIDITY_CHILD_INDEX; use crate::layouts::struct_::Struct; use crate::layouts::struct_::StructLayout; use crate::plan::ConcatPlan; -use crate::plan::FilterPlan; use crate::plan::ListPackPlan; use crate::plan::PackPlan; use crate::plan::PlanChildren; @@ -37,12 +36,9 @@ use crate::plan::TakePlan; /// /// The root operator is built immediately. Its child container owns a hidden clone of the source /// layout and lowers each child independently on first access. -/// -/// A flat layout lowers to a [`Filter`](crate::plan::Filter) over its segment scan, so the scan -/// returns only the rows it is executed with. pub fn lower(layout: &LayoutRef) -> VortexResult { if let Some(layout) = layout.as_opt::() { - return Ok(FilterPlan::new(lower_flat(layout).into_plan()).into_plan()); + return Ok(lower_flat(layout).into_plan()); } if let Some(layout) = layout.as_opt::() { return Ok(lower_chunked(layout)?.into_plan()); diff --git a/vortex-layout/src/plan/pipeline/tests.rs b/vortex-layout/src/plan/pipeline/tests.rs index 3aa732b7e45..41d44dcc3a1 100644 --- a/vortex-layout/src/plan/pipeline/tests.rs +++ b/vortex-layout/src/plan/pipeline/tests.rs @@ -33,7 +33,6 @@ use crate::OwnedLayoutChildren; use crate::layouts::chunked::ChunkedLayout; use crate::layouts::flat::FlatLayout; use crate::layouts::struct_::StructLayout; -use crate::plan::Filter; use crate::plan::PlanRef; use crate::plan::SegmentScan; use crate::plan::lower; @@ -521,42 +520,34 @@ fn aligned_fields_stream_one_struct_per_chunk() -> VortexResult<()> { assert_view(&expected, &(0..ROWS), &mask, run.arrays) } -/// A bare segment scan returns every row of its range whatever it is told to care about, and a -/// filter over the same scan returns only the selected rows. +/// A segment scan keeps the rows its split selects, and so does a filter over the same scan, +/// which runs as the same source. #[rstest] #[case::every_other(Sel::EveryOther)] #[case::sparse(Sel::Rows(&[1, 4]))] #[case::nothing(Sel::None)] -fn bare_scan_is_dense_and_filter_keeps_the_selection(#[case] sel: Sel) -> VortexResult<()> { +fn scan_keeps_the_selection_with_or_without_a_filter(#[case] sel: Sel) -> VortexResult<()> { let mut store = Store::default(); let values = PrimitiveArray::from_iter(0..ROWS as i32).into_array(); - let filtered = lower(&store.flat(&values)?)?; - assert!(filtered.is::()); - let scan = filtered.child_required(0)?; + let scan = lower(&store.flat(&values)?)?; assert!(scan.is::()); + let filtered = crate::plan::FilterPlan::new(scan.clone()).into_plan(); let rows = 3..9; let mask = sel.mask(6); - let mut ctx = SESSION.create_execution_ctx(); - - let dense = run( - &store, - &scan, - rows.clone(), - mask.clone(), - delivery(Delivery::Fifo), - )?; - assert_eq!(dense.arrays.len(), 1); - assert_arrays_eq!(dense.arrays[0], values.slice(3..9)?, &mut ctx); - - let kept = run( - &store, - &filtered, - rows.clone(), - mask.clone(), - delivery(Delivery::Fifo), - )?; - assert_view(&values, &rows, &mask, kept.arrays)?; + for plan in [&scan, &filtered] { + let kept = run( + &store, + plan, + rows.clone(), + mask.clone(), + delivery(Delivery::Fifo), + )?; + if mask.all_false() { + assert_eq!(reads(&kept.events), 0, "nothing selected reads nothing"); + } + assert_view(&values, &rows, &mask, kept.arrays)?; + } Ok(()) } diff --git a/vortex-layout/src/plan/plans/segment_scan.rs b/vortex-layout/src/plan/plans/segment_scan.rs index 5f455fae561..9ace7e43ff9 100644 --- a/vortex-layout/src/plan/plans/segment_scan.rs +++ b/vortex-layout/src/plan/plans/segment_scan.rs @@ -105,9 +105,8 @@ impl PlanVTable for SegmentScan { mask: &Mask, compiler: &mut Compiler<'_>, ) -> VortexResult> { - let _ = mask; - // A bare scan produces every row of its range; a filter over it keeps the selected ones. - compiler.scan(plan, rows, None) + // A scan keeps the selected rows itself, as every plan produces only those. + compiler.scan(plan, rows, Some(mask.clone())) } fn reach( diff --git a/vortex-layout/src/plan/tests.rs b/vortex-layout/src/plan/tests.rs index fccc06ebb2f..783eb455327 100644 --- a/vortex-layout/src/plan/tests.rs +++ b/vortex-layout/src/plan/tests.rs @@ -94,14 +94,12 @@ fn unsupported_layout_has_no_plan() -> VortexResult<()> { } #[test] -fn flat_plan_is_a_filtered_segment_scan_without_children() -> VortexResult<()> { +fn flat_plan_has_no_children() -> VortexResult<()> { let plan = make_plan(flat(3, primitive(PType::I32, Nullability::NonNullable), 0))?; - assert!(plan.is::()); - let scan = child_of(&plan, 0)?; - assert!(scan.is::()); - assert_eq!(scan.child_count(), 0); - assert!(scan.child(0)?.is_none()); + assert!(plan.is::()); + assert_eq!(plan.child_count(), 0); + assert!(plan.child(0)?.is_none()); Ok(()) } @@ -325,7 +323,7 @@ fn optimize_drops_identity_expressions() -> VortexResult<()> { let expression = root().bind(child.dtype())?; let plan: PlanRef = EvalPlan::try_new(expression, child)?.into_plan(); - assert!(optimize(plan)?.is::()); + assert!(optimize(plan)?.is::()); Ok(()) } @@ -348,8 +346,8 @@ fn optimize_rewrites_nested_children() -> VortexResult<()> { let optimized = optimize(wrapped)?; assert!(optimized.is::()); - assert!(child_of(&optimized, 0)?.is::()); - assert!(child_of(&optimized, 1)?.is::()); + assert!(child_of(&optimized, 0)?.is::()); + assert!(child_of(&optimized, 1)?.is::()); Ok(()) } @@ -373,13 +371,11 @@ fn plan_display_matches_array_tree_display_shape() -> VortexResult<()> { let plan: PlanRef = plan.into_plan(); assert_eq!(plan.to_string(), "vortex.plan.eval(i32, rows=3) expr=$.a"); - insta::assert_snapshot!(plan.display_tree(), @" + insta::assert_snapshot!(plan.display_tree(), @r" root: vortex.plan.eval(i32, rows=3) expr=$.a child: vortex.plan.pack({a=i32, b=i32}, rows=3) - a: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) - b: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) + a: vortex.plan.segment_scan(i32, rows=3) + b: vortex.plan.segment_scan(i32, rows=3) "); struct DepthExtractor; @@ -395,13 +391,11 @@ fn plan_display_matches_array_tree_display_shape() -> VortexResult<()> { } } - insta::assert_snapshot!(plan.tree_display_builder().with(DepthExtractor), @" + insta::assert_snapshot!(plan.tree_display_builder().with(DepthExtractor), @r" root: depth=0 child: depth=1 a: depth=2 - child: depth=3 b: depth=2 - child: depth=3 "); let nullable_fields = StructFields::from_iter([ @@ -419,14 +413,11 @@ fn plan_display_matches_array_tree_display_shape() -> VortexResult<()> { ) .into_layout(); let nullable = make_plan(nullable_layout)?; - insta::assert_snapshot!(nullable.tree_display_builder(), @" + insta::assert_snapshot!(nullable.tree_display_builder(), @r" root: a: - child: b: - child: validity: - child: "); Ok(()) } @@ -442,12 +433,10 @@ fn chunked_plan_display_names_chunks() -> VortexResult<()> { .into_layout(); let plan = make_plan(layout)?; - insta::assert_snapshot!(plan.display_tree(), @" + insta::assert_snapshot!(plan.display_tree(), @r" root: vortex.plan.concat(i32, rows=3) - chunks[0]: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) - chunks[1]: vortex.plan.filter(i32, rows=1) - child: vortex.plan.segment_scan(i32, rows=1) + chunks[0]: vortex.plan.segment_scan(i32, rows=2) + chunks[1]: vortex.plan.segment_scan(i32, rows=1) "); Ok(()) } @@ -461,12 +450,10 @@ fn dict_plan_display_names_logical_children() -> VortexResult<()> { .into_layout(); let plan = make_plan(layout)?; - insta::assert_snapshot!(plan.display_tree(), @" + insta::assert_snapshot!(plan.display_tree(), @r" root: vortex.plan.take(i32, rows=3) - codes: vortex.plan.filter(u8, rows=3) - child: vortex.plan.segment_scan(u8, rows=3) - values: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) + codes: vortex.plan.segment_scan(u8, rows=3) + values: vortex.plan.segment_scan(i32, rows=2) "); Ok(()) } @@ -484,12 +471,10 @@ fn list_plan_display_handles_optional_validity() -> VortexResult<()> { .into_layout(); let non_nullable = make_plan(non_nullable_layout)?; - insta::assert_snapshot!(non_nullable.display_tree(), @" + insta::assert_snapshot!(non_nullable.display_tree(), @r" root: vortex.plan.list_pack(list(i32), rows=2) - elements: vortex.plan.filter(i32, rows=4) - child: vortex.plan.segment_scan(i32, rows=4) - offsets: vortex.plan.filter(u32, rows=3) - child: vortex.plan.segment_scan(u32, rows=3) + elements: vortex.plan.segment_scan(i32, rows=4) + offsets: vortex.plan.segment_scan(u32, rows=3) "); let nullable_layout = ListLayout::new( @@ -501,14 +486,11 @@ fn list_plan_display_handles_optional_validity() -> VortexResult<()> { .into_layout(); let nullable = make_plan(nullable_layout)?; - insta::assert_snapshot!(nullable.display_tree(), @" + insta::assert_snapshot!(nullable.display_tree(), @r" root: vortex.plan.list_pack(list(i32)?, rows=2) - elements: vortex.plan.filter(i32, rows=4) - child: vortex.plan.segment_scan(i32, rows=4) - offsets: vortex.plan.filter(u32, rows=3) - child: vortex.plan.segment_scan(u32, rows=3) - validity: vortex.plan.filter(bool, rows=2) - child: vortex.plan.segment_scan(bool, rows=2) + elements: vortex.plan.segment_scan(i32, rows=4) + offsets: vortex.plan.segment_scan(u32, rows=3) + validity: vortex.plan.segment_scan(bool, rows=2) "); Ok(()) } @@ -564,14 +546,11 @@ fn expression_partitions_across_row_idx_and_struct() -> VortexResult<()> { child: vortex.plan.eval({child_0=bool, child_1=bool}, rows=3) expr=pack(child_0: $.a, child_1: $.b) child: vortex.plan.pack({a=bool, b=bool}, rows=3) a: vortex.plan.take(bool, rows=3) - codes: vortex.plan.filter(u8, rows=3) - child: vortex.plan.segment_scan(u8, rows=3) + codes: vortex.plan.segment_scan(u8, rows=3) values: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) b: vortex.plan.eval(bool, rows=3) expr=($ > 7i32) - child: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); Ok(()) } @@ -621,18 +600,16 @@ fn row_idx_and_data_expression_pushes_data_into_chunks() -> VortexResult<()> { let expression = and(gt(row_idx(), lit(10_u64)), gt(root(), lit(5_i32))); let plan = optimize(make_row_idx_plan(expression, child)?)?; - insta::assert_snapshot!(plan.display_tree(), @" + insta::assert_snapshot!(plan.display_tree(), @r" root: vortex.plan.eval(bool, rows=3) expr=($.row_idx and $.child) child: vortex.plan.pack({row_idx=bool, child=bool}, rows=3) row_idx: vortex.plan.eval(bool, rows=3) expr=($ > 10u64) child: vortex.plan.row_idx(u64, rows=3) child: vortex.plan.concat(bool, rows=3) chunks[0]: vortex.plan.eval(bool, rows=1) expr=($ > 5i32) - child: vortex.plan.filter(i32, rows=1) - child: vortex.plan.segment_scan(i32, rows=1) + child: vortex.plan.segment_scan(i32, rows=1) chunks[1]: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) "); Ok(()) } @@ -658,22 +635,17 @@ fn expression_pushes_through_struct_field_and_dictionary_values() -> VortexResul root: vortex.plan.eval(bool, rows=3) expr=($.a > 5i32) child: vortex.plan.pack({a=i32, b=i32}, rows=3) a: vortex.plan.take(i32, rows=3) - codes: vortex.plan.filter(u8, rows=3) - child: vortex.plan.segment_scan(u8, rows=3) - values: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) - b: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) + codes: vortex.plan.segment_scan(u8, rows=3) + values: vortex.plan.segment_scan(i32, rows=2) + b: vortex.plan.segment_scan(i32, rows=3) "); let optimized = optimize(plan)?; insta::assert_snapshot!(optimized.display_tree(), @" root: vortex.plan.take(bool, rows=3) - codes: vortex.plan.filter(u8, rows=3) - child: vortex.plan.segment_scan(u8, rows=3) + codes: vortex.plan.segment_scan(u8, rows=3) values: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) "); Ok(()) } @@ -708,28 +680,21 @@ fn expression_pushes_through_struct_field_with_heterogeneous_chunks() -> VortexR child: vortex.plan.pack({a=i32, b=i32}, rows=5) a: vortex.plan.concat(i32, rows=5) chunks[0]: vortex.plan.take(i32, rows=3) - codes: vortex.plan.filter(u8, rows=3) - child: vortex.plan.segment_scan(u8, rows=3) - values: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) - chunks[1]: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) - b: vortex.plan.filter(i32, rows=5) - child: vortex.plan.segment_scan(i32, rows=5) + codes: vortex.plan.segment_scan(u8, rows=3) + values: vortex.plan.segment_scan(i32, rows=2) + chunks[1]: vortex.plan.segment_scan(i32, rows=2) + b: vortex.plan.segment_scan(i32, rows=5) "); let optimized = optimize(plan)?; insta::assert_snapshot!(optimized.display_tree(), @" root: vortex.plan.concat(bool, rows=5) chunks[0]: vortex.plan.take(bool, rows=3) - codes: vortex.plan.filter(u8, rows=3) - child: vortex.plan.segment_scan(u8, rows=3) + codes: vortex.plan.segment_scan(u8, rows=3) values: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) - chunks[1]: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.filter(i32, rows=2) child: vortex.plan.segment_scan(i32, rows=2) + chunks[1]: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) + child: vortex.plan.segment_scan(i32, rows=2) "); Ok(()) } @@ -763,10 +728,9 @@ fn expression_pushes_through_nested_struct_fields_in_one_pass() -> VortexResult< let plan = make_eval(expression, make_plan(layout)?)?.into_plan(); let optimized = optimize(plan)?; - insta::assert_snapshot!(optimized.display_tree(), @" + insta::assert_snapshot!(optimized.display_tree(), @r" root: vortex.plan.eval(bool, rows=3) expr=($ > 5i32) - child: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); Ok(()) } @@ -792,19 +756,17 @@ fn expression_pushes_through_single_field_nested_structs() -> VortexResult<()> { let expression = gt(get_item("b", get_item("a", root())), lit(5_i32)); let plan = make_eval(expression, make_plan(layout)?)?.into_plan(); - insta::assert_snapshot!(plan.display_tree(), @" + insta::assert_snapshot!(plan.display_tree(), @r" root: vortex.plan.eval(bool, rows=3) expr=($.a.b > 5i32) child: vortex.plan.pack({a={b=i32}}, rows=3) a: vortex.plan.pack({b=i32}, rows=3) - b: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) + b: vortex.plan.segment_scan(i32, rows=3) "); let optimized = optimize(plan)?; - insta::assert_snapshot!(optimized.display_tree(), @" + insta::assert_snapshot!(optimized.display_tree(), @r" root: vortex.plan.eval(bool, rows=3) expr=($ > 5i32) - child: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); Ok(()) } @@ -838,10 +800,9 @@ fn expression_pushes_through_three_single_field_structs() -> VortexResult<()> { let plan = make_eval(expression, make_plan(layout)?)?.into_plan(); let optimized = optimize(plan)?; - insta::assert_snapshot!(optimized.display_tree(), @" + insta::assert_snapshot!(optimized.display_tree(), @r" root: vortex.plan.eval(bool, rows=3) expr=($ > 5i32) - child: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); Ok(()) } @@ -886,18 +847,15 @@ fn compound_expression_pushes_through_nested_struct_and_dictionary() -> VortexRe let plan = make_eval(expression, make_plan(layout)?)?.into_plan(); let optimized = optimize(plan)?; - insta::assert_snapshot!(optimized.display_tree(), @" + insta::assert_snapshot!(optimized.display_tree(), @r" root: vortex.plan.eval(bool, rows=3) expr=($.b and $.c) child: vortex.plan.pack({b=bool, c=bool}, rows=3) b: vortex.plan.take(bool, rows=3) - codes: vortex.plan.filter(u8, rows=3) - child: vortex.plan.segment_scan(u8, rows=3) + codes: vortex.plan.segment_scan(u8, rows=3) values: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) c: vortex.plan.eval(bool, rows=3) expr=($ > 7i32) - child: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); Ok(()) } @@ -936,13 +894,11 @@ fn expression_pushes_through_dictionary_of_struct_values() -> VortexResult<()> { let plan = make_eval(expression, make_plan(layout)?)?.into_plan(); let optimized = optimize(plan)?; - insta::assert_snapshot!(optimized.display_tree(), @" + insta::assert_snapshot!(optimized.display_tree(), @r" root: vortex.plan.take(bool, rows=3) - codes: vortex.plan.filter(u8, rows=3) - child: vortex.plan.segment_scan(u8, rows=3) + codes: vortex.plan.segment_scan(u8, rows=3) values: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) "); Ok(()) } @@ -982,14 +938,10 @@ fn multi_field_struct_expression_pushes_into_each_field() -> VortexResult<()> { root: vortex.plan.eval(bool, rows=3) expr=(($.a > 5i32) and ($.b > 7i32)) child: vortex.plan.pack({a=i32, b=i32, c=i32}, rows=3) a: vortex.plan.take(i32, rows=3) - codes: vortex.plan.filter(u8, rows=3) - child: vortex.plan.segment_scan(u8, rows=3) - values: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) - b: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) - c: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) + codes: vortex.plan.segment_scan(u8, rows=3) + values: vortex.plan.segment_scan(i32, rows=2) + b: vortex.plan.segment_scan(i32, rows=3) + c: vortex.plan.segment_scan(i32, rows=3) "); let optimized = optimize(plan)?; @@ -997,14 +949,11 @@ fn multi_field_struct_expression_pushes_into_each_field() -> VortexResult<()> { root: vortex.plan.eval(bool, rows=3) expr=($.a and $.b) child: vortex.plan.pack({a=bool, b=bool}, rows=3) a: vortex.plan.take(bool, rows=3) - codes: vortex.plan.filter(u8, rows=3) - child: vortex.plan.segment_scan(u8, rows=3) + codes: vortex.plan.segment_scan(u8, rows=3) values: vortex.plan.eval(bool, rows=2) expr=($ > 5i32) - child: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) + child: vortex.plan.segment_scan(i32, rows=2) b: vortex.plan.eval(bool, rows=3) expr=($ > 7i32) - child: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) + child: vortex.plan.segment_scan(i32, rows=3) "); let reoptimized = optimize(optimized.clone())?; assert!(PlanRef::ptr_eq(&optimized, &reoptimized)); @@ -1073,12 +1022,9 @@ fn multi_field_struct_expression_keeps_cross_field_refinement() -> VortexResult< root: vortex.plan.eval(bool, rows=3) expr=(($.a + $.b) > 10i32) child: vortex.plan.pack({a=i32, b=i32}, rows=3) a: vortex.plan.take(i32, rows=3) - codes: vortex.plan.filter(u8, rows=3) - child: vortex.plan.segment_scan(u8, rows=3) - values: vortex.plan.filter(i32, rows=2) - child: vortex.plan.segment_scan(i32, rows=2) - b: vortex.plan.filter(i32, rows=3) - child: vortex.plan.segment_scan(i32, rows=3) + codes: vortex.plan.segment_scan(u8, rows=3) + values: vortex.plan.segment_scan(i32, rows=2) + b: vortex.plan.segment_scan(i32, rows=3) "); assert_eq!( optimized.display_tree().to_string(), From fd843a9a5f9d0a5b004b235d903a279606301e75 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 9 Oct 2026 22:07:04 +0000 Subject: [PATCH 4/4] perf(layout): pack joins misaligned fields instead of cutting at every boundary Pack used to emit a struct as soon as every field had a front batch, cut to the shortest one. Fields chunked differently were sliced at every boundary of every other field: 10,000 misaligned columns made 8,192 structs and 80 million slices per scan. Now, when every field's front batch has the same length, Pack takes those batches whole, so aligned fields still stream one struct per chunk with no slicing. Otherwise it waits on the one field with the fewest queued rows until that field is full or closed, then takes that many rows from every field. Whole batches are joined as a ChunkedArray, not copied, and at most one batch per field is sliced. Memory stays bounded by the inlet capacities. Waiting costs O(1) per wake. Queued rows only grow until a struct is built, so Pack resumes at the first inlet it found empty instead of rescanning, and while filling it checks only the inlet it waits on. Signed-off-by: Claude Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01NWYmfHTFtaEdhdyd3u5hv3 --- vortex-layout/src/plan/pipeline/ops/mod.rs | 36 +++-- vortex-layout/src/plan/pipeline/ops/pack.rs | 140 +++++++++++++++----- vortex-layout/src/plan/pipeline/tests.rs | 64 +++++---- vortex-layout/src/plan/plans/concat.rs | 2 +- vortex-layout/src/plan/plans/pack.rs | 4 +- 5 files changed, 177 insertions(+), 69 deletions(-) diff --git a/vortex-layout/src/plan/pipeline/ops/mod.rs b/vortex-layout/src/plan/pipeline/ops/mod.rs index 7ed5d8943a0..f5506551a3d 100644 --- a/vortex-layout/src/plan/pipeline/ops/mod.rs +++ b/vortex-layout/src/plan/pipeline/ops/mod.rs @@ -17,23 +17,37 @@ pub(crate) use filter::*; pub(crate) use pack::*; pub(crate) use port::*; pub(crate) use scan::*; +use smallvec::SmallVec; use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::arrays::ChunkedArray; use vortex_error::VortexResult; use vortex_error::vortex_err; use crate::plan::pipeline::Inlet; -/// Takes the first `len` rows of the inlet's front batch: the batch itself when it is that -/// long, otherwise a slice, leaving the rest in place. +/// Takes the first `len` rows of the inlet, across as many batches as hold them, slicing the last +/// and leaving its rest in place. Rows spanning batches come back chunked, not copied. pub(crate) fn take_rows(inlet: &mut Inlet<'_>, len: usize) -> VortexResult { - let front = inlet - .peek_mut() - .ok_or_else(|| vortex_err!("Taking rows from an empty inlet"))?; - if front.len() == len { - return inlet - .take() - .ok_or_else(|| vortex_err!("Taking rows from an empty inlet")); + let mut pieces: SmallVec<[ArrayRef; 2]> = SmallVec::new(); + let mut left = len; + while left > 0 { + let front = inlet + .peek_mut() + .ok_or_else(|| vortex_err!("Taking {len} rows from an inlet holding fewer"))?; + if front.len() <= left { + left -= front.len(); + pieces.extend(inlet.take()); + } else { + let rest = front.slice(left..front.len())?; + pieces.push(std::mem::replace(front, rest).slice(0..left)?); + left = 0; + } } - let rest = front.slice(len..front.len())?; - std::mem::replace(front, rest).slice(0..len) + if pieces.len() == 1 { + return Ok(pieces.swap_remove(0)); + } + let dtype = pieces[0].dtype().clone(); + // SAFETY: the pieces come from one inlet, whose batches share its writer's dtype. + Ok(unsafe { ChunkedArray::new_unchecked(pieces, dtype) }.into_array()) } diff --git a/vortex-layout/src/plan/pipeline/ops/pack.rs b/vortex-layout/src/plan/pipeline/ops/pack.rs index a35cb50fbeb..3021d41dce8 100644 --- a/vortex-layout/src/plan/pipeline/ops/pack.rs +++ b/vortex-layout/src/plan/pipeline/ops/pack.rs @@ -6,7 +6,6 @@ use vortex_array::ArrayRef; use vortex_array::IntoArray; use vortex_array::arrays::StructArray; -use vortex_array::dtype::FieldNames; use vortex_array::dtype::Nullability; use vortex_array::dtype::StructFields; use vortex_array::validity::Validity; @@ -22,12 +21,31 @@ use crate::plan::pipeline::Operator; use crate::plan::pipeline::Source; use crate::plan::pipeline::Step; -/// Zips its inlets, one per field and the validity last when nullable, into structs, taking as -/// many rows from each as the shortest front batch holds. +/// Zips its inlets, one per field and the validity last when nullable, into structs. +/// +/// When every inlet's front batch has the same length, the fields are chunked alike there and the +/// struct takes those batches whole. Otherwise the struct takes as many rows as every inlet has +/// queued, once the inlet with the fewest can queue no more, so fields chunked differently are +/// joined, not cut at every boundary of every field. +/// +/// Waiting costs nothing per wake. Queued rows only grow until the struct is built, so the source +/// resumes at the first inlet it found empty instead of rescanning the others, and then waits on +/// the one inlet with the fewest rows. pub(crate) struct PackSource { fields: StructFields, nullable: bool, inlets: usize, + /// Inlets before this one hold rows. + cursor: usize, + /// The fewest rows queued on an inlet before `cursor` when the scan saw it, and that inlet. + fewest: (usize, usize), + /// The length of the front batch of every inlet before `cursor`, while they agree. + front: Option, + /// Inlets before `cursor` that had ended. + ended: usize, + /// Every inlet holds rows, and the source waits for this one, which holds the fewest, to + /// fill or close. + filling: Option, } impl PackSource { @@ -36,27 +54,77 @@ impl PackSource { fields, nullable, inlets, + cursor: 0, + fewest: (usize::MAX, 0), + front: None, + ended: 0, + filling: None, } } + + /// Finds whether every inlet holds rows, resuming at `cursor`. + fn scan(&mut self, cx: &mut Cx<'_>) -> Option { + while self.cursor < self.inlets { + let mut inlet = cx.inlet(self.cursor); + match inlet.peek_mut().map(|front| front.len()) { + None if inlet.closed() => self.ended += 1, + None => return Some(Step::Blocked(Blocked::Inlet(self.cursor))), + Some(front) => { + self.front = match self.cursor { + 0 => Some(front), + _ => self.front.filter(|&agreed| agreed == front), + }; + let rows = inlet.rows(); + if rows < self.fewest.0 { + self.fewest = (rows, self.cursor); + } + } + } + self.cursor += 1; + } + None + } } impl Operator for PackSource { fn compute(&mut self, _input: Input, cx: &mut Cx<'_>) -> VortexResult { - let mut len = usize::MAX; - let mut ended = 0; - for index in 0..self.inlets { - let mut inlet = cx.inlet(index); - let closed = inlet.closed(); - match inlet.peek_mut().map(|front| front.len()) { - Some(front) => len = len.min(front), - None if closed => ended += 1, - None => return Ok(Step::Blocked(Blocked::Inlet(index))), + if let Some(index) = self.filling { + let inlet = cx.inlet(index); + if !inlet.full() && !inlet.closed() { + return Ok(Step::Blocked(Blocked::Inlet(index))); } + // Every inlet may have gained rows while the source waited. + self.filling = None; + self.cursor = 0; + self.fewest = (usize::MAX, 0); + self.front = None; } - if ended == self.inlets { + if let Some(blocked) = self.scan(cx) { + return Ok(blocked); + } + if self.ended == self.inlets { return Ok(Step::Finished); } - vortex_ensure!(ended == 0, "Pack fields ended at different rows"); + vortex_ensure!(self.ended == 0, "Pack fields ended at different rows"); + let fewest = self.fewest.1; + let len = match self.front { + Some(front) => front, + None => { + let inlet = cx.inlet(fewest); + if !inlet.full() && !inlet.closed() { + self.filling = Some(fewest); + return Ok(Step::Blocked(Blocked::Inlet(fewest))); + } + // The counts the scan saw are lower bounds by now, so count again. + (0..self.inlets) + .map(|index| cx.inlet(index).rows()) + .min() + .unwrap_or(0) + } + }; + self.cursor = 0; + self.fewest = (usize::MAX, 0); + self.front = None; let mut arrays = Vec::with_capacity(self.inlets); let mut more = true; for index in 0..self.inlets { @@ -73,8 +141,7 @@ impl Operator for PackSource { } else { Validity::NonNullable }; - let array = StructArray::try_new_with_dtype(arrays, self.fields.clone(), len, validity)? - .into_array(); + let array = assemble(&self.fields, arrays, validity).into_array(); Ok(if more { Step::More(array) } else { @@ -91,36 +158,47 @@ impl Source for PackSource { /// Wraps each batch of a pack's only field into a struct. pub(crate) struct WrapStage { - names: FieldNames, + fields: StructFields, } impl WrapStage { - pub(crate) fn new(names: FieldNames) -> Self { - Self { names } + pub(crate) fn new(fields: StructFields) -> Self { + Self { fields } } } impl Operator for WrapStage { fn compute(&mut self, input: Input, _cx: &mut Cx<'_>) -> VortexResult { match input { - Input::Chunk(batch) => { - let len = batch.len(); - Ok(Step::Last( - StructArray::try_new( - self.names.clone(), - vec![batch], - len, - Validity::NonNullable, - )? - .into_array(), - )) - } + Input::Chunk(batch) => Ok(Step::Last( + assemble(&self.fields, vec![batch], Validity::NonNullable).into_array(), + )), Input::End => Ok(Step::Finished), Input::None => Ok(Step::Consumed), } } } +/// A struct of `fields` over arrays of one length. +/// +/// The type is not derived or validated per batch: a `Pack` plan checks when it is built that +/// each field's plan produces the field's dtype over the struct's rows, and a pack source takes +/// the same rows from every field. +fn assemble(fields: &StructFields, arrays: Vec, validity: Validity) -> StructArray { + let len = arrays.first().map_or(0, |array| array.len()); + debug_assert!( + arrays.iter().all(|array| array.len() == len) + && fields + .fields() + .zip(&arrays) + .all(|(dtype, array)| &dtype == array.dtype()), + "pack fields do not match the struct" + ); + // SAFETY: the plan validated every field's dtype against `fields`, and every array holds the + // same rows, checked above in debug builds. + unsafe { StructArray::new_unchecked(arrays, fields.clone(), len, validity) } +} + /// The struct of no fields over `len` rows. pub(crate) fn empty_struct( fields: StructFields, diff --git a/vortex-layout/src/plan/pipeline/tests.rs b/vortex-layout/src/plan/pipeline/tests.rs index 41d44dcc3a1..3d8d4845bd1 100644 --- a/vortex-layout/src/plan/pipeline/tests.rs +++ b/vortex-layout/src/plan/pipeline/tests.rs @@ -456,37 +456,55 @@ fn two_columns(store: &mut Store) -> VortexResult<(PlanRef, ArrayRef)> { Ok((lower(&layout)?, expected)) } -/// Pack emits a struct as soon as every field has rows, however the reads complete, and the -/// structs come out in row order, one per chunk of `a`. +/// Fields chunked differently are joined, not cut at each other's boundaries: Pack waits for the +/// field with the fewest rows to fill or close, then emits every row all fields hold, in row +/// order, however the reads complete. /// /// Segments: a0=0, a1=1, a2=2, b=3. #[rstest] -// `b` arrives last: nothing can be emitted before it. -#[case::b_last(&[1, 0, 2, 3], 4)] -// `b` and `a0` first: rows 0..7 go out before `a2` and `a1` arrive. -#[case::a_chunk_last(&[3, 0, 2, 1], 2)] -// Chunks in order: each chunk's struct goes out as it lands. -#[case::in_order(&[3, 0, 1, 2], 2)] -fn pack_streams_in_row_order( - #[case] order: &[u32], - #[case] deliveries_before_first_piece: usize, -) -> VortexResult<()> { +#[case::b_last(&[1, 0, 2, 3])] +#[case::a_chunk_last(&[3, 0, 2, 1])] +#[case::in_order(&[3, 0, 1, 2])] +fn pack_joins_misaligned_fields(#[case] order: &[u32]) -> VortexResult<()> { let mut store = Store::default(); let (plan, expected) = two_columns(&mut store)?; let mask = Mask::new_true(ROWS as usize); let run = run(&store, &plan, 0..ROWS, mask.clone(), scripted(order))?; - assert_eq!(pieces(&run.events), [7, 5, 8]); - let first_piece = run - .events - .iter() - .position(|e| matches!(e, Event::Piece(_))) - .vortex_expect("a piece"); - let delivered = run.events[..first_piece] - .iter() - .filter(|e| matches!(e, Event::Delivered(_))) - .count(); - assert_eq!(delivered, deliveries_before_first_piece); + assert_eq!(pieces(&run.events), [ROWS as usize]); + assert_view(&expected, &(0..ROWS), &mask, run.arrays) +} + +/// A misaligned field that fills its inlet makes Pack emit what every field holds, so a struct +/// never waits for more rows than the inlets can queue. +#[test] +fn pack_emits_when_the_shortest_field_fills() -> VortexResult<()> { + let mut store = Store::default(); + let a = PrimitiveArray::from_iter(0..ROWS as i32).into_array(); + let b = PrimitiveArray::from_iter((0..ROWS as i64).map(|v| v * 10)).into_array(); + let expected = StructArray::from_fields(&[("a", a.clone()), ("b", b.clone())])?.into_array(); + let layout = StructLayout::new( + ROWS, + expected.dtype().clone(), + vec![store.chunked(&a, &[1; ROWS as usize])?, store.flat(&b)?], + ) + .into_layout(); + let plan = lower(&layout)?; + let mask = Mask::new_true(ROWS as usize); + + let run = run( + &store, + &plan, + 0..ROWS, + mask.clone(), + delivery(Delivery::Fifo), + )?; + let pieces = pieces(&run.events); + assert!(pieces.len() > 1, "{pieces:?}"); + assert!( + pieces.iter().all(|&piece| piece <= DEFAULT_CAPACITY), + "{pieces:?}" + ); assert_view(&expected, &(0..ROWS), &mask, run.arrays) } diff --git a/vortex-layout/src/plan/plans/concat.rs b/vortex-layout/src/plan/plans/concat.rs index 852e58fda4b..d2562a6001f 100644 --- a/vortex-layout/src/plan/plans/concat.rs +++ b/vortex-layout/src/plan/plans/concat.rs @@ -157,7 +157,7 @@ impl PlanVTable for Concat { mask: &Mask, compiler: &mut Compiler<'_>, ) -> VortexResult> { - let mut chains = Vec::new(); + let mut chains = Vec::with_capacity(overlapping(plan, &rows).size_hint().1.unwrap_or(0)); for (index, start, local) in overlapping(plan, &rows) { let local_mask = mask.slice( usize::try_from(local.start - rows.start)? diff --git a/vortex-layout/src/plan/plans/pack.rs b/vortex-layout/src/plan/plans/pack.rs index 991bd9e03c9..89f786ba7fb 100644 --- a/vortex-layout/src/plan/plans/pack.rs +++ b/vortex-layout/src/plan/plans/pack.rs @@ -246,9 +246,7 @@ impl PlanVTable for Pack { if count == 1 && !nullable { // One field needs no zip: each of its batches is wrapped as it passes. let chain = chains.remove(0); - return Ok(Some( - chain.with(WrapStage::new(plan.fields().names().clone())), - )); + return Ok(Some(chain.with(WrapStage::new(plan.fields().clone())))); } let source = PackSource::new(plan.fields().clone(), nullable, count); Ok(Some(compiler.join(chains, source)))