rudb-exec 0.3.66

Operators, morsels, the scheduler, hash tables, sorting and spilling.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
//! Sorting when only the first few rows are wanted.
//!
//! `ORDER BY x LIMIT 10` over a hundred million rows has the same answer as a sort with a limit over
//! it and nothing like the same cost. A sort holds every row it was given, because the last row of
//! the input can be the first row of the output, so the memory it takes is the size of the input and
//! the time it takes is a full sort of it. Ten rows are wanted. This holds the rows that could still
//! come out and throws the rest away as it goes.
//!
//! What has to be held is `count + offset` rows, not `count`, since the rows that are skipped still
//! have to be found before there is anything to skip them from.
//!
//! # A sorted bound rather than a heap
//!
//! The textbook answer is a binary heap of the bound, pushing every row and popping the worst. This
//! keeps the candidates sorted instead. A row first compares with the worst candidate and is
//! discarded without materializing its payload when it loses. A winner is inserted after existing
//! equal keys, which preserves input order for ties the same way the stable full sort does.
//!
//! Binary search makes the comparison cost `O(m log n)`, as it is for a heap. Inserting moves `n`
//! small row handles, but only a shrinking share of the input wins after the first `n` rows. The
//! important saving is that almost every row pays for its key and never becomes a heap-allocated
//! payload row. A large offset makes those moves expensive, so bounds above 64 keep the batched
//! sort and trim path until normalized keys make a heap cheap.
//!
//! # Rejecting a chunk in one pass instead of a row at a time
//!
//! Almost every row of a million loses to the ten already held, and the cheap way to find that out
//! is the comparison kernel rather than a loop. Once the candidates are full, the worst of them is
//! a constant for the length of a chunk, so one vectorized comparison of the first key column
//! against that constant says which rows are still worth looking at, and on `ORDER BY EventTime
//! LIMIT 10` over ClickBench that is almost none of them after the first chunk. The rows it hands
//! back go down the same row at a time path as before, which is what keeps the answer the same.
//!
//! The bound moves while the chunk is being walked, since a winner replaces the worst candidate, so
//! what the pass produces is a superset of the rows that really win. That is the point: it is a
//! filter and not the decision, and every row it keeps is compared again properly.
//!
//! The batched path gets the same pass, one step behind. It does not hold its candidates in order,
//! so it only knows its worst one just after a trim, and that key is what it rejects against until
//! the next trim. It stays a true bound in between, because the running only improves: once the
//! bound of candidates are all at least as good as that key, nothing worse than it can finish
//! inside the bound. This matters more here than it does above, not less. A row this path keeps
//! costs a `Vec` for its key and a `Vec` for the whole row with a `Value` per column, so on
//! ClickBench's four `LIMIT 10 OFFSET 1000` queries, which come out of a group by with a `URL` in
//! the key, building those for every row was most of what the operator did.
//!
//! Three things make it step aside and look at every row instead. A worst candidate whose first key
//! is null, because then what beats it depends on where the query puts nulls and the comparison
//! kernel answers null rather than true. A `NULLS FIRST` ordering over a column that has nulls, for
//! the same reason read the other way: those rows win and a comparison against a value says nothing
//! about them. And a comparison the kernel refuses, which is left to the row path so that the error
//! comes out of the same place it came out of before. With more than one sort key the pass keeps
//! ties on the first one, since the second key can still turn a tie into a win.
//!
//! # The shape a sink has
//!
//! Like the sort, this is a [`Sink`]. The difference between them is where the trimming happens: an
//! instance trims what it holds as it goes, and `combine` trims again over what two instances
//! brought, which is what keeps the bound the bound rather than the bound times the number of
//! threads. On one thread there is one instance and the second trim does nothing.
//!
//! Every candidate carries where it arrived, for the reason the sort's own documentation gives: a
//! tie on every key is settled by that rather than by which thread got there first, so the ten rows
//! this hands back are the ten a single thread would have handed back. It costs sixteen bytes per
//! candidate and the bound is ten.

use std::cmp::Ordering;
use std::sync::{Arc, Mutex};

