Skip to main content

rudb_exec/
query.rs

1//! A built query: the pipelines it runs, in the order they have to run, and the queue the rows come
2//! out of.
3//!
4//! This is what replaced the pull tree. A plan used to become a tree of operators whose root was
5//! pulled from, and each pipeline breaker in it drained the tree below it on the first pull. The
6//! order that produced was right, because a breaker cannot answer until its input is finished, but
7//! it was an order the call stack happened to have rather than one anybody wrote down. Here it is
8//! written down: [`Query::run`] takes the pipelines in dependency order and runs each of them to
9//! completion.
10//!
11//! # Where the threads are
12//!
13//! Inside one pipeline and not across them. Each pipeline runs on as many threads as
14//! [`Pipeline::degree`] says, which is bounded by what the database's [`Pool`] will lend, by
15//! whether every operator in it will run as more than one instance, and by how many morsels its
16//! source has. Then the next one starts.
17//!
18//! Running two pipelines of one query at the same time is the other kind of parallelism and it is
19//! not here. The dependency edges say which pairs could overlap, so the information is already
20//! written down, and what is missing is a scheduler that holds several pipelines at once rather
21//! than a driver that is handed one. It is also worth much less: the shapes in ClickBench are a
22//! scan feeding an aggregate feeding a sort, which is a chain, and a chain has nothing to overlap.
23
24use std::sync::Arc;
25
26use rudb_common::{Cancel, Error, Result};
27use rudb_metrics::Driver;
28use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
29
30use rudb_pipeline::{Pipeline, Pool, RootReader, run_parallel};
31use rudb_vector::Chunk;
32
33use crate::schema::Schema;
34
35/// A plan that has been built and is ready to run.
36///
37/// It borrows the plan and the catalog it was built from, which is what `'a` is. A scan reads its
38/// rows out of the catalog's table rather than copying them and an expression reads its constants
39/// out of the plan's arena, so a query cannot outlive either.
40#[derive(Debug)]
41pub struct Query<'a> {
42    /// The pipelines, in an order where everything a pipeline waits for comes before it.
43    pipelines: Vec<Pipeline<'a>>,
44    /// The driver counters for each pipeline, in the same order.
45    drivers: Vec<Arc<Driver>>,
46    /// Where the last pipeline puts its rows.
47    reader: Option<RootReader>,
48    /// What the query produces.
49    schema: Schema,
50    /// CPU nanoseconds burned on threads other than the one that called [`Query::run`].
51    worker_cpu_ns: AtomicU64,
52    /// The most instances any one pipeline ran as.
53    widest: AtomicUsize,
54    /// Whether a pipeline reads a stored table or a file, rather than only statistics kept about
55    /// one, constants and the rows of other pipelines.
56    reads: bool,
57}
58
59impl<'a> Query<'a> {
60    /// A query over pipelines that are already in dependency order, each paired with its driver.
61    ///
62    /// # Errors
63    ///
64    /// [`ErrorCode::Internal`](rudb_common::ErrorCode::Internal) if a pipeline waits for one that
65    /// does not come before it. That is a builder bug rather than anything a query can cause, and it
66    /// is checked here because running the pipelines in the wrong order reads a buffer nobody has
67    /// filled yet and answers with no rows rather than failing.
68    pub(crate) fn new(
69        pipelines: Vec<Pipeline<'a>>,
70        drivers: Vec<Arc<Driver>>,
71        reader: Option<RootReader>,
72        schema: Schema,
73    ) -> Result<Self> {
74        for (at, pipeline) in pipelines.iter().enumerate() {
75            for waited in pipeline.depends_on() {
76                let before = pipelines[..at].iter().any(|earlier| earlier.id() == *waited);
77                if !before {
78                    return Err(Error::internal(format!(
79                        "{} waits for {waited}, which the builder did not put before it",
80                        pipeline.id()
81                    )));
82                }
83            }
84        }
85        Ok(Self {
86            pipelines,
87            drivers,
88            reader,
89            schema,
90            worker_cpu_ns: AtomicU64::new(0),
91            widest: AtomicUsize::new(0),
92            reads: true,
93        })
94    }
95
96    /// The same query, saying whether it reads a stored table or a file.
97    pub(crate) fn reading(mut self, reads: bool) -> Self {
98        self.reads = reads;
99        self
100    }
101
102    /// Whether running the query reads a stored table or a file. A query that does not is answered
103    /// out of the statistics kept about its tables, or out of constants, and runs in no time.
104    #[must_use]
105    pub fn reads_tables(&self) -> bool {
106        self.reads
107    }
108
109    /// The columns this query produces.
110    #[must_use]
111    pub fn schema(&self) -> &Schema {
112        &self.schema
113    }
114
115    /// How many pipelines the query runs.
116    #[must_use]
117    pub fn pipelines(&self) -> usize {
118        self.pipelines.len()
119    }
120
121    /// Runs every pipeline, stopping at the first one that fails.
122    ///
123    /// Each one is timed against its own driver, which is the loop that runs a pipeline rather than
124    /// any operator in it. That time is not nothing: on a scan of ten million rows the loop goes
125    /// round ten thousand times, and none of it sits inside an operator's own span, so without a
126    /// driver it is time the metrics document cannot account for.
127    ///
128    /// The lease is taken per pipeline and given back at the end of it, so a query whose scan uses
129    /// nine threads and whose sort uses one holds nine for as long as the scan and one after that,
130    /// and the threads it is not using are there for whatever else the database is running.
131    ///
132    /// How many threads it borrows and how many instances it runs are two numbers. The instances
133    /// are what the source has work for. The borrow is the wider of that and what the sink says it
134    /// can finish on, because the finish happens on the same threads with every instance already
135    /// joined, and a hash aggregate merging a million groups is not the same width as the scan that
136    /// fed it.
137    ///
138    /// # Errors
139    ///
140    /// Whatever any operator reports, or [`ErrorCode::Interrupt`](rudb_common::ErrorCode::Interrupt)
141    /// if the token says to stop. The check is per chunk, in the driver, which is why no operator
142    /// here holds a token of its own except the join, whose nested loop can outlive a chunk.
143    pub fn run(&self, cancel: &Cancel, pool: &Pool) -> Result<()> {
144        for (pipeline, driver) in self.pipelines.iter().zip(&self.drivers) {
145            // Both numbers out of one call, because asking is what makes the source read its
146            // statistics and cut its morsels. See [`Pipeline::widths`].
147            let (wanted, width) = pipeline.widths(pool.threads());
148            let lease = pool.lease(width);
149            let degree = wanted.min(lease.degree());
150            let spread = {
151                let _running = driver.running();
152                run_parallel(pipeline, cancel, &lease, degree)?
153            };
154            driver.ran(degree, spread.worker_cpu_ns);
155            driver.waited(
156                spread.slowest_ns,
157                spread.slowest_cpu_ns,
158                spread.finalize_ns,
159                spread.stagger_ns,
160            );
161            self.worker_cpu_ns.fetch_add(spread.worker_cpu_ns, Ordering::Relaxed);
162            self.widest.fetch_max(degree, Ordering::Relaxed);
163        }
164        Ok(())
165    }
166
167    /// CPU nanoseconds this query burned on threads other than the one that ran it.
168    ///
169    /// A caller timing the execution reads its own thread's CPU clock, which is the only clock
170    /// there is that attributes work to the thread that did it, and which therefore cannot see the
171    /// workers. This is what it missed.
172    #[must_use]
173    pub fn worker_cpu_ns(&self) -> u64 {
174        self.worker_cpu_ns.load(Ordering::Relaxed)
175    }
176
177    /// The most instances any one pipeline of this query ran as.
178    ///
179    /// Not the setting and not an average. A query whose scan ran on nine threads and whose sort ran
180    /// on one reports nine, because the question this answers is what the query was able to use.
181    #[must_use]
182    pub fn widest(&self) -> usize {
183        self.widest.load(Ordering::Relaxed)
184    }
185
186    /// Say that these rows are going out of the engine, so every column is flat when it arrives.
187    ///
188    /// A query whose rows go into a table does not call this and gets the forms the operators
189    /// produced, which is what storage wants and is why this is asked for rather than always done.
190    /// A query whose rows go to a caller calls it before [`Query::run`], and then the flattening
191    /// happens on the worker that produced the chunk. Doing it afterwards, on the one thread that
192    /// drains the queue, is the same work in the one place in a parallel query where the rest of
193    /// the pool has nothing to do but wait for it.
194    pub fn for_a_caller(&self) {
195        if let Some(reader) = &self.reader {
196            reader.flattening();
197        }
198    }
199
200    /// How long turning the answer into flat columns for a caller took, summed over threads.
201    ///
202    /// Zero unless [`Query::for_a_caller`] asked for it.
203    #[must_use]
204    pub fn flattened_ns(&self) -> u64 {
205        self.reader.as_ref().map_or(0, RootReader::flattened_ns)
206    }
207
208    /// The next chunk of the answer, or `None` when there are no more.
209    ///
210    /// Only meaningful after [`Query::run`] has returned. The serial driver runs a pipeline to
211    /// completion, so everything the query produced is queued by then, and taking a chunk here
212    /// removes it from the queue rather than copying it out.
213    ///
214    /// # Errors
215    ///
216    /// [`ErrorCode::Internal`](rudb_common::ErrorCode::Internal) if a thread panicked while holding
217    /// the queue.
218    pub fn next_chunk(&self) -> Result<Option<Chunk>> {
219        let chunk = self
220            .reader
221            .as_ref()
222            .ok_or_else(|| Error::internal("a query built into a sink has no result reader"))?
223            .next_chunk()?;
224        if let Some(chunk) = &chunk {
225            chunk.validate_external()?;
226        }
227        Ok(chunk)
228    }
229
230    /// Runs the query and collects everything it produced.
231    ///
232    /// The convenience the tests and the simple callers want. A caller that cares about holding one
233    /// chunk at a time calls [`Query::run`] and [`Query::next_chunk`] itself.
234    ///
235    /// # Errors
236    ///
237    /// The same as [`Query::run`].
238    pub fn collect(&self, cancel: &Cancel, pool: &Pool) -> Result<Vec<Chunk>> {
239        self.run(cancel, pool)?;
240        let mut chunks = Vec::new();
241        while let Some(chunk) = self.next_chunk()? {
242            chunks.push(chunk);
243        }
244        Ok(chunks)
245    }
246}