Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions vortex-layout/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ pin-project-lite = { workspace = true }
prost = { workspace = true }
rustc-hash = { workspace = true }
sketches-ddsketch = { workspace = true }
smallvec = { workspace = true }
termtree = { workspace = true }
tokio = { workspace = true, features = ["rt"], optional = true }
tracing = { workspace = true }
Expand Down
1 change: 1 addition & 0 deletions vortex-layout/src/plan/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ mod display;
mod lower;
mod optimize;
pub mod optimizer;
pub mod pipeline;
mod plans;
mod typed;
mod vtable;
Expand Down
141 changes: 141 additions & 0 deletions vortex-layout/src/plan/pipeline/compile.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,141 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright the Vortex contributors

//! Compiling plans into pipelines, and counting the readers of each segment so segments read
//! more than once are decoded once.

use std::ops::Range;

use vortex_error::VortexResult;
use vortex_mask::Mask;
use vortex_session::VortexSession;

use super::Operator;
use super::Source;
use super::port::PortId;
use super::port::Reader;
use super::scan::Core;
use crate::plan::PlanRef;

/// 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.
pub struct Chain {
pub(crate) source: Box<dyn Source>,
pub(crate) stages: Vec<Box<dyn Operator>>,
pub(crate) inlets: Vec<PortId>,
}

impl Chain {
/// A chain of `source` alone.
pub fn new(source: impl Source + 'static) -> Self {
Self {
source: Box::new(source),
stages: Vec::new(),
inlets: Vec::new(),
}
}

/// This chain with `stage` after its last stage.
pub fn with(mut self, stage: impl Operator + 'static) -> Self {
self.stages.push(Box::new(stage));
self
}
}

impl Core {
/// Compiles `plan` over `rows`, restricted to `mask`. See [`Compiler::compile`].
pub(crate) fn compile(
&mut self,
plan: &PlanRef,
rows: Range<u64>,
mask: &Mask,
) -> VortexResult<Option<Chain>> {
Compiler { core: self }.compile(plan, rows, mask)
}
}

/// What a plan compiles its children through. See
/// [`PlanVTable::compile`](crate::plan::PlanVTable::compile).
pub struct Compiler<'a> {
pub(crate) core: &'a mut Core,
}

impl Compiler<'_> {
/// Compiles `plan` over `rows` of its domain, restricted to `mask`, into a chain producing
/// the rows `mask` selects, in order, or `None` when that is no row.
pub fn compile(
&mut self,
plan: &PlanRef,
rows: Range<u64>,
mask: &Mask,
) -> VortexResult<Option<Chain>> {
if rows.start >= rows.end {
return Ok(None);
}
plan.compile(rows, mask, self)
}

/// A chain reading `chains` through `source`: each chain becomes a pipeline writing a new
/// port, inlet `i` of `source` reading chain `i`, with the capacity `source` asks for.
pub fn join(&mut self, chains: Vec<Chain>, source: impl Source + 'static) -> Chain {
let inlets = chains
.into_iter()
.enumerate()
.map(|(index, chain)| {
let port = self
.core
.arena
.create(source.capacity(index), None, Reader::Unclaimed);
self.core.add_pipeline(chain, smallvec::smallvec![port]);
port
})
.collect();
Chain {
source: Box::new(source),
stages: Vec::new(),
inlets,
}
}

/// The global row index of the scanned plan's row zero.
pub fn row_offset(&self) -> u64 {
self.core.row_offset
}

/// The session used for decoding and evaluation.
pub fn session(&self) -> &VortexSession {
&self.core.session
}
}

/// How a plan's rows map to the rows of the plan a scan's splits range over.
#[derive(Clone, Debug)]
pub enum Reach {
/// Row `r` is scanned row `r + offset`.
Offset(u64),
/// Every row is read for these scanned rows, as a dictionary's values are for its codes.
Fixed(Range<u64>),
}

impl Reach {
/// The scanned rows `rows` are read for.
pub fn root(&self, rows: &Range<u64>) -> Range<u64> {
match self {
Reach::Offset(offset) => rows.start + offset..rows.end + offset,
Reach::Fixed(root) => root.clone(),
}
}

/// The mapping of a child whose row zero is this plan's row `by`.
pub fn shift(&self, by: u64) -> Self {
match self {
Reach::Offset(offset) => Reach::Offset(offset + by),
Reach::Fixed(root) => Reach::Fixed(root.clone()),
}
}

/// The mapping of a child read whole for `rows` of this plan, as dictionary values are.
pub fn fixed(&self, rows: &Range<u64>) -> Self {
Reach::Fixed(self.root(rows))
}
}
168 changes: 168 additions & 0 deletions vortex-layout/src/plan/pipeline/mod.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,168 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright the Vortex contributors

//! Push-based execution of physical plans as pipelines.
//!
//! A [`Scan`] runs a [`PlanRef`](crate::plan::PlanRef) over a list of splits, row ranges of the
//! plan's domain. Each split is compiled into pipelines: a pipeline is one [`Source`], zero or
//! more [`Operator`] stages, and the ports its last stage writes. The driver asks the source for
//! a batch and pushes it through the stages into the ports. Every port has exactly one writer and
//! one reader, and the reader's source declares how many batches it may hold.
//!
//! Only a source blocks, and only on one of three things: a segment read it requested, an empty
//! inlet, or a full outlet, which the driver detects. The scheduler re-runs a blocked pipeline
//! when that one condition clears.
//!
//! 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.