use rudb_common::bounds::Bound;
use rudb_common::{Error, LogicalType, Memory, Reservation, Result, Session, Value};
use rudb_kernels::{Comparison, compare as compare_vectors, selection};
use rudb_pipeline::{Lease, Progress, Sink};
use rudb_plan::{Plan, Slice, SortKey};
use rudb_vector::{Chunk, Selection, Vector};

use crate::buffer::Buffered;
use crate::cutoff::Cutoff;
use crate::prepared::{Prepared, Scratch};
use crate::rows;
use crate::schema::Schema;
use crate::sort::{Place, compare, rank, settled};

/// One row in the running: the values of its keys, the row itself, and where it arrived.
type Sortable = crate::sort::Sortable;

/// Above this bound, moving a sorted candidate array costs more than trimming in batches.
const SORTED_BOUND: usize = 64;

/// The first rows of an ordering, without holding the rest.
#[derive(Debug)]
pub(crate) struct TopN {
    keys: Vec<SortKey>,
    /// The key expressions, evaluated against the input's schema.
    exprs: Prepared,
    /// The input's types, which are also the output's.
    types: Vec<LogicalType>,
    /// How many rows to emit, once the ones to skip have been skipped.
    count: usize,
    /// How many rows to skip first.
    offset: usize,
    /// `count + offset`, which is how many rows can still turn out to be wanted.
    bound: usize,
    memory: Memory,
    /// What every instance brought, already trimmed to the bound.
    rows: Mutex<Vec<Sortable>>,
    /// What those rows are charged, given back once the finished chunks are charged instead.
    charged: Mutex<Vec<Reservation>>,
    /// What the finished chunks are charged, held for as long as they are readable.
    held: Mutex<Reservation>,
    /// How good a row has to be to still be wanted, told to the scan below as this fills up.
    ///
    /// Empty for a top N with nothing under it that could use one, which is every shape but a scan
    /// under a filter under a projection. See [`crate::cutoff`].
    cutoff: Option<Arc<Cutoff>>,
    out: Buffered,
}

/// What one instance of a top N holds while it runs.
#[derive(Debug)]
pub(crate) struct Running {
    kept: Vec<Sortable>,
    scratch: Scratch,
    charged: Reservation,
    failure: Option<Error>,
    /// The morsel this instance is reading and how many of its rows have arrived.
    place: Place,
    /// The key of the worst candidate this instance has room for, once it has the bound of them.
    ///
    /// `None` until the first trim fills the running, and set at every trim after that. Only the
    /// batched path keeps it, since the sorted path reads the same key straight out of `kept`, which
    /// it holds in order at all times.
    ///
    /// It stays true between trims because the running only improves. Once the bound of candidates
    /// are all at least as good as this key, a row worse than it cannot finish inside the bound, and
    /// neither can one that ties it, because everything it ties arrived first.
    cut: Option<Vec<Value>>,
}

impl TopN {
    /// Applies the session semantics to the prepared sort keys.
    #[must_use]
    pub(crate) fn in_session(mut self, session: &Session) -> Self {
        self.exprs = self.exprs.in_session(session);
        self
    }

    /// # Errors
    ///
    /// If a sort key does not resolve against the input's schema.
    pub(crate) fn new(
        plan: &Plan,
        input: &Schema,
        keys: Slice,
        count: u64,
        offset: u64,
        memory: &Memory,
    ) -> Result<(Self, Buffered)> {
        // A limit past what a `Vec` can hold is a limit nothing reaches, so saturating here turns a
        // bound nobody can hit into the largest one this machine has room for, and the operator
        // degenerates into the sort it would have been.
        let count = usize::try_from(count).unwrap_or(usize::MAX);
        let offset = usize::try_from(offset).unwrap_or(usize::MAX);
        let keys = plan.sort_key_list(keys).to_vec();
        let exprs: Vec<_> = keys.iter().map(|key| key.expr).collect();
        let out = Buffered::new();
        let top = Self {
            exprs: Prepared::new(plan, &exprs, input)?,
            keys,
            types: input.types(),
            count,
            offset,
            bound: count.saturating_add(offset),
            memory: memory.clone(),
            rows: Mutex::new(Vec::new()),
            charged: Mutex::new(Vec::new()),
            held: Mutex::new(memory.reservation()),
            cutoff: None,
            out: out.clone(),
        };
        Ok((top, out))
    }

