diff --git a/vortex-layout/src/plan/mod.rs b/vortex-layout/src/plan/mod.rs index 92ece0c38dd..36cb5047149 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/pipeline/compile.rs b/vortex-layout/src/plan/pipeline/compile.rs index 0f5735b86aa..28fc9a52bbd 100644 --- a/vortex-layout/src/plan/pipeline/compile.rs +++ b/vortex-layout/src/plan/pipeline/compile.rs @@ -6,16 +6,24 @@ use std::ops::Range; +use rustc_hash::FxHashMap; +use smallvec::SmallVec; use vortex_error::VortexResult; use vortex_mask::Mask; use vortex_session::VortexSession; use super::Operator; use super::Source; +use super::ops::PortSource; +use super::ops::ScanSource; +use super::ops::SelectStage; use super::port::PortId; use super::port::Reader; use super::scan::Core; +use crate::plan::ConcatPlan; use crate::plan::PlanRef; +use crate::plan::SegmentScanPlan; +use crate::segments::SegmentId; /// A pipeline under construction: a source and the stages after it, not yet given outlets, so a /// parent plan can add stages before it becomes a pipeline. @@ -106,6 +114,192 @@ impl Compiler<'_> { pub fn session(&self) -> &VortexSession { &self.core.session } + + /// The rows `rows` of `scan`'s segment, filtered by `filter`, read by the split compiling. + /// A segment several readers read is decoded once and shared. + pub(crate) fn scan( + &mut self, + scan: &SegmentScanPlan, + rows: Range, + filter: Option, + ) -> VortexResult> { + let filter = match filter { + Some(filter) if filter.all_false() => return Ok(None), + Some(filter) if filter.all_true() => None, + filter => filter, + }; + let slice = if rows.start == 0 && rows.end == scan.row_count() { + None + } else { + Some(usize::try_from(rows.start)?..usize::try_from(rows.end)?) + }; + let split = self.core.split; + Ok(Some(match self.claim(scan, split) { + None => Chain::new(ScanSource::new(scan.clone(), slice, filter)), + Some(port) => { + let chain = Chain { + source: Box::new(PortSource), + stages: Vec::new(), + inlets: vec![port], + }; + if slice.is_none() && filter.is_none() { + chain + } else { + chain.with(SelectStage::new(slice, filter)) + } + } + })) + } + + /// Claims split `split`'s reader of `scan`'s segment: its port when the segment is + /// shared, building the pipeline that decodes it for the first reader and a port for every + /// reader in every split not yet finished, or `None` when this is its only reader. + fn claim(&mut self, scan: &SegmentScanPlan, split: usize) -> Option { + if split == usize::MAX { + return None; + } + let segment = scan.segment_id(); + let shares = &mut self.core.shares; + let ports = &mut shares.ports[split]; + if let Some(waiting) = ports.get_mut(&segment) { + let port = waiting.pop(); + if waiting.is_empty() { + ports.remove(&segment); + } + return port; + } + let ranges = std::mem::take(shares.reaches.get_mut(*segment as usize)?); + let mut readers: SmallVec<[usize; 4]> = SmallVec::new(); + for rows in &ranges { + readers.extend( + shares + .overlapping(rows) + .filter(|&reader| !shares.finished[reader] && shares.overlaps(reader, rows)), + ); + } + if readers.len() <= 1 { + return None; + } + let mut outlets: SmallVec<[PortId; 1]> = SmallVec::with_capacity(readers.len()); + let mut mine = None; + for reader in readers { + let port = self.core.arena.create(usize::MAX, None, Reader::Unclaimed); + outlets.push(port); + if reader == split && mine.is_none() { + mine = Some(port); + } else { + self.core.shares.ports[reader] + .entry(segment) + .or_default() + .push(port); + } + } + self.core.add_pipeline( + Chain::new(ScanSource::new(scan.clone(), None, None)), + outlets, + ); + mine + } +} + +/// Where each segment read more than once is read from, and the ports of the segments whose +/// decoding pipeline is built. +#[derive(Default)] +pub(crate) struct Shares { + /// The splits' row ranges, in split order. + splits: Vec>, + /// Whether the splits are sorted and disjoint, so the splits overlapping a range are found + /// by binary search. + ordered: bool, + /// Splits that have finished, and so will claim nothing. + finished: Vec, + /// By segment id, which a file numbers densely: for each segment more than one reader will + /// read and whose decoding pipeline is not built yet, the range of plan rows each of its + /// readers reads it for. Every split overlapping a range reads it once for that range. + reaches: Vec; 1]>>, + /// By split: the unclaimed ports of each shared segment the split reads. + ports: Vec>>, +} + +impl Shares { + pub(crate) fn new(splits: Vec>) -> Self { + let ordered = splits.windows(2).all(|pair| pair[0].end <= pair[1].start); + Self { + finished: vec![false; splits.len()], + ports: (0..splits.len()).map(|_| FxHashMap::default()).collect(), + splits, + ordered, + reaches: Vec::new(), + } + } + + /// The splits overlapping `rows`. + fn overlapping(&self, rows: &Range) -> Range { + if self.ordered { + let first = self.splits.partition_point(|split| split.end <= rows.start); + let end = self.splits.partition_point(|split| split.start < rows.end); + return first..end.max(first); + } + let mut overlapping = self + .splits + .iter() + .enumerate() + .filter(|(_, split)| split.start < rows.end && rows.start < split.end) + .map(|(index, _)| index); + match overlapping.next() { + Some(first) => first..overlapping.next_back().unwrap_or(first) + 1, + None => 0..0, + } + } + + fn overlaps(&self, split: usize, rows: &Range) -> bool { + let split = &self.splits[split]; + split.start < rows.end && rows.start < split.end + } + + /// Records the ranges `plan` reads each of its segments for, read over `rows` by every + /// split overlapping them. Call once per plan a split stage runs, before the scan starts, + /// then [`retain_shared`](Self::retain_shared). + pub(crate) fn add(&mut self, plan: &PlanRef, rows: Range) -> VortexResult<()> { + let reaches = &mut self.reaches; + plan.reach(rows, &Reach::Offset(0), &mut |segment, rows| { + record(reaches, segment, rows) + }) + } + + /// Keeps only the segments more than one reader reads. The rest are read by the one + /// pipeline that needs them. + pub(crate) fn retain_shared(&mut self) { + let mut reaches = std::mem::take(&mut self.reaches); + for ranges in &mut reaches { + let shared = match ranges.as_slice() { + [] => false, + // One reader per split overlapping the range: shared when the range crosses a + // split boundary. + [rows] => self + .overlapping(rows) + .filter(|&split| self.overlaps(split, rows)) + .nth(1) + .is_some(), + _ => true, + }; + if !shared { + *ranges = SmallVec::new(); + } + } + self.reaches = reaches; + } + + /// Drops the ports split `split` did not claim: its readers that never came, as a chunk + /// its mask ruled out or a stage it never ran. + pub(crate) fn finish(&mut self, split: usize, arena: &mut super::port::Arena) { + self.finished[split] = true; + for (_, ports) in std::mem::take(&mut self.ports[split]) { + for port in ports { + arena.drop_reader(port); + } + } + } } /// How a plan's rows map to the rows of the plan a scan's splits range over. @@ -139,3 +333,48 @@ impl Reach { Reach::Fixed(self.root(rows)) } } + +/// Adds a reader of `segment` reading it for `rows`. +fn record(reaches: &mut Vec; 1]>>, segment: SegmentId, rows: Range) { + let index = *segment as usize; + if index >= reaches.len() { + reaches.resize_with(index + 1, SmallVec::new); + } + reaches[index].push(rows); +} + +/// The chunks of `concat` overlapping `rows`: each chunk's index, where it starts, and the rows +/// of `rows` it holds, in the concatenation's domain. Found by binary search, so a split pays +/// for the chunks it overlaps, not for the chunks of the file. +pub(crate) fn overlapping<'a>( + concat: &'a ConcatPlan, + rows: &Range, +) -> impl Iterator)> + 'a { + let offsets = concat.row_offsets(); + let first = offsets + .partition_point(|&offset| offset <= rows.start) + .saturating_sub(1); + let end = offsets.partition_point(|&offset| offset < rows.end); + let rows = rows.clone(); + (first..end).filter_map(move |index| { + let start = offsets[index]; + let chunk_end = offsets + .get(index + 1) + .copied() + .unwrap_or_else(|| concat.row_count()); + let local = rows.start.max(start)..rows.end.min(chunk_end); + (local.start < local.end).then_some((index, start, local)) + }) +} + +/// Unclaimed shared ports, for tests. +#[cfg(test)] +impl Shares { + pub(crate) fn pending(&self) -> usize { + self.ports + .iter() + .flat_map(|ports| ports.values()) + .map(SmallVec::len) + .sum() + } +} diff --git a/vortex-layout/src/plan/pipeline/mod.rs b/vortex-layout/src/plan/pipeline/mod.rs index 1beed39fad1..f1a220657be 100644 --- a/vortex-layout/src/plan/pipeline/mod.rs +++ b/vortex-layout/src/plan/pipeline/mod.rs @@ -15,8 +15,28 @@ //! //! The owner of a scan answers reads: [`Scan::step`] returns each read as a [`Turn::Read`], and //! [`Scan::deliver`] hands its bytes back. Nothing in a scan awaits. +//! +//! # Masks +//! +//! Every mask a scan applies is known when the pipelines that apply it are compiled: the split's +//! selection, and the rows the earlier conjuncts kept. The +//! compiler places the mask in the stage that reads the rows, so chunks a mask selects nothing +//! of are never built and never read. +//! +//! # Sharing +//! +//! A segment several readers decode is decoded once. Before the scan starts, one pass over the +//! plans every split's stages run records the rows each segment is read for, so the splits that +//! will read it are found by binary search over the splits. The first reader of a segment read +//! more than once builds one pipeline that decodes it and fans the decoded array out into one +//! single-use port per reader, tagged with the reader's split: a reader in the same stage, a +//! later stage of a query, or a later split claims its own port. A split's ports it never +//! claims, because its rows were pruned or a stage never ran, are dropped when it finishes, and +//! a claimed port is dropped once its reader has drained it, so a decoded segment lives only as +//! long as a split that may still read it. mod compile; +pub(crate) mod ops; mod port; mod scan; pub mod synthetic; @@ -32,6 +52,7 @@ use vortex_session::VortexSession; pub use self::compile::Chain; pub use self::compile::Compiler; pub use self::compile::Reach; +pub(crate) use self::compile::overlapping; use self::port::Arena; pub use self::port::Inlet; use self::port::PortId; @@ -166,3 +187,5 @@ impl Cx<'_> { #[cfg(test)] mod scheduling_tests; +#[cfg(test)] +mod tests; diff --git a/vortex-layout/src/plan/pipeline/ops/concat.rs b/vortex-layout/src/plan/pipeline/ops/concat.rs new file mode 100644 index 00000000000..765efbb391e --- /dev/null +++ b/vortex-layout/src/plan/pipeline/ops/concat.rs @@ -0,0 +1,51 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! Concatenating inlets in order. + +use vortex_error::VortexResult; + +use crate::plan::pipeline::Blocked; +use crate::plan::pipeline::Cx; +use crate::plan::pipeline::Input; +use crate::plan::pipeline::Operator; +use crate::plan::pipeline::Source; +use crate::plan::pipeline::Step; + +/// Emits its inlets one after another, each to its end. +pub(crate) struct ConcatSource { + inlets: usize, + current: usize, +} + +impl ConcatSource { + pub(crate) fn new(inlets: usize) -> Self { + Self { inlets, current: 0 } + } +} + +impl Operator for ConcatSource { + fn compute(&mut self, _input: Input, cx: &mut Cx<'_>) -> VortexResult { + while self.current < self.inlets { + let mut inlet = cx.inlet(self.current); + if let Some(batch) = inlet.take() { + return Ok(if inlet.is_empty() && !inlet.closed() { + Step::Last(batch) + } else { + Step::More(batch) + }); + } + if !inlet.closed() { + return Ok(Step::Blocked(Blocked::Inlet(self.current))); + } + self.current += 1; + } + Ok(Step::Finished) + } +} + +impl Source for ConcatSource { + fn inlet_count(&self) -> usize { + self.inlets + } +} diff --git a/vortex-layout/src/plan/pipeline/ops/filter.rs b/vortex-layout/src/plan/pipeline/ops/filter.rs new file mode 100644 index 00000000000..d2ea761515b --- /dev/null +++ b/vortex-layout/src/plan/pipeline/ops/filter.rs @@ -0,0 +1,76 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! Keeping the selected rows of dense batches. + +use vortex_array::ArrayRef; +use vortex_array::Canonical; +use vortex_array::IntoArray; +use vortex_error::VortexResult; +use vortex_error::vortex_ensure; +use vortex_mask::Mask; + +use crate::plan::pipeline::Cx; +use crate::plan::pipeline::Input; +use crate::plan::pipeline::Operator; +use crate::plan::pipeline::Step; + +/// The selected fraction of a predicate's batch at or above which it is executed whole and then +/// filtered, as the default scan's flat reader does. +const EXPR_EVAL_THRESHOLD: f64 = 0.2; + +/// Keeps the rows of `array` that `mask` selects. `predicate` says the array is a predicate's +/// result, which, mostly selected, is cheaper to execute whole than to filter lazily. +pub(crate) fn keep_selected( + array: ArrayRef, + mask: Mask, + predicate: bool, + cx: &mut Cx<'_>, +) -> VortexResult { + if mask.all_true() { + return Ok(array); + } + if predicate && mask.density() >= EXPR_EVAL_THRESHOLD { + return array + .execute::(cx.exec())? + .into_array() + .filter(mask); + } + array.filter(mask) +} + +/// Keeps the selected rows of the dense batches below it, each by its own slice of the mask. +pub(crate) struct MaskStage { + mask: Mask, + cursor: usize, + predicate: bool, +} + +impl MaskStage { + pub(crate) fn new(mask: Mask, predicate: bool) -> Self { + Self { + mask, + cursor: 0, + predicate, + } + } +} + +impl Operator for MaskStage { + fn compute(&mut self, input: Input, cx: &mut Cx<'_>) -> VortexResult { + match input { + Input::Chunk(batch) => { + let end = self.cursor + batch.len(); + vortex_ensure!( + end <= self.mask.len(), + "Filter input is longer than its mask" + ); + let mask = self.mask.slice(self.cursor..end); + self.cursor = end; + Ok(Step::Last(keep_selected(batch, mask, self.predicate, cx)?)) + } + Input::End => Ok(Step::Finished), + Input::None => Ok(Step::Consumed), + } + } +} diff --git a/vortex-layout/src/plan/pipeline/ops/mod.rs b/vortex-layout/src/plan/pipeline/ops/mod.rs new file mode 100644 index 00000000000..f5506551a3d --- /dev/null +++ b/vortex-layout/src/plan/pipeline/ops/mod.rs @@ -0,0 +1,53 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! The sources and stages plans compile to. +//! +//! Every operator does work proportional to the rows it is handed: no operator looks at a batch +//! twice, holds more than the batches it must join, or scans inlets it does not read. + +mod concat; +mod filter; +mod pack; +mod port; +mod scan; + +pub(crate) use concat::*; +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, 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 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; + } + } + 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 new file mode 100644 index 00000000000..3021d41dce8 --- /dev/null +++ b/vortex-layout/src/plan/pipeline/ops/pack.rs @@ -0,0 +1,213 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! Zipping fields into structs. + +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::arrays::StructArray; +use vortex_array::dtype::Nullability; +use vortex_array::dtype::StructFields; +use vortex_array::validity::Validity; +use vortex_error::VortexResult; +use vortex_error::vortex_ensure; +use vortex_error::vortex_err; + +use super::*; +use crate::plan::pipeline::Blocked; +use crate::plan::pipeline::Cx; +use crate::plan::pipeline::Input; +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. +/// +/// 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 { + pub(crate) fn new(fields: StructFields, nullable: bool, inlets: usize) -> Self { + Self { + 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 { + 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 let Some(blocked) = self.scan(cx) { + return Ok(blocked); + } + if self.ended == self.inlets { + return Ok(Step::Finished); + } + 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 { + let mut inlet = cx.inlet(index); + arrays.push(take_rows(&mut inlet, len)?); + more &= !inlet.is_empty(); + } + let validity = if self.nullable { + Validity::Array( + arrays + .pop() + .ok_or_else(|| vortex_err!("Nullable Pack is missing its validity"))?, + ) + } else { + Validity::NonNullable + }; + let array = assemble(&self.fields, arrays, validity).into_array(); + Ok(if more { + Step::More(array) + } else { + Step::Last(array) + }) + } +} + +impl Source for PackSource { + fn inlet_count(&self) -> usize { + self.inlets + } +} + +/// Wraps each batch of a pack's only field into a struct. +pub(crate) struct WrapStage { + fields: StructFields, +} + +impl WrapStage { + 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) => 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, + nullability: Nullability, + len: usize, +) -> VortexResult { + let validity = match nullability { + Nullability::NonNullable => Validity::NonNullable, + Nullability::Nullable => Validity::AllValid, + }; + Ok(StructArray::try_new_with_dtype(Vec::new(), fields, len, validity)?.into_array()) +} diff --git a/vortex-layout/src/plan/pipeline/ops/port.rs b/vortex-layout/src/plan/pipeline/ops/port.rs new file mode 100644 index 00000000000..9a64629fa49 --- /dev/null +++ b/vortex-layout/src/plan/pipeline/ops/port.rs @@ -0,0 +1,55 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! Sources that read one port, or emit one batch. + +use vortex_array::ArrayRef; +use vortex_error::VortexResult; + +use crate::plan::pipeline::Blocked; +use crate::plan::pipeline::Cx; +use crate::plan::pipeline::Input; +use crate::plan::pipeline::Operator; +use crate::plan::pipeline::Source; +use crate::plan::pipeline::Step; + +/// Emits its inlet's batches as they arrive. +pub(crate) struct PortSource; + +impl Operator for PortSource { + fn compute(&mut self, _input: Input, cx: &mut Cx<'_>) -> VortexResult { + let mut inlet = cx.inlet(0); + match inlet.take() { + Some(batch) if inlet.is_empty() => Ok(Step::Last(batch)), + Some(batch) => Ok(Step::More(batch)), + None if inlet.closed() => Ok(Step::Finished), + None => Ok(Step::Blocked(Blocked::Inlet(0))), + } + } +} + +impl Source for PortSource { + fn inlet_count(&self) -> usize { + 1 + } +} + +/// Emits one batch it was built with. +pub(crate) struct OnceSource(Option); + +impl OnceSource { + pub(crate) fn new(batch: ArrayRef) -> Self { + Self(Some(batch)) + } +} + +impl Operator for OnceSource { + fn compute(&mut self, _input: Input, _cx: &mut Cx<'_>) -> VortexResult { + Ok(match self.0.take() { + Some(batch) => Step::Last(batch), + None => Step::Finished, + }) + } +} + +impl Source for OnceSource {} diff --git a/vortex-layout/src/plan/pipeline/ops/scan.rs b/vortex-layout/src/plan/pipeline/ops/scan.rs new file mode 100644 index 00000000000..7eba1b323fe --- /dev/null +++ b/vortex-layout/src/plan/pipeline/ops/scan.rs @@ -0,0 +1,144 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! Reading and decoding one segment. + +use std::ops::Range; + +use vortex_array::ArrayRef; +use vortex_array::serde::SerializedArray; +use vortex_error::VortexResult; +use vortex_error::vortex_bail; +use vortex_error::vortex_err; +use vortex_mask::Mask; +use vortex_session::VortexSession; + +use crate::plan::SegmentScanPlan; +use crate::plan::pipeline::Cx; +use crate::plan::pipeline::Input; +use crate::plan::pipeline::Operator; +use crate::plan::pipeline::Source; +use crate::plan::pipeline::Step; +use crate::segments::SegmentId; + +/// Decodes the whole of a scan's segment. +pub(crate) fn decode( + plan: &SegmentScanPlan, + segment: vortex_array::buffer::BufferHandle, + session: &VortexSession, +) -> VortexResult { + let serialized = match plan.array_tree() { + Some(tree) => SerializedArray::from_flatbuffer_and_segment(tree.clone(), segment)?, + None => SerializedArray::try_from(segment)?, + }; + serialized.decode( + plan.dtype(), + usize::try_from(plan.row_count())?, + plan.array_ctx(), + session, + ) +} + +/// Slices a whole decoded segment to `slice`, then keeps the rows `filter` selects. +fn select( + array: ArrayRef, + slice: Option<&Range>, + filter: Option<&Mask>, +) -> VortexResult { + let array = match slice { + Some(slice) => array.slice(slice.clone())?, + None => array, + }; + match filter { + Some(filter) => array.filter(filter.clone()), + None => Ok(array), + } +} + +enum ScanState { + Request, + Waiting, + Done, +} + +/// Reads one segment, decodes it, and emits its rows of `slice`, filtered by `filter`. +pub(crate) struct ScanSource { + plan: SegmentScanPlan, + slice: Option>, + filter: Option, + state: ScanState, +} + +impl ScanSource { + pub(crate) fn new( + plan: SegmentScanPlan, + slice: Option>, + filter: Option, + ) -> Self { + Self { + plan, + slice, + filter, + state: ScanState::Request, + } + } +} + +impl Operator for ScanSource { + fn compute(&mut self, _input: Input, cx: &mut Cx<'_>) -> VortexResult { + match self.state { + ScanState::Request => vortex_bail!("Segment scan computed before its read"), + ScanState::Waiting => { + let bytes = cx + .take_bytes() + .ok_or_else(|| vortex_err!("Segment scan ran without its bytes"))?; + self.state = ScanState::Done; + let array = decode(&self.plan, bytes, cx.session())?; + Ok(Step::Last(select( + array, + self.slice.as_ref(), + self.filter.as_ref(), + )?)) + } + ScanState::Done => Ok(Step::Finished), + } + } +} + +impl Source for ScanSource { + fn request(&mut self) -> Option { + match self.state { + ScanState::Request => { + self.state = ScanState::Waiting; + Some(self.plan.segment_id()) + } + ScanState::Waiting | ScanState::Done => None, + } + } +} + +/// Narrows a whole decoded segment, the one batch of a share's port, to a reader's rows. +pub(crate) struct SelectStage { + slice: Option>, + filter: Option, +} + +impl SelectStage { + pub(crate) fn new(slice: Option>, filter: Option) -> Self { + Self { slice, filter } + } +} + +impl Operator for SelectStage { + fn compute(&mut self, input: Input, _cx: &mut Cx<'_>) -> VortexResult { + match input { + Input::Chunk(batch) => Ok(Step::Last(select( + batch, + self.slice.as_ref(), + self.filter.as_ref(), + )?)), + Input::End => Ok(Step::Finished), + Input::None => Ok(Step::Consumed), + } + } +} diff --git a/vortex-layout/src/plan/pipeline/scan.rs b/vortex-layout/src/plan/pipeline/scan.rs index 2433fe1c23f..c90fc73c995 100644 --- a/vortex-layout/src/plan/pipeline/scan.rs +++ b/vortex-layout/src/plan/pipeline/scan.rs @@ -29,6 +29,7 @@ use super::Operator; use super::Source; use super::Step; use super::compile::Chain; +use super::compile::Shares; use super::port::Arena; use super::port::PipelineId; use super::port::PortId; @@ -252,10 +253,14 @@ pub(crate) struct Core { /// queued once. dirty: Vec, dirty_flags: Vec, + pub(crate) shares: Shares, + /// The split being compiled, whose shared readers a compile claims. `usize::MAX` when no + /// split is, as when a source asks for a plan mid-run. + pub(crate) split: usize, } impl Core { - fn new(session: VortexSession, row_offset: u64) -> Self { + fn new(session: VortexSession, row_offset: u64, splits: Vec>) -> Self { Self { exec: session.create_execution_ctx(), session, @@ -269,6 +274,8 @@ impl Core { new_reads: VecDeque::new(), dirty: Vec::new(), dirty_flags: Vec::new(), + shares: Shares::new(splits), + split: usize::MAX, } } @@ -316,8 +323,10 @@ impl Core { mask: &Mask, (slot, split): (usize, usize), ) -> VortexResult> { - let _ = split; - let Some(chain) = self.compile(plan, rows, mask)? else { + self.split = split; + let chain = self.compile(plan, rows, mask); + self.split = usize::MAX; + let Some(chain) = chain? else { return Ok(None); }; let port = self @@ -544,7 +553,15 @@ impl Scan { ); } let root = Root::Plan(plan); - let mut core = Core::new(session, 0); + let mut core = Core::new( + session, + 0, + splits.iter().map(|split| split.rows.clone()).collect(), + ); + match &root { + Root::Plan(plan) => core.shares.add(plan, 0..plan.row_count())?, + } + core.shares.retain_shared(); core.dirty_flags = vec![false; 1]; Ok(Self { core, @@ -630,6 +647,8 @@ impl Scan { if output.is_some() { self.active[slot] = Some(Active { index, output }); self.live += 1; + } else { + self.core.shares.finish(index, &mut self.core.arena); } Ok(()) } @@ -651,6 +670,7 @@ impl Scan { self.core.arena.drop_reader(port); active.output = None; } + self.core.shares.finish(active.index, &mut self.core.arena); self.active[slot] = None; self.live -= 1; Ok(()) @@ -661,4 +681,10 @@ impl Scan { pub(crate) fn live_ports(&self) -> usize { self.core.arena.live() } + + /// Shared ports not yet claimed or dropped. + #[cfg(test)] + pub(crate) fn pending_shares(&self) -> usize { + self.core.shares.pending() + } } diff --git a/vortex-layout/src/plan/pipeline/tests.rs b/vortex-layout/src/plan/pipeline/tests.rs new file mode 100644 index 00000000000..3d8d4845bd1 --- /dev/null +++ b/vortex-layout/src/plan/pipeline/tests.rs @@ -0,0 +1,604 @@ +// 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::ChunkedArray; +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::VortexExpect; +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::PlanRef; +use crate::plan::SegmentScan; +use crate::plan::lower; +use crate::segments::SegmentId; +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 read the scripted IO service completes next. +#[derive(Clone, Copy, Debug)] +enum Delivery { + /// Oldest read first. + Fifo, + /// Newest read first. + Lifo, +} + +#[derive(Debug, PartialEq, Eq)] +enum Event { + /// Reads the scan asked for before it next waited, as segment ids. + Io(Vec), + /// A read delivered for this segment. + Delivered(u32), + /// An array of this many rows came out. + Piece(usize), +} + +struct Run { + /// Every array, in the order it came out. + arrays: Vec, + /// The arrays of each split, in order. + splits: Vec>, + events: Vec, + /// Ports and shared readers still held once the scan finished. + leftover: (usize, usize), +} + +/// Drives a scan the way an owner does, completing reads only while the scan waits. +fn drive( + store: &Store, + mut scan: Scan, + mut pick: impl FnMut(&[ReadRequest]) -> usize, +) -> VortexResult { + let mut inflight: Vec = Vec::new(); + let mut batch: Vec = Vec::new(); + let mut arrays = Vec::new(); + let mut splits: Vec> = Vec::new(); + let mut events = Vec::new(); + loop { + match scan.step()? { + Turn::Read(read) => { + batch.push(*read.segment_id); + inflight.push(read); + } + Turn::Output(split, array) => { + if !batch.is_empty() { + events.push(Event::Io(std::mem::take(&mut batch))); + } + assert!(!array.is_empty(), "a scan never emits an empty array"); + events.push(Event::Piece(array.len())); + if split >= splits.len() { + splits.resize_with(split + 1, Vec::new); + } + splits[split].push(array.clone()); + arrays.push(array); + } + Turn::Waiting => { + if !batch.is_empty() { + events.push(Event::Io(std::mem::take(&mut batch))); + } + if inflight.is_empty() { + return Err(vortex_err!("scan waits with no reads in flight")); + } + let read = inflight.remove(pick(&inflight)); + events.push(Event::Delivered(*read.segment_id)); + scan.deliver(read.id, store.read(read.segment_id))?; + } + Turn::Done => break, + } + } + assert!( + batch.is_empty() && inflight.is_empty(), + "scan finished with reads in flight" + ); + let leftover = (scan.live_ports(), scan.pending_shares()); + Ok(Run { + arrays, + splits, + events, + leftover, + }) +} + +/// Runs `plan` over one split. +fn run( + store: &Store, + plan: &PlanRef, + rows: Range, + mask: Mask, + pick: impl FnMut(&[ReadRequest]) -> usize, +) -> VortexResult { + let scan = Scan::try_new(SESSION.clone(), plan.clone(), vec![Split { rows, mask }])?; + let run = drive(store, scan, pick)?; + assert_eq!(run.leftover, (0, 0), "a finished scan holds no port"); + Ok(run) +} + +fn delivery(order: Delivery) -> impl FnMut(&[ReadRequest]) -> 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(&[ReadRequest]) -> 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() +} + +/// Every segment read, sorted. +fn segments_read(events: &[Event]) -> Vec { + let mut read: Vec = events + .iter() + .flat_map(|event| match event { + Event::Io(batch) => batch.clone(), + _ => Vec::new(), + }) + .collect(); + read.sort_unstable(); + read +} + +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()), + } + } +} + +/// Joins arrays covering consecutive rows into one. +fn join(dtype: &vortex_array::dtype::DType, arrays: Vec) -> VortexResult { + Ok(ChunkedArray::try_new(arrays, dtype.clone())?.into_array()) +} + +/// 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) +} + +/// Splits cutting every column at different places, run as one scan, each produce their own +/// rows, and a segment several splits read is read and decoded once, then dropped. +#[rstest] +fn splits_of_one_scan_read_each_segment_once( + #[values([0, 4, 11, 20], [0, 9, 13, 20], [0, 1, 2, 20])] cuts: [u64; 4], + #[values(1, 3)] active: usize, + #[values(Delivery::Fifo, Delivery::Lifo)] order: Delivery, +) -> VortexResult<()> { + let mut store = Store::default(); + let (plan, expected) = fixture(&mut store)?; + let splits: Vec = cuts + .windows(2) + .map(|w| Split { + rows: w[0]..w[1], + mask: Sel::EveryOther.mask((w[1] - w[0]) as usize), + }) + .collect(); + let scan = Scan::try_new(SESSION.clone(), plan, splits.clone())?.with_max_active(active); + let run = drive(&store, scan, delivery(order))?; + assert_eq!(run.leftover, (0, 0), "every shared segment was dropped"); + assert_eq!( + segments_read(&run.events), + (0..store.segments.len() as u32).collect::>() + ); + for (index, split) in splits.iter().enumerate() { + assert_view( + &expected, + &split.rows, + &split.mask, + run.splits.get(index).cloned().unwrap_or_default(), + )?; + } + Ok(()) +} + +#[test] +fn every_read_is_issued_before_any_is_delivered() -> 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)) +} + +/// 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] +#[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), [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) +} + +/// 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) +} + +/// 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 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 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); + 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(()) +} + +/// A column read by two fields of one struct is read and decoded once, and both fields get it. +#[test] +fn a_segment_two_readers_need_is_read_once() -> VortexResult<()> { + let mut store = Store::default(); + let a = PrimitiveArray::from_iter(0..ROWS as i32).into_array(); + let column = lower(&store.chunked(&a, &[7, 13])?)?; + let expected = StructArray::from_fields(&[("x", a.clone()), ("y", a)])?.into_array(); + let fields = expected + .dtype() + .as_struct_fields_opt() + .vortex_expect("struct") + .clone(); + let plan = crate::plan::PackPlan::try_new( + fields, + vortex_array::dtype::Nullability::NonNullable, + ROWS, + vec![column.clone(), column], + None, + )? + .into_plan(); + let rows = 3..17; + let mask = Sel::EveryOther.mask(14); + + let run = run( + &store, + &plan, + rows.clone(), + mask.clone(), + delivery(Delivery::Lifo), + )?; + assert_eq!(segments_read(&run.events), [0, 1]); + assert_view(&expected, &rows, &mask, run.arrays) +} diff --git a/vortex-layout/src/plan/plans/concat.rs b/vortex-layout/src/plan/plans/concat.rs index 2e814156309..d2562a6001f 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; @@ -19,6 +21,12 @@ use crate::plan::PlanParts; use crate::plan::PlanRef; use crate::plan::PlanVTable; use crate::plan::optimizer::PlanParentReduceRule; +use crate::plan::pipeline::Chain; +use crate::plan::pipeline::Compiler; +use crate::plan::pipeline::Reach; +use crate::plan::pipeline::ops::ConcatSource; +use crate::plan::pipeline::overlapping; +use crate::segments::SegmentId; /// Concatenates its children row-wise. #[derive(Clone, Debug)] @@ -142,6 +150,53 @@ impl PlanVTable for Concat { fn child_name(_plan: &Plan, index: usize) -> Cow<'_, str> { Cow::Owned(format!("chunks[{index}]")) } + + fn compile( + plan: &Plan, + rows: Range, + mask: &Mask, + compiler: &mut Compiler<'_>, + ) -> VortexResult> { + 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)? + ..usize::try_from(local.end - rows.start)?, + ); + // A chunk the mask selects nothing of is never built, so never read. + if local_mask.all_false() { + continue; + } + let chunk = plan.child_required(index)?; + if let Some(chain) = + compiler.compile(&chunk, local.start - start..local.end - start, &local_mask)? + { + chains.push(chain); + } + } + Ok(match chains.len() { + 0 => None, + // Rows inside one chunk are that chunk's rows: no concatenation is built. + 1 => chains.pop(), + count => Some(compiler.join(chains, ConcatSource::new(count))), + }) + } + + fn reach( + plan: &Plan, + rows: Range, + at: &Reach, + visit: &mut dyn FnMut(SegmentId, Range), + ) -> VortexResult<()> { + for (index, start, local) in overlapping(plan, &rows) { + plan.child_required(index)?.reach( + local.start - start..local.end - start, + &at.shift(start), + visit, + )?; + } + Ok(()) + } } /// 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 new file mode 100644 index 00000000000..8a38a2c06d6 --- /dev/null +++ b/vortex-layout/src/plan/plans/filter.rs @@ -0,0 +1,133 @@ +// SPDX-License-Identifier: Apache-2.0 +// 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; +use crate::plan::PlanChildren; +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::pipeline::Chain; +use crate::plan::pipeline::Compiler; +use crate::plan::pipeline::Reach; +use crate::plan::pipeline::ops::MaskStage; +use crate::segments::SegmentId; + +/// 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}]")) + } + } + + /// Fuses with a segment-scan child, whose source keeps the selected rows itself; over any + /// other child, compiles the child over every row and keeps the selected rows of each batch. + fn compile( + plan: &Plan, + rows: Range, + mask: &Mask, + compiler: &mut Compiler<'_>, + ) -> VortexResult> { + let child = plan.child_plan()?; + if let Some(scan) = child.as_opt::() { + return compiler.scan(scan, rows, Some(mask.clone())); + } + if mask.all_false() { + return Ok(None); + } + let len = usize::try_from(rows.end - rows.start)?; + let predicate = plan.dtype().is_boolean(); + let chain = compiler.compile(&child, rows, &Mask::new_true(len))?; + Ok(chain.map(|chain| { + if mask.all_true() { + chain + } else { + chain.with(MaskStage::new(mask.clone(), predicate)) + } + })) + } + + fn reach( + plan: &Plan, + rows: Range, + at: &Reach, + visit: &mut dyn FnMut(SegmentId, Range), + ) -> VortexResult<()> { + plan.child_plan()?.reach(rows, at, visit) + } +} 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/plans/pack.rs b/vortex-layout/src/plan/plans/pack.rs index 6497cc63d45..89f786ba7fb 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; @@ -38,6 +40,14 @@ use crate::plan::PlanParts; use crate::plan::PlanRef; use crate::plan::PlanVTable; use crate::plan::optimizer::PlanParentReduceRule; +use crate::plan::pipeline::Chain; +use crate::plan::pipeline::Compiler; +use crate::plan::pipeline::Reach; +use crate::plan::pipeline::ops::OnceSource; +use crate::plan::pipeline::ops::PackSource; +use crate::plan::pipeline::ops::WrapStage; +use crate::plan::pipeline::ops::empty_struct; +use crate::segments::SegmentId; /// Assembles a struct from one child per field, plus an optional trailing validity child. #[derive(Clone, Debug)] @@ -204,6 +214,55 @@ impl PlanVTable for Pack { } Cow::Borrowed("validity") } + + fn compile( + plan: &Plan, + rows: Range, + mask: &Mask, + compiler: &mut Compiler<'_>, + ) -> VortexResult> { + let count = plan.children().len(); + if count == 0 { + // A struct with no fields still has rows, so no child can carry them. + let array = empty_struct( + plan.fields().clone(), + plan.dtype().nullability(), + mask.true_count(), + )?; + return Ok(Some(Chain::new(OnceSource::new(array)))); + } + if mask.all_false() { + return Ok(None); + } + let mut chains = Vec::with_capacity(count); + for child in plan.children().iter() { + chains.push( + compiler + .compile(&child?, rows.clone(), mask)? + .ok_or_else(|| vortex_err!("A Pack field produced no rows"))?, + ); + } + let nullable = plan.dtype().is_nullable(); + 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().clone())))); + } + let source = PackSource::new(plan.fields().clone(), nullable, count); + Ok(Some(compiler.join(chains, source))) + } + + fn reach( + plan: &Plan, + rows: Range, + at: &Reach, + visit: &mut dyn FnMut(SegmentId, Range), + ) -> VortexResult<()> { + for child in plan.children().iter() { + child?.reach(rows.clone(), at, visit)?; + } + Ok(()) + } } 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..9ace7e43ff9 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,9 @@ use crate::plan::PlanId; use crate::plan::PlanParts; use crate::plan::PlanVTable; use crate::plan::check_child_count; +use crate::plan::pipeline::Chain; +use crate::plan::pipeline::Compiler; +use crate::plan::pipeline::Reach; use crate::segments::SegmentId; /// Reads one serialized array segment. @@ -92,4 +98,24 @@ impl PlanVTable for SegmentScan { check_child_count("SegmentScan", children, 0)?; Ok(()) } + + fn compile( + plan: &Plan, + rows: Range, + mask: &Mask, + compiler: &mut Compiler<'_>, + ) -> VortexResult> { + // A scan keeps the selected rows itself, as every plan produces only those. + compiler.scan(plan, rows, Some(mask.clone())) + } + + fn reach( + plan: &Plan, + rows: Range, + at: &Reach, + visit: &mut dyn FnMut(SegmentId, Range), + ) -> VortexResult<()> { + visit(plan.segment_id(), at.root(&rows)); + Ok(()) + } }