mod compile;
mod port;
mod scan;
pub mod synthetic;

use std::collections::VecDeque;

use vortex_array::ArrayRef;
use vortex_array::ExecutionCtx;
use vortex_array::buffer::BufferHandle;
use vortex_error::VortexResult;
use vortex_session::VortexSession;

pub use self::compile::Chain;
pub use self::compile::Compiler;
pub use self::compile::Reach;
use self::port::Arena;
pub use self::port::Inlet;
use self::port::PortId;
pub use self::scan::ReadId;
pub use self::scan::ReadRequest;
pub use self::scan::Scan;
pub use self::scan::Split;
pub use self::scan::Turn;
use crate::segments::SegmentId;

/// Batches an inlet holds before its writer is blocked, unless its reader says otherwise.
pub const DEFAULT_CAPACITY: usize = 8;

/// What a pipeline stage is handed on each call.
#[derive(Debug)]
pub enum Input {
/// One batch from the stage below.
Chunk(ArrayRef),
/// The stage below has finished. Flush whatever is held.
End,
/// Nothing new. Sent to a source on every call, and to a stage that returned
/// [`Step::More`].
None,
}

/// What a stage says back.
#[derive(Debug)]
pub enum Step {
/// A batch, and more is ready now without further input. The driver pushes the batch on
/// and calls again with [`Input::None`].
More(ArrayRef),
/// A batch, and the input is spent.
Last(ArrayRef),
/// The input was absorbed and nothing is ready.
Consumed,
/// Nothing can happen until the condition clears. Sources only.
Blocked(Blocked),
/// No more output will ever come. The driver sends [`Input::End`] to the next stage.
Finished,
}

/// Why a pipeline cannot run.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Blocked {
/// A read was requested and has not been delivered.
Io,
/// The inlet is empty and its writer has not closed it.
Inlet(usize),
/// An outlet is full. Produced by the driver, never by a stage.
Outlet,
}

/// A stage of a pipeline: takes one input at a time and pushes what it produces on.
///
/// A stage never blocks. A stage returning [`Step::More`] must eventually return something
/// else when called again with [`Input::None`]. A stage that has returned [`Step::Finished`] is
/// not called again. A stage receiving [`Input::End`] flushes, returning [`Step::More`] while it
/// has more to flush.
pub trait Operator: Send {
/// Takes `input` and says what happens next.
fn compute(&mut self, input: Input, cx: &mut Cx<'_>) -> VortexResult<Step>;
}

/// The start of a pipeline. Reads its inlets and requests segments rather than taking input
/// from a stage below. Always called with [`Input::None`].
pub trait Source: Operator {
/// The inlets this source reads, numbered from zero.
fn inlet_count(&self) -> usize {
0
}

/// How many batches inlet `inlet` may hold before its writer is blocked. Asked once, when
/// the port is made. An empty port always has room.
fn capacity(&self, inlet: usize) -> usize {
let _ = inlet;
DEFAULT_CAPACITY
}

/// A segment this source wants read, asked until it answers `None` before every
/// [`compute`](Operator::compute), so a source may have several reads in flight. A source
/// that asked for any is not computed again until one of them arrives; the bytes come
/// through [`Cx::take_segment`] in the order they arrive.
fn request(&mut self) -> Option<SegmentId> {
None
}
}

/// What a stage may touch during one call: its pipeline's inlets, the bytes it asked for, and
/// the session.
pub struct Cx<'a> {
arena: &'a mut Arena,
inlets: &'a [PortId],
bytes: &'a mut VecDeque<(SegmentId, BufferHandle)>,
session: &'a VortexSession,
exec: &'a mut ExecutionCtx,
/// The inlets read during this run, each once, so only their writers are checked for room
/// after it, and whether each inlet is listed.
touched: &'a mut Vec<usize>,
listed: &'a mut [bool],
}

impl Cx<'_> {
/// The reading side of inlet `index` of this pipeline's source.
pub fn inlet(&mut self, index: usize) -> Inlet<'_> {
if !self.listed[index] {
self.listed[index] = true;
self.touched.push(index);
}
self.arena.inlet(self.inlets[index])
}

/// The bytes of a segment the source asked for that has arrived, oldest arrival first.
pub fn take_bytes(&mut self) -> Option<BufferHandle> {
self.take_segment().map(|(_, bytes)| bytes)
}

/// Like [`take_bytes`](Self::take_bytes), with the segment the bytes are of.
pub fn take_segment(&mut self) -> Option<(SegmentId, BufferHandle)> {
self.bytes.pop_front()
}

/// The session used for decoding and evaluation.
pub fn session(&self) -> &VortexSession {
self.session
}

/// The execution context for executing arrays.
pub fn exec(&mut self) -> &mut ExecutionCtx {
self.exec
}
}

#[cfg(test)]
mod scheduling_tests;
Loading
Loading