    /// Tells this to publish its worst candidate to `cutoff` as it goes.
    ///
    /// Taken whether the cutoff was armed or not, because the builder makes one before it walks into
    /// the input and only finds out afterwards whether the walk reached a scan. An unarmed one costs
    /// a load and a branch per chunk and is never read by anybody.
    #[must_use]
    pub(crate) fn telling(mut self, cutoff: Arc<Cutoff>) -> Self {
        self.cutoff = Some(cutoff);
        self
    }

    /// Says how good a row now has to be, given the worst of a full set of candidates.
    ///
    /// Only the first key, because a part of the file whose first key is all worse than this is worse
    /// whatever its later keys hold, and one that ties on the first key says nothing. A null first
    /// key publishes nothing, which under the `NULLS LAST` this is only ever armed for means the
    /// candidates do not yet rule out any value at all.
    ///
    /// The arming is checked before the key is turned into a bound rather than after, because a
    /// string key would be copied to find out that nobody was listening.
    fn reached(&self, worst: &[Value]) {
        let Some(cutoff) = self.cutoff.as_ref().filter(|cutoff| cutoff.armed()) else { return };
        let Some(bound) = worst.first().and_then(Bound::of_value) else { return };
        cutoff.reached(bound);
    }
}

impl Sink for TopN {
    type Local = Running;

    fn local(&self) -> Running {
        Running {
            kept: Vec::new(),
            scratch: self.exprs.scratch(),
            charged: self.memory.reservation(),
            failure: None,
            place: Place::default(),
            cut: None,
        }
    }

    fn at(&self, morsel: &rudb_pipeline::Morsel, local: &mut Running) -> Result<()> {
        local.place.start(morsel.index());
        Ok(())
    }

    fn sink(&self, chunk: &Chunk, local: &mut Running) -> Result<Progress> {
        let mut keys = Vec::with_capacity(self.keys.len());
        self.exprs.evaluate(chunk, &mut local.scratch, &mut keys)?;
        if self.bound <= SORTED_BOUND {
            let full = self.bound > 0 && local.kept.len() == self.bound;
            let narrowed = full
                .then(|| {
                    worth_looking_at(&self.keys, &keys, &local.kept[self.bound - 1].0, chunk.len())
                })
                .flatten();
            let failure = &mut local.failure;
            match narrowed {
                // row at a time: the rows the pass kept are the ones that can still win, and each
                // of them has to be placed among the candidates rather than counted.
                Some(rows) => {
                    for row in rows.iter() {
                        let arrival = local.place.of(row);
                        keep(
                            Where { keys: &self.keys, columns: &keys, chunk, row, arrival },
                            &mut local.kept,
                            self.bound,
                            failure,
                        );
                    }
                }
                // row at a time: the key still has the same Value layout the sort holds, and 2i
                // (#63) replaces it with one normalized comparable byte string per row.
                None => {
                    for row in 0..chunk.len() {
                        let arrival = local.place.of(row);
                        keep(
                            Where { keys: &self.keys, columns: &keys, chunk, row, arrival },
                            &mut local.kept,
                            self.bound,
                            failure,
                        );
                    }
                }
            }
            local.place.past(chunk.len());
            recharge(&local.kept, &mut local.charged)?;
            if self.bound > 0 && local.kept.len() == self.bound {
                self.reached(&local.kept[self.bound - 1].0);
            }
            return Ok(Progress::More);
        }
        // The same question the sorted path asks, asked once the running is full rather than on
        // every chunk from the start, because this path only knows its worst candidate after a trim.
        // Without it a `LIMIT 10 OFFSET 1000` builds a `Vec<Value>` for the key and another for
        // every column of every row that reaches it, however hopeless the row is, and on a wide row
        // with a string in it that is most of what the operator does.
        let narrowed = local
            .cut
            .as_ref()
            .and_then(|cut| worth_looking_at(&self.keys, &keys, cut, chunk.len()));
        let mut taken = 0;
        match narrowed {
            // row at a time: the rows the pass kept are the ones that can still win, and each of
            // them has to be read out of the columns rather than counted.
            Some(rows) => {
                for row in rows.iter() {
                    let arrival = local.place.of(row);
                    taken += hold(&keys, chunk, row, arrival, &mut local.kept)?;
                }
            }
            // row at a time: the key still has the same Value layout the sort holds, and 2i (#63)
            // replaces it with one normalized comparable byte string per row.
            None => {
                for row in 0..chunk.len() {
                    let arrival = local.place.of(row);
                    taken += hold(&keys, chunk, row, arrival, &mut local.kept)?;
                }
            }
        }
        local.place.past(chunk.len());
        local.charged.grow(taken)?;
        if local.kept.len() > self.bound.saturating_mul(2) {
            trim(&self.keys, &mut local.kept, self.bound, &mut local.failure);
            recharge(&local.kept, &mut local.charged)?;
            local.cut = (local.kept.len() == self.bound && self.bound > 0)
                .then(|| local.kept[self.bound - 1].0.clone());
            if let Some(cut) = local.cut.as_ref() {
                self.reached(cut);
            }
        }
        Ok(Progress::More)
    }

