Skip to content
3 changes: 3 additions & 0 deletions vortex-layout/src/plan/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
239 changes: 239 additions & 0 deletions vortex-layout/src/plan/pipeline/compile.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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<u64>,
filter: Option<Mask>,
) -> VortexResult<Option<Chain>> {
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<PortId> {
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<Range<u64>>,
/// 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<bool>,
/// 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<SmallVec<[Range<u64>; 1]>>,
/// By split: the unclaimed ports of each shared segment the split reads.
ports: Vec<FxHashMap<SegmentId, SmallVec<[PortId; 1]>>>,
}

impl Shares {
pub(crate) fn new(splits: Vec<Range<u64>>) -> 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<u64>) -> Range<usize> {
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<u64>) -> 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<u64>) -> 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.
Expand Down Expand Up @@ -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<SmallVec<[Range<u64>; 1]>>, segment: SegmentId, rows: Range<u64>) {
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<u64>,
) -> impl Iterator<Item = (usize, u64, Range<u64>)> + '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()
}
}
23 changes: 23 additions & 0 deletions vortex-layout/src/plan/pipeline/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -166,3 +187,5 @@ impl Cx<'_> {

#[cfg(test)]
mod scheduling_tests;
#[cfg(test)]
mod tests;
51 changes: 51 additions & 0 deletions vortex-layout/src/plan/pipeline/ops/concat.rs
Original file line number Diff line number Diff line change
@@ -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<Step> {
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
}
}
Loading
Loading