    fn combine(&self, mut local: Running) -> Result<()> {
        if let Some(error) = local.failure {
            return Err(error);
        }
        let mut rows = self.rows.lock().map_err(poisoned)?;
        rows.extend(local.kept);
        // Trimmed here as well as in the instance, so that combining thirty two instances holding
        // the bound each leaves the bound and not thirty two times it.
        let mut failure = None;
        trim(&self.keys, &mut rows, self.bound, &mut failure);
        if let Some(error) = failure {
            return Err(error);
        }
        recharge(&rows, &mut local.charged)?;
        self.charged.lock().map_err(poisoned)?.push(local.charged);
        Ok(())
    }

    fn finalize(&self, _threads: &Lease<'_>) -> Result<()> {
        let kept = std::mem::take(&mut *self.rows.lock().map_err(poisoned)?);
        let wanted = kept.into_iter().skip(self.offset).take(self.count);
        let ordered: Vec<Vec<Value>> = wanted.map(|(_, row, _)| row).collect();
        let mut held = self.held.lock().map_err(poisoned)?;
        let chunks = rows::chunks(&self.types, &ordered, &mut held)?;
        self.out.fill(chunks)?;
        self.charged.lock().map_err(poisoned)?.clear();
        Ok(())
    }
}

fn poisoned<T>(_: T) -> Error {
    Error::internal("a thread panicked while holding the rows a top N is keeping")
}

/// Orders what is held and keeps the first `bound` of it.
///
/// Ties are settled by where the rows arrived rather than by the order they were handed over, which
/// is what makes the trim in `combine` give the same answer whichever thread combined first.
fn trim(keys: &[SortKey], kept: &mut Vec<Sortable>, bound: usize, failure: &mut Option<Error>) {
    kept.sort_by(|left, right| settled(keys, left, right, failure));
    kept.truncate(bound);
}

/// Reads one row out of the columns and puts it among the candidates, unordered.
///
/// What the batched path does with a row it has decided to keep, and the reason it costs what it
/// costs: a `Vec` for the key, a `Vec` for the row, and a `Value` per column of each, which for a
/// string column is a `String`. The answer is what those two together are charged.
fn hold(
    keys: &[Vector],
    chunk: &Chunk,
    row: usize,
    arrival: crate::sort::Arrival,
    kept: &mut Vec<Sortable>,
) -> Result<u64> {
    let key: Vec<Value> =
        keys.iter().map(|column| column.try_value_at(row)).collect::<Result<_>>()?;
    let values: Vec<Value> =
        (0..chunk.width()).map(|column| chunk.try_value_at(row, column)).collect::<Result<_>>()?;
    let taken = rows::footprint(&key) + rows::footprint(&values);
    kept.push((key, values, arrival));
    Ok(taken)
}

/// Keeps one row when its key belongs in the ordered prefix.
///
/// The key is read out of the columns a value at a time and only as far as the first key that
/// separates it from the worst candidate, so a row that loses on the first of three keys costs one
/// value rather than three and never allocates the `Vec` that holds them. Almost every row loses.
fn keep(
    Where { keys, columns, chunk, row, arrival }: Where<'_>,
    kept: &mut Vec<Sortable>,
    bound: usize,
    failure: &mut Option<Error>,
) {
    if bound == 0 {
        return;
    }
    if kept.len() == bound
        && against(keys, columns, row, &kept[bound - 1].0, failure) != Ordering::Less
    {
        return;
    }
    let key: Vec<Value> = match columns.iter().map(|column| column.try_value_at(row)).collect() {
        Ok(key) => key,
        Err(error) => {
            failure.get_or_insert(error);
            return;
        }
    };
    let values: Vec<Value> =
        match (0..chunk.width()).map(|column| chunk.try_value_at(row, column)).collect() {
            Ok(values) => values,
            Err(error) => {
                failure.get_or_insert(error);
                return;
            }
        };
    // After every candidate whose key it ties, which is where its arrival puts it too: an instance
    // reads the morsels it is given in order and each of them from the start, so a row reaching
    // here arrived after everything already held.
    let at = kept.partition_point(|candidate| {
        compare(keys, &candidate.0, &key, failure) != Ordering::Greater
    });
    kept.insert(at, (key, values, arrival));
    kept.truncate(bound);
}

/// One row being offered to the candidates, which is five things that only travel together.
struct Where<'a> {
    keys: &'a [SortKey],
    columns: &'a [Vector],
    chunk: &'a Chunk,
    row: usize,
    arrival: crate::sort::Arrival,
}

/// Where one row of the key columns sits against a key already held.
///
/// The same answer [`compare`] gives for the same two keys, read straight out of the columns rather
/// than out of a `Vec` built for the purpose.
fn against(
    keys: &[SortKey],
    columns: &[Vector],
    row: usize,
    held: &[Value],
    failure: &mut Option<Error>,
) -> Ordering {
    for (at, key) in keys.iter().enumerate() {
        let value = match columns[at].try_value_at(row) {
            Ok(value) => value,
            Err(error) => {
                failure.get_or_insert(error);
                return Ordering::Equal;
            }
        };
        let ordering = match rank(&value, &held[at], *key) {
            Ok(ordering) => ordering,
            Err(error) => {
                failure.get_or_insert(error);
                Ordering::Equal
            }
        };
        if ordering != Ordering::Equal {
            return ordering;
        }
    }
    Ordering::Equal
}

/// The rows of a chunk that can still beat `worst`, or nothing when every row has to be looked at.
///
/// One comparison of the first key column against a constant, for the reasons in the module doc.
fn worth_looking_at(
    keys: &[SortKey],
    columns: &[Vector],
    worst: &[Value],
    rows: usize,
) -> Option<Selection> {
    let key = *keys.first()?;
    let bound = worst.first()?;
    let column = columns.first()?;
    if bound.is_null() || (key.nulls_first && column.validity().has_nulls(rows)) {
        return None;
    }
    let op = match (key.descending, keys.len() == 1) {
        (false, true) => Comparison::Less,
        (false, false) => Comparison::LessOrEqual,
        (true, true) => Comparison::Greater,
        (true, false) => Comparison::GreaterOrEqual,
    };
    let against = Vector::constant(column.logical_type().clone(), bound.clone(), rows);
    // The whole chunk is in play here, so this asks the kernel that reads its operands where they
    // lie rather than the one that reads them through a selection. Handing the threaded kernel an
    // identity selection would build a `u32` a row to say "all of them", check every one of them is
    // in range, and then put a load and an indirection in front of each of the comparisons that
    // loop is otherwise three instructions long. On `ORDER BY EventTime LIMIT 10` over ClickBench
    // that was more instructions than decoding the column cost.
    let flags = compare_vectors(op, column, &against).ok()?;
    Some(selection(&flags, rows))
}

/// Charges the scratch reservation for what is still held after a trim.
///
/// Released and taken again rather than shrunk, because a reservation gives everything back at once
/// and has no partial release. Nothing else can be holding the difference at this point, since the
/// operator is between two reads of its input.
fn recharge(kept: &[Sortable], scratch: &mut Reservation) -> Result<()> {
    let footprint =
        kept.iter().map(|(key, values, _)| rows::footprint(key) + rows::footprint(values)).sum();
    scratch.release();
    scratch.grow(footprint)
}