rudb_native/stats.rs
1//! Building a table's statistics sections from the table's own columns.
2//!
3//! The same meeting place `graph` is, for the other document. `rudb-stats` at rank 5 knows what a
4//! column summary says and knows nothing about a file; the rest of this crate knows how to put an
5//! opaque payload in a file and nothing about what one means. Building a summary for a real table
6//! means reading the column back, so it happens here, in the crate allowed to see both.
7//!
8//! Everything here obeys `spec/stats/03-the-file-format.md` section 3.1, which is the graph
9//! document's section 3.1 applied to a second kind of payload: delete every statistics section and
10//! no query changes its answer, only the time. That is why [`summary`] and [`sketches`] answer with
11//! an [`Option`] and not a [`Result`]. There is no failure they could report that is not answered
12//! by planning the query the way it was planned before the section existed.
13//!
14//! # The invariant has teeth here that it does not have in the graph layer
15//!
16//! A key map can only make a join faster. A summary can answer a query: a `COUNT(DISTINCT c)` comes
17//! out of one without the column being touched. So the thing that has to survive is not only *is
18//! the section there* but *is the number in it exact*, and [`Summary::distinct_class`] is where that
19//! lives. This module's job is to never write [`Class::Exact`] onto a number that is not, which in
20//! practice means one rule: the sketch says whether it overflowed, and everything else follows from
21//! that answer rather than from what the writer hoped.
22//!
23//! # One pass, and what that costs
24//!
25//! Section 3.7 gives the statistics build ten percent of the native write time, and the way to stay
26//! inside it is not to be clever but to read the column once. [`build_summary`] takes one scan and
27//! computes every field of the summary and the sketch from it, so the cost of statistics on a write
28//! is the cost of one more read of each column asked for, and no column is read twice.
29//!
30//! # The per stripe rule
31//!
32//! Section 3.8 says per stripe structures are written only for the columns that get read, and the
33//! arithmetic behind that is not close: sixteen `lineitem` columns at SF100, sketched per stripe
34//! even at the small k a stripe sketch keeps, come to several hundred megabytes against a budget of
35//! two percent. So the default is a merged sketch and nothing else, and [`Sketches::stripes`] being
36//! empty is the state the rule says most columns are in rather than a degraded one.
37//!
38//! [`read_columns`] is what this build promotes a column with. It reads the promoted set off the
39//! file, which today means the columns that already carry a key map or a forward link, because
40//! those are the columns something has declared a relationship or a key over and section 3.8 names
41//! them directly. Document 06's observation log is the other source the spec names and it is not
42//! built yet, so when it arrives it adds columns to this list and changes nothing else here.
43//!
44//! Promotion costs no extra hashing. The column is read once and hashed once either way, and what
45//! changes is where the counting is reset. [`build_summary_for`] has the argument.
46
47use std::cmp::Ordering;
48use std::path::Path;
49use std::sync::Arc;
50use std::time::{Duration, Instant};
51
52use rudb_common::bounds::{self, Bound};
53use rudb_common::stat::Class;
54use rudb_common::{LogicalType, Result, Value};
55use rudb_encoding::sketch::{DEFAULT_K, Sketch};
56use rudb_stats::{Order, STRIPE_K, Sketches, Summary, sketches::HEADER_BYTES as SKETCH_HEADER};
57use rudb_storage::count::{Counts, countable};
58use rudb_vector::{Data, Form, StringColumn, Validity, Vector};
59
60use crate::section::{self, Attachment};
61use crate::{Catalog, Reader, invalid};
62
63/// The share of a table's stored column bytes its statistics sections are allowed to cost together.
64///
65/// Two percent, per section 3.8, and kept apart from the graph layer's ten percent rather than
66/// pooled with it. Two budgets that share a pot are two budgets where the one that runs first wins,
67/// and a table whose key maps happened to be built before its summaries would then have no
68/// summaries for a reason that has nothing to do with summaries. They are counted separately for the
69/// same reason they are two documents.
70pub const BUDGET_SHARE: u64 = 2;
71
72/// The size below which a table's statistics sections always fit, whatever the share works out to.
73///
74/// The same floor and the same argument as `graph::BUDGET_FLOOR`. A summary is a few hundred bytes
75/// on a table of any size and two percent of a small, well compressed column is less than that, so
76/// the pure rule would throw away the cheapest structure in the system for being expensive.
77pub const BUDGET_FLOOR: u64 = 64 * 1024;
78
79/// What one column's statistics cost and what they say.
80#[derive(Debug, Clone)]
81pub struct Built {
82 /// Which column was summarized.
83 pub column: usize,
84 /// Rows in the column, nulls included.
85 pub rows: u64,
86 /// Distinct non-null values, as the summary reports them.
87 pub distinct: u64,
88 /// Whether that distinct count is exact rather than a sketch estimate.
89 pub exact: bool,
90 /// Which way the values run.
91 pub order: Order,
92 /// What the summary section takes in the file.
93 pub summary_bytes: usize,
94 /// What the sketches section takes in the file.
95 pub sketch_bytes: usize,
96 /// How many per stripe sketches went in it, which is zero for a column the per stripe rule did
97 /// not promote and is most of them.
98 pub stripes: usize,
99 /// What the column takes in the file, which is what the budget is a share of.
100 pub column_bytes: u64,
101 /// Whether the sections were kept. False means they were built, measured, and found to cost more
102 /// than section 3.8 allows, so the file does not have them and every query plans as though
103 /// statistics had never been implemented.
104 pub built: bool,
105 /// How long the build took, the reading of the column included.
106 pub build: Duration,
107}
108
109impl Built {
110 /// Both sections together, which is what the budget spends.
111 #[must_use]
112 pub fn bytes(&self) -> usize {
113 self.summary_bytes + self.sketch_bytes
114 }
115}
116
117/// A column's summary and its sketches, which are built together because they are one pass.
118#[derive(Debug, Clone)]
119pub struct Stats {
120 /// What the column says about itself.
121 pub summary: Summary,
122 /// The sketch the distinct count came out of.
123 pub sketches: Sketches,
124}
125
126/// Builds the summary and the sketches for one column of a committed table.
127///
128/// # Errors
129///
130/// If the column cannot be read, is past the end of the table, or is of a type with no hash rule.
131/// The last one is refused by name rather than approximated: the types without a rule are the
132/// interval and the nested ones, a summary of one would carry a distinct count of zero that nothing
133/// could tell from a column of nulls, and none of TPC-H or ClickBench has one.
134pub fn build_summary(reader: &Reader, column: usize) -> Result<Stats> {
135 build_summary_for(reader, column, false)
136}
137
138/// The same, keeping a sketch per stripe as well as the merged one when `per_stripe` is set.
139///
140/// Whether to set it is section 3.8's rule and not a caller's taste: per stripe structures are
141/// written only for the columns that get read, because sixteen `lineitem` columns at SF100 come to
142/// several hundred megabytes of them against a budget of two percent. [`read_columns`] is what this
143/// build answers that question with.
144///
145/// The extra sketches cost no extra hashing. Each stripe is counted into its own [`Counts`] at the
146/// column's k, the merged sketch is the union of those, which is exact because they are all at the
147/// same k, and each one is written down at [`rudb_stats::STRIPE_K`] through [`Sketch::narrowed`],
148/// which is exact because a bottom-k of a bottom-k is a bottom-k. So the column is read once and
149/// hashed once either way, and the difference between a promoted column and an ordinary one is
150/// where the counting is reset and how much of it is written.
151///
152/// # Errors
153///
154/// If the column cannot be read, is past the end of the table, or is of a type with no hash rule.
155/// The last one is refused by name rather than approximated: the types without a rule are the
156/// interval and the nested ones, a summary of one would carry a distinct count of zero that nothing
157/// could tell from a column of nulls, and none of TPC-H or ClickBench has one.
158pub fn build_summary_for(reader: &Reader, column: usize, per_stripe: bool) -> Result<Stats> {
159 let fields = reader.table().fields();
160 let Some(field) = fields.get(column) else {
161 return Err(invalid(&format!(
162 "column {column} is past the {} of table {}",
163 fields.len(),
164 reader.table().name()
165 )));
166 };
167 if !countable(&field.ty) {
168 return Err(invalid(&format!(
169 "a summary of {} needs a hash rule, and {} has none",
170 field.name, field.ty
171 )));
172 }
173 let blind = || {
174 // A blind column: a form `rudb_storage::count` has no arm for turned up, so its sketch is
175 // missing rows and says nothing about which. A distinct count that is too low is the one
176 // error an estimator has no defence against, so the column gets no summary at all rather
177 // than a summary with a number in it nothing can check.
178 invalid(&format!(
179 "column {} of {} holds a form with no hash rule, so it has no sketch",
180 field.name,
181 reader.table().name()
182 ))
183 };
184
185 let mut whole = Counts::new(1);
186 let mut stripes = Vec::new();
187 let mut pass = Pass::new(&field.ty, reader.table().generation());
188 for (at, stripe) in reader.stripe_parts().into_iter().enumerate() {
189 pass.open_stripe((at as u64, 0));
190 let mut counted = per_stripe.then(|| Counts::new(1));
191 for part in stripe {
192 let chunk = reader.read(part, &[column])?;
193 match counted.as_mut() {
194 Some(counted) => counted.add(&chunk),
195 None => whole.add(&chunk),
196 }
197 pass.scan(chunk.column(0)?);
198 }
199 pass.close_stripe();
200 if let Some(counted) = counted {
201 stripes.push(counted.sketch(0).ok_or_else(blind)?);
202 }
203 }
204 if !per_stripe {
205 return Ok(pass.finish(whole.sketch(0).ok_or_else(blind)?, Vec::new()));
206 }
207 let mut merged = Sketch::new(DEFAULT_K)?;
208 for stripe in &stripes {
209 merged = merged.union(stripe)?;
210 }
211 let narrowed =
212 stripes.iter().map(|stripe| stripe.narrowed(STRIPE_K)).collect::<Result<Vec<_>>>()?;
213 Ok(pass.finish(merged, narrowed))
214}
215
216/// What one stripe says about the order of its rows, kept until every stripe is in.
217#[derive(Debug)]
218struct Piece {
219 key: (u64, u64),
220 first: Option<Bound>,
221 last: Option<Bound>,
222 ascending: bool,
223 descending: bool,
224 runs: u64,
225}
226
227/// One scan of one column, in `rid` order, for everything the sketch does not answer.
228///
229/// In `rid` order because the order fields depend on it. A pass that read the parts in any other
230/// order would report a column as unordered that is sorted, which costs a plan and not an answer,
231/// and would report the run count of a shuffle, which is worse because it is a number rather than a
232/// flag and looks like it was measured.
233///
234/// The distinct count is not here. That is `rudb_storage::count::Counts`, which walks a vector by
235/// its form rather than a row at a time and which a column of a million runs costs one hash. Doing
236/// it twice would double the expensive half of the build and the budget is ten percent of the write.
237#[derive(Debug)]
238struct Pass {
239 rows: u64,
240 nulls: u64,
241 low: Option<Bound>,
242 high: Option<Bound>,
243 /// False once a non-null value turned up that has no ordered bound, which makes both ends
244 /// unusable rather than merely absent.
245 bounded: bool,
246 ascending: bool,
247 descending: bool,
248 runs: u64,
249 previous: Option<Bound>,
250 bytes: u64,
251 widest: u64,
252 generation: u64,
253 /// The two ends of the stripe being read, kept apart rather than as a pair so that each one can
254 /// be compared against and refilled on its own. A pair would have to be taken out and put back
255 /// whole, which is the move that made this pass allocate.
256 stripe_low: Option<Bound>,
257 stripe_high: Option<Bound>,
258 stripes: Vec<(Bound, Bound)>,
259 /// Where the stripe being read sits in the table, and the first value it held.
260 ///
261 /// A writer fed by several pipeline instances gets its stripes in the order they finished
262 /// rather than the order they sit in, and sorts them by this key when it commits. The order
263 /// fields are about adjacent rows, so they are read a stripe at a time into [`Piece`]s and put
264 /// together in key order at the end, which is the rid order the reader will see.
265 key: (u64, u64),
266 first: Option<Bound>,
267 pieces: Vec<Piece>,
268 /// What one value of this column takes, when every value takes the same.
269 ///
270 /// Read off the type once rather than off each value, because for every fixed width column it is
271 /// a constant and asking a value for it is a branch a hundred million times to hear the same
272 /// number. `None` is a variable width type and those are measured per value.
273 fixed: Option<u64>,
274 /// The scale of a decimal column, so that an integer read out of a vector becomes the bound the
275 /// column's other writers would have written for the same value.
276 scale: Option<u8>,
277 /// The last dictionary this pass read, so that a column whose vectors share one reads it once.
278 coded: Option<Coded>,
279}
280
281impl Pass {
282 fn new(ty: &LogicalType, generation: u64) -> Self {
283 Self {
284 rows: 0,
285 nulls: 0,
286 low: None,
287 high: None,
288 bounded: true,
289 ascending: true,
290 descending: true,
291 runs: 0,
292 previous: None,
293 bytes: 0,
294 widest: 0,
295 generation,
296 stripe_low: None,
297 stripe_high: None,
298 stripes: Vec::new(),
299 key: (0, 0),
300 first: None,
301 pieces: Vec::new(),
302 fixed: fixed_width(ty),
303 scale: bounds::scale_of(ty),
304 coded: None,
305 }
306 }
307
308 /// One vector of the column, a vector at a time where the layout allows it and a row at a time
309 /// where it does not.
310 fn scan(&mut self, vector: &Vector) {
311 if self.scan_flat(vector) || self.scan_gathered(vector) || self.scan_dictionary(vector) {
312 return;
313 }
314 self.scan_rows(vector);
315 }
316
317 /// One vector of a flat signed column, with the layout matched on once instead of once a row.
318 ///
319 /// `false` if the vector is not one of those, and the caller falls back to [`Self::scan_rows`].
320 ///
321 /// This is where the build's time went. [`Self::scan_rows`] asks `Vector::signed_at` for every
322 /// row, and that is a validity test, a match over the body forms and a second match over the
323 /// dozen layouts, and then the answer is wrapped in a [`Bound`] and compared through
324 /// [`Bound::order`], which is another match, four times. Measured on TPC-H SF1 that came to
325 /// about 355 instructions for a value whose whole job is three comparisons: 53.4 G instructions
326 /// of the 60.8 G the statistics added to the write, against 7.9 G for the sketch that hashes
327 /// every one of the same values. The sketch was never the expensive half.
328 ///
329 /// Matched once, the loop underneath is a validity bit and three integer compares. The layouts
330 /// are the signed group and not the unsigned one, because `Vector::signed_at` reads the signed
331 /// group and this has to agree with the path it is replacing rather than be better than it.
332 fn scan_flat(&mut self, vector: &Vector) -> bool {
333 if vector.form() != Form::Flat {
334 return false;
335 }
336 let Some(data) = vector.data() else { return false };
337 let rows = vector.len();
338 let validity = vector.validity();
339 macro_rules! signed {
340 ($(($variant:ident, $native:ty, $zero:expr)),+ $(,)?) => {
341 match data {
342 $(Data::$variant(held) => {
343 let held: &[$native] = held;
344 if held.len() < rows {
345 return false;
346 }
347 let spread = spread(rows, validity, |row| i128::from(held[row]));
348 self.fold_spread(&spread);
349 return true;
350 })+
351 Data::Float64(held) => self.scan_reals(rows, validity, held),
352 Data::Float32(held) => self.scan_reals(rows, validity, held),
353 Data::Varlen(held) => self.scan_strings(rows, validity, held),
354 _ => false,
355 }
356 };
357 }
358 rudb_vector::for_each_layout!(signed, signed)
359 }
360
361 /// One vector of a flat float column, which [`Self::scan_rows`] read by building a `Value` a row
362 /// and comparing it as a [`Bound`] four times.
363 ///
364 /// On a `lineitem` load from CSV the four price and quantity columns come in as `DOUBLE`, and
365 /// that was about 5% of the load's cycles. `false` when the vector holds a NaN, and the caller
366 /// reads it a row at a time as before.
367 fn scan_reals<T: Copy + Into<f64>>(
368 &mut self,
369 rows: usize,
370 validity: &Validity,
371 held: &[T],
372 ) -> bool {
373 if held.len() < rows {
374 return false;
375 }
376 let Some(spread) = real_spread(rows, validity, |row| held[row].into()) else {
377 return false;
378 };
379 let width = self.fixed.unwrap_or(8);
380 self.fold(Reduced {
381 rows: spread.rows,
382 nulls: spread.nulls,
383 values: spread.values,
384 bytes: width.saturating_mul(spread.values),
385 widest: if spread.values > 0 { width } else { 0 },
386 ascents: spread.ascents,
387 descents: spread.descents,
388 ends: (spread.values > 0).then_some(Ends {
389 low: Bound::Real(spread.low),
390 high: Bound::Real(spread.high),
391 first: Bound::Real(spread.first),
392 last: Bound::Real(spread.last),
393 }),
394 });
395 true
396 }
397
398 /// One vector of a flat string column, compared against itself and then folded in once.
399 ///
400 /// [`Self::bytes_value`] compares every row with the row before it and with all four ends, and
401 /// copies it in as the previous row, which on a `lineitem` load from CSV was about 3% of the
402 /// load's cycles in `memcmp`. Within a vector the row before is still a borrow, so nothing is
403 /// copied, and a row is only compared with the end it can move: one above the row before it
404 /// cannot be the lowest yet, and one below cannot be the highest. The vector's own ends go
405 /// through [`Self::fold`] once, which is the only place they are copied.
406 fn scan_strings(&mut self, rows: usize, validity: &Validity, held: &StringColumn) -> bool {
407 if held.len() < rows {
408 return false;
409 }
410 let nullable = validity.has_nulls(rows);
411 let mut out = Reduced::empty(rows as u64);
412 let mut ends: Option<(&[u8], &[u8], &[u8])> = None;
413 let mut last: &[u8] = &[];
414 // row at a time: the ascents, the descents and which end a row can move are all about the
415 // row before it. The bytes are borrowed out of the column and no `Value` is built.
416 for row in 0..rows {
417 if nullable && !validity.is_valid(row) {
418 out.nulls += 1;
419 continue;
420 }
421 let Some(bytes) = held.bytes(row) else { return false };
422 let width = bytes.len() as u64;
423 out.bytes = out.bytes.saturating_add(width);
424 out.widest = out.widest.max(width);
425 match &mut ends {
426 None => ends = Some((bytes, bytes, bytes)),
427 Some((low, high, _)) => match bytes.cmp(last) {
428 Ordering::Greater => {
429 out.ascents += 1;
430 if bytes > *high {
431 *high = bytes;
432 }
433 }
434 Ordering::Less => {
435 out.descents += 1;
436 if bytes < *low {
437 *low = bytes;
438 }
439 }
440 Ordering::Equal => {}
441 },
442 }
443 last = bytes;
444 out.values += 1;
445 }
446 out.ends = ends.map(|(low, high, first)| Ends {
447 low: Bound::Bytes(low.to_vec()),
448 high: Bound::Bytes(high.to_vec()),
449 first: Bound::Bytes(first.to_vec()),
450 last: Bound::Bytes(last.to_vec()),
451 });
452 self.fold(out);
453 true
454 }
455
456 /// What [`spread`] made of one vector of a signed column, folded into the pass.
457 fn fold_spread(&mut self, spread: &Spread) {
458 let width = self.fixed.unwrap_or(8);
459 self.fold(Reduced {
460 rows: spread.rows,
461 nulls: spread.nulls,
462 values: spread.values,
463 bytes: width.saturating_mul(spread.values),
464 widest: if spread.values > 0 { width } else { 0 },
465 ascents: spread.ascents,
466 descents: spread.descents,
467 ends: (spread.values > 0).then(|| Ends {
468 low: self.bound(spread.low),
469 high: self.bound(spread.high),
470 first: self.bound(spread.first),
471 last: self.bound(spread.last),
472 }),
473 });
474 }
475
476 /// One vector of a signed column coded against a flat dictionary, read through its codes.
477 ///
478 /// `false` if the vector is not one of those, or if its dictionary holds a null, and the caller
479 /// tries [`Self::scan_dictionary`] and then [`Self::scan_rows`].
480 ///
481 /// A Parquet column chunk arrives with one dictionary for the whole row group, which is tens of
482 /// thousands of entries behind vectors of a couple of thousand rows. That is too big for
483 /// [`Self::scan_dictionary`] to order, so every one of those vectors went a row at a time, with
484 /// a code lookup, a `Bound` and four calls to [`Bound::order`] a row. On the `hits_0` load that
485 /// was about 5% of the write's cycles. A signed entry needs no ordering to be compared, so the
486 /// flat loop reads it through the code instead.
487 fn scan_gathered(&mut self, vector: &Vector) -> bool {
488 let Some((codes, values)) = vector.shared_dictionary_parts() else { return false };
489 let rows = vector.len();
490 if codes.len() < rows
491 || values.form() != Form::Flat
492 || values.validity().has_nulls(values.len())
493 {
494 return false;
495 }
496 let Some(data) = values.data() else { return false };
497 let codes = &codes[..rows];
498 let validity = vector.validity();
499 macro_rules! signed {
500 ($(($variant:ident, $native:ty, $zero:expr)),+ $(,)?) => {
501 match data {
502 $(Data::$variant(held) => {
503 let held: &[$native] = held;
504 if codes.iter().any(|&code| code as usize >= held.len()) {
505 return false;
506 }
507 let spread =
508 spread(rows, validity, |row| i128::from(held[codes[row] as usize]));
509 self.fold_spread(&spread);
510 return true;
511 })+
512 _ => false,
513 }
514 };
515 }
516 rudb_vector::for_each_layout!(signed, signed)
517 }
518
519 /// One vector of a dictionary column, with the values compared once each instead of once a row.
520 ///
521 /// `false` if the vector is not one, or if the dictionary is too big for this to be worth it, or
522 /// if its entries turn out not to be orderable against each other.
523 ///
524 /// A dictionary vector is where the rest of the build's time went, and it is most of what a load
525 /// hands the writer: on a TPC-H SF1 `lineitem` about seven vectors in ten arrive dictionary
526 /// coded, the five string columns among them. Reading one a row at a time costs a code lookup
527 /// and then all the work the flat path was doing, and for a string column it costs a byte
528 /// comparison against the value before it, for a column whose whole point is that it holds a few
529 /// dozen distinct values.
530 ///
531 /// So the dictionary is read once and then the rows are read against what it came to. See
532 /// [`Coded`] for what that is and [`Self::read_dictionary`] for how it is built.
533 fn scan_dictionary(&mut self, vector: &Vector) -> bool {
534 let Some((codes, values)) = vector.shared_dictionary_parts() else { return false };
535 let rows = vector.len();
536 if codes.len() < rows {
537 return false;
538 }
539 let held = match self.coded.take() {
540 Some(held) if Arc::ptr_eq(&held.values, values) => held,
541 // A dictionary this pass has not read. Wider than the vector it codes means reading it
542 // costs more than the rows it is about are worth, so that one goes back to the row at a
543 // time pass rather than being read at all.
544 _ => {
545 if values.len() > rows {
546 return false;
547 }
548 match self.read_dictionary(values) {
549 Some(read) => read,
550 None => return false,
551 }
552 }
553 };
554 let out = held.reduce(codes, rows, vector.validity());
555 self.coded = Some(held);
556 self.fold(out);
557 true
558 }
559
560 /// Reads a dictionary into the positions and the widths its codes stand for.
561 ///
562 /// `None` for a dictionary holding a value with no ordered bound, or a pair this build cannot
563 /// order against each other. Either way the vector goes back to [`Self::scan_rows`], which has
564 /// the rule for what a value like that does to a column's ends and is the one place it lives.
565 ///
566 /// Every entry is ordered against every other, which is one sort of a few dozen things, and then
567 /// each code carries the position its value holds in that order. Entries that order equal share
568 /// a position, so a dictionary that happens to hold one value twice says what the row at a time
569 /// pass says rather than seeing a step between the two copies of it.
570 fn read_dictionary(&self, values: &Arc<Vector>) -> Option<Coded> {
571 let mut entries = Vec::with_capacity(values.len());
572 // row at a time: these are a dictionary's entries rather than a column's rows, and there are
573 // a few dozen of them behind the thousands of rows that code against them. The third arm
574 // builds a `Value` and is the one the checker is looking for, and it runs for a float
575 // dictionary and for nothing else.
576 for at in 0..values.len() {
577 if values.is_null_at(at) {
578 entries.push(None);
579 continue;
580 }
581 // Derived the way `scan_rows` derives it, arm for arm, because the two have to agree on
582 // what bound a value has. A date read as a signed integer and a date read through
583 // `Bound::of_value` are not required to be the same bound, and a column whose vectors
584 // took different paths would be comparing one against the other.
585 let entry = match values.signed_at(at) {
586 Some(signed) => Some((self.bound(signed), self.fixed.unwrap_or(8))),
587 None => match values.bytes_at(at) {
588 Some(bytes) => Some((Bound::Bytes(bytes.to_vec()), bytes.len() as u64)),
589 None => {
590 let value = values.value_at(at);
591 let wide = self.fixed.unwrap_or_else(|| width(&value));
592 Bound::of_value(&value).map(|bound| (bound, wide))
593 }
594 },
595 };
596 entries.push(Some(entry?));
597 }
598 let mut order = (0..entries.len()).filter(|&at| entries[at].is_some()).collect::<Vec<_>>();
599 order.sort_by(|&one, &other| {
600 bound_of(&entries, one).order(bound_of(&entries, other)).unwrap_or(Ordering::Equal)
601 });
602 // The walk that hands out the positions is also what checks the sort meant anything: a pair
603 // this build cannot order sorted to wherever it happened to sit, so an unordered pair here
604 // is the whole dictionary going back to the row at a time pass.
605 let mut codes = vec![None; entries.len()];
606 let mut bounds = Vec::new();
607 for (at, &code) in order.iter().enumerate() {
608 if at > 0 {
609 match bound_of(&entries, order[at - 1]).order(bound_of(&entries, code)) {
610 Some(Ordering::Less) => bounds.push(bound_of(&entries, code).clone()),
611 Some(Ordering::Equal) => {}
612 Some(Ordering::Greater) | None => return None,
613 }
614 } else {
615 bounds.push(bound_of(&entries, code).clone());
616 }
617 let width = entries[code].as_ref().map_or(0, |(_, width)| *width);
618 codes[code] = Some(((bounds.len() - 1) as u32, width));
619 }
620 Some(Coded { values: Arc::clone(values), codes, bounds })
621 }
622
623 /// Folds what one vector came to into the pass, which is where the sequential half is settled.
624 ///
625 /// The order flags and the run count are a question about adjacent rows, so a vector at a time
626 /// pass cannot answer them alone. It can answer them about its own rows and hand back the two
627 /// ends of itself, and then one comparison against the value before the vector joins the two
628 /// halves. That is what this does, and it is the whole of the sequential dependency.
629 fn fold(&mut self, one: Reduced) {
630 self.rows += one.rows;
631 self.nulls += one.nulls;
632 self.bytes = self.bytes.saturating_add(one.bytes);
633 self.widest = self.widest.max(one.widest);
634 let Some(ends) = one.ends else { return };
635 match self.previous.take() {
636 None => {
637 self.runs = 1;
638 self.first = Some(ends.first.clone());
639 }
640 Some(previous) => self.run(Some(previous.order(&ends.first))),
641 }
642 self.runs += one.descents;
643 if one.descents > 0 {
644 self.ascending = false;
645 }
646 if one.ascents > 0 {
647 self.descending = false;
648 }
649 if takes(&self.low, &ends.low, Ordering::Less) {
650 self.low = Some(ends.low.clone());
651 }
652 if takes(&self.stripe_low, &ends.low, Ordering::Less) {
653 self.stripe_low = Some(ends.low);
654 }
655 if takes(&self.high, &ends.high, Ordering::Greater) {
656 self.high = Some(ends.high.clone());
657 }
658 if takes(&self.stripe_high, &ends.high, Ordering::Greater) {
659 self.stripe_high = Some(ends.high);
660 }
661 self.previous = Some(ends.last);
662 }
663
664 /// The bound this column writes for a signed value, which a decimal column spells differently.
665 fn bound(&self, signed: i128) -> Bound {
666 match self.scale {
667 Some(scale) => Bound::Scaled { unscaled: signed, scale },
668 None => Bound::Int(signed),
669 }
670 }
671
672 /// One vector, a row at a time, for every column the fast path above does not read.
673 ///
674 /// The floats, the unsigned widths, the strings, and every form that is not flat. A string
675 /// column is here rather than in the fast path because its values are not a slice of one width
676 /// and its ends are byte comparisons, and `bytes_value` is already written to not allocate.
677 fn scan_rows(&mut self, vector: &Vector) {
678 // row at a time: the run count and the order flags are a sequential dependency. Whether this
679 // value is below the one before it is a question about a pair of adjacent rows, so there is
680 // no shape of this loop that answers it a vector at a time, and the two typed accessors
681 // below are loads against a slice rather than value construction. What the checker is
682 // looking for is the third arm, which does build a `Value`, and that one runs for a float
683 // column and for a form the first two cannot read and for nothing else.
684 for row in 0..vector.len() {
685 self.rows += 1;
686 if vector.is_null_at(row) {
687 self.nulls += 1;
688 continue;
689 }
690 if let Some(signed) = vector.signed_at(row) {
691 let bound = match self.scale {
692 Some(scale) => Bound::Scaled { unscaled: signed, scale },
693 None => Bound::Int(signed),
694 };
695 self.value(bound, self.fixed.unwrap_or(8));
696 continue;
697 }
698 if let Some(bytes) = vector.bytes_at(row) {
699 self.bytes_value(bytes);
700 continue;
701 }
702 // row at a time: a float and a form neither typed accessor above can read have no slice
703 // to walk, so the value is built for this row and for no other.
704 let value = vector.value_at(row);
705 let width = self.fixed.unwrap_or_else(|| width(&value));
706 match Bound::of_value(&value) {
707 Some(bound) => self.value(bound, width),
708 None => {
709 // A non-null value with no ordered bound. Both ends go rather than the value
710 // being skipped, because an end computed from only the values that had bounds is
711 // an end that answers a MIN with a value the column does not hold.
712 self.bytes = self.bytes.saturating_add(width);
713 self.widest = self.widest.max(width);
714 self.bounded = false;
715 self.ascending = false;
716 self.descending = false;
717 }
718 }
719 }
720 }
721
722 /// One non-null value, as its bound and its width.
723 ///
724 /// Every end is compared before it is copied. The obvious way to write this is to hand the
725 /// bound to each end and let the end keep whichever is smaller, and that costs a clone a row per
726 /// end whether or not the row is one. For an integer that is four copies of a machine word and
727 /// hardly matters. For a string it is four allocations a row, and on SF1 `l_comment` that is
728 /// twenty four million of them for a column with two ends. Compared first, an end is copied once
729 /// on a sorted column and about log n times on a shuffled one.
730 fn value(&mut self, bound: Bound, width: u64) {
731 self.measure(width);
732 let ordering = self.previous.as_ref().map(|previous| previous.order(&bound));
733 if ordering.is_none() {
734 self.first = Some(bound.clone());
735 }
736 self.run(ordering);
737 if takes(&self.low, &bound, Ordering::Less) {
738 self.low = Some(bound.clone());
739 }
740 if takes(&self.high, &bound, Ordering::Greater) {
741 self.high = Some(bound.clone());
742 }
743 if takes(&self.stripe_low, &bound, Ordering::Less) {
744 self.stripe_low = Some(bound.clone());
745 }
746 if takes(&self.stripe_high, &bound, Ordering::Greater) {
747 self.stripe_high = Some(bound.clone());
748 }
749 self.previous = Some(bound);
750 }
751
752 /// The same for a byte string, without a `Vec` a row.
753 ///
754 /// A string column is where the pass above still allocates, because the bound it is handed had
755 /// to be built out of the slice before it could be compared to anything, and the row it keeps as
756 /// the previous one is a new `Vec` every row whether or not any end moved. Here nothing is built
757 /// to be compared, and the buffer the previous row owns is refilled rather than replaced, which
758 /// is an allocation on the first row of the column and none after it.
759 ///
760 /// This is the difference between statistics costing a tenth of the write and costing as much as
761 /// it. At SF1, `lineitem`'s five string columns took nineteen of the pass's twenty eight seconds
762 /// before this and its eleven numeric columns took the other nine.
763 fn bytes_value(&mut self, bytes: &[u8]) {
764 self.measure(bytes.len() as u64);
765 let ordering = match &self.previous {
766 None => {
767 self.first = Some(Bound::Bytes(bytes.to_vec()));
768 None
769 }
770 Some(Bound::Bytes(previous)) => Some(Some(previous.as_slice().cmp(bytes))),
771 // A bound of another domain in a byte column, which a column of one type cannot hold.
772 Some(_) => Some(None),
773 };
774 self.run(ordering);
775 if takes_bytes(&self.low, bytes, Ordering::Less) {
776 fill(&mut self.low, bytes);
777 }
778 if takes_bytes(&self.high, bytes, Ordering::Greater) {
779 fill(&mut self.high, bytes);
780 }
781 if takes_bytes(&self.stripe_low, bytes, Ordering::Less) {
782 fill(&mut self.stripe_low, bytes);
783 }
784 if takes_bytes(&self.stripe_high, bytes, Ordering::Greater) {
785 fill(&mut self.stripe_high, bytes);
786 }
787 fill(&mut self.previous, bytes);
788 }
789
790 /// What one value costs, which is the byte total and the widest of them.
791 fn measure(&mut self, width: u64) {
792 self.bytes = self.bytes.saturating_add(width);
793 self.widest = self.widest.max(width);
794 }
795
796 /// What this value standing above, below or level with the one before it does to the order flags.
797 ///
798 /// The outer `None` is the first value of the column. The inner one is a pair this build cannot
799 /// order, which a column of one type cannot produce and which costs an order claim rather than
800 /// being assumed away.
801 fn run(&mut self, ordering: Option<Option<Ordering>>) {
802 match ordering {
803 None => self.runs = 1,
804 Some(Some(Ordering::Less)) => self.descending = false,
805 Some(Some(Ordering::Greater)) => {
806 self.ascending = false;
807 self.runs += 1;
808 }
809 Some(Some(Ordering::Equal)) => {}
810 Some(None) => {
811 self.ascending = false;
812 self.descending = false;
813 }
814 }
815 }
816
817 /// Starts a stripe, which is `key` in the order the table will be read in.
818 fn open_stripe(&mut self, key: (u64, u64)) {
819 self.stripe_low = None;
820 self.stripe_high = None;
821 self.key = key;
822 self.first = None;
823 self.previous = None;
824 self.ascending = true;
825 self.descending = true;
826 self.runs = 0;
827 }
828
829 fn close_stripe(&mut self) {
830 // Both taken whatever happens, so that a stripe of nothing but nulls leaves neither end
831 // behind for the next stripe to be compared against.
832 if let (Some(low), Some(high)) = (self.stripe_low.take(), self.stripe_high.take()) {
833 self.stripes.push((low, high));
834 }
835 self.pieces.push(Piece {
836 key: self.key,
837 first: self.first.take(),
838 last: self.previous.take(),
839 ascending: self.ascending,
840 descending: self.descending,
841 runs: self.runs,
842 });
843 }
844
845 /// Takes in a pass that read whole stripes of the same column on its own. See [`Gather::absorb`].
846 ///
847 /// Only between stripes, which is the only place a pass is ever handed over: the fields that
848 /// describe the stripe being read are empty then on both sides.
849 fn absorb(&mut self, later: Pass) {
850 self.rows += later.rows;
851 self.nulls += later.nulls;
852 self.bounded &= later.bounded;
853 self.bytes = self.bytes.saturating_add(later.bytes);
854 self.widest = self.widest.max(later.widest);
855 if let Some(low) = later.low {
856 if takes(&self.low, &low, Ordering::Less) {
857 self.low = Some(low);
858 }
859 }
860 if let Some(high) = later.high {
861 if takes(&self.high, &high, Ordering::Greater) {
862 self.high = Some(high);
863 }
864 }
865 self.stripes.extend(later.stripes);
866 self.pieces.extend(later.pieces);
867 }
868
869 /// Puts the stripes' order fields together in the order the table is read in.
870 ///
871 /// A pass that never opened a stripe has nothing here and keeps what it counted as it went.
872 /// Otherwise the pieces are laid end to end by key: each one's own flags hold, and the seam
873 /// between two is one comparison of the last value of the first against the first value of the
874 /// second, which is the comparison the pass would have made had the rows come in that order.
875 /// Every piece that held a value started its run count at one, so a seam that is not a descent
876 /// joins two runs into one and gives one back.
877 fn settle(&mut self) {
878 if self.pieces.is_empty() {
879 return;
880 }
881 let mut pieces = std::mem::take(&mut self.pieces);
882 pieces.sort_by_key(|piece| piece.key);
883 let (mut ascending, mut descending, mut runs) = (true, true, 0_u64);
884 let mut previous: Option<Bound> = None;
885 for piece in pieces {
886 ascending &= piece.ascending;
887 descending &= piece.descending;
888 let (Some(first), Some(last)) = (piece.first, piece.last) else { continue };
889 runs += piece.runs;
890 if let Some(previous) = &previous {
891 match previous.order(&first) {
892 Some(Ordering::Less) => descending = false,
893 Some(Ordering::Greater) => ascending = false,
894 Some(Ordering::Equal) => {}
895 None => {
896 ascending = false;
897 descending = false;
898 }
899 }
900 if previous.order(&first) != Some(Ordering::Greater) {
901 runs = runs.saturating_sub(1);
902 }
903 }
904 previous = Some(last);
905 }
906 self.ascending = ascending;
907 self.descending = descending;
908 self.runs = runs;
909 }
910
911 fn finish(mut self, sketch: Sketch, stripes: Vec<Sketch>) -> Stats {
912 self.settle();
913 let present = self.rows - self.nulls;
914 // The one rule the module doc names. An exact distinct count is one the sketch never had to
915 // throw a value away to keep, and everything downstream of the count follows from this
916 // answer rather than from what the writer hoped.
917 let exact = sketch.is_exact();
918 let distinct = if exact {
919 sketch.len() as u64
920 } else {
921 // Rounded rather than truncated, and clamped under the rows it cannot exceed. An
922 // estimate above the row count is arithmetically possible and is always wrong, and a
923 // planner that sees one concludes a column has more distinct values than rows.
924 #[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss)]
925 let estimate = sketch.distinct().round().max(0.0) as u64;
926 estimate.min(present)
927 };
928 let summary = Summary {
929 rows: self.rows,
930 nulls: self.nulls,
931 low: if self.bounded { self.low } else { None },
932 high: if self.bounded { self.high } else { None },
933 // Every end here came from a value the column holds, because this pass read them all.
934 // That is the whole difference between a summary and a zone map, which is allowed to be
935 // wider than its column and so can only skip and never answer.
936 ends_exact: self.bounded,
937 distinct,
938 // Exact or estimated, and never certified. A KMV sketch's relative error is about one
939 // over the square root of k, which is a standard error and not a bound, and Certified
940 // in this codebase means a bound that holds. Calling a one and a half percent standard
941 // error a guarantee is how an estimate gets treated as an answer.
942 distinct_class: if exact { Class::Exact } else { Class::Estimated },
943 // Only from an exact count. A sketch that overflowed cannot tell a column of a million
944 // unique values from one where two of them repeat, and uniqueness is the claim a key
945 // map is built on.
946 unique: exact && distinct == present,
947 order: if present == 0 {
948 Order::Neither
949 } else if self.ascending {
950 Order::Ascending
951 } else if self.descending {
952 Order::Descending
953 } else {
954 Order::Neither
955 },
956 runs: self.runs,
957 overlapping: overlapping(&self.stripes),
958 bytes: self.bytes,
959 widest: self.widest,
960 newest: self.generation,
961 };
962 // `new` rather than `merged` even for the empty case, because the two differ only in
963 // whether the list is checked and an empty list passes. A stripe sketch that is not at
964 // STRIPE_K is a bug in this file and is worth hearing about here rather than at the read.
965 let sketches = match Sketches::new(sketch.clone(), stripes) {
966 Ok(sketches) => sketches,
967 // Unreachable, since every stripe sketch above came out of `narrowed(STRIPE_K)` and a
968 // table cannot hold a million stripes. The merged sketch alone is the answer anyway:
969 // per stripe sketches are an optimization over a summary that is complete without
970 // them, so losing them costs a skipped stripe and never an answer.
971 Err(_) => Sketches::merged(sketch),
972 };
973 Stats { summary, sketches }
974 }
975}
976
977/// Whether an end has to become this bound, which is the question that replaces a clone.
978///
979/// `want` is [`Ordering::Less`] for a low end and [`Ordering::Greater`] for a high one. An end that
980/// is not there yet takes any value. A pair this build cannot order leaves the end alone, which is
981/// what [`Bound::smaller`] does across domains and which a column of one type cannot reach anyway.
982/// A dictionary the pass has read, and the positions and widths its codes stand for.
983///
984/// Kept from one vector to the next, and kept by the identity of the values it was read from rather
985/// than by a guess about what the caller is doing. One Parquet dictionary page serves every data
986/// page of its column chunk, so a load hands over a hundred vectors that share a dictionary, and
987/// reading it once instead of a hundred times is most of what this arm is worth. The `Arc` is held
988/// rather than its address noted, because a freed allocation's address is one a later dictionary can
989/// be handed and a cache keyed on that would read the wrong values and never know.
990#[derive(Debug)]
991struct Coded {
992 values: Arc<Vector>,
993 /// Position and width per code, `None` for a code whose entry is null.
994 codes: Vec<Option<(u32, u64)>>,
995 /// The distinct bounds in ascending order, which is what a position indexes.
996 bounds: Vec<Bound>,
997}
998
999impl Coded {
1000 /// One vector of codes, reduced against this dictionary.
1001 ///
1002 /// The loop this whole arm is for. A code lookup, a bounds check and an integer comparison, for
1003 /// a column whose row at a time path was comparing byte strings.
1004 fn reduce(&self, codes: &[u32], rows: usize, validity: &Validity) -> Reduced {
1005 let nullable = validity.has_nulls(rows);
1006 let mut out = Reduced::empty(rows as u64);
1007 let (mut low, mut high, mut first, mut last) = (0_u32, 0_u32, 0_u32, 0_u32);
1008 for (row, &code) in codes.iter().take(rows).enumerate() {
1009 let entry = if nullable && !validity.is_valid(row) {
1010 None
1011 } else {
1012 self.codes.get(code as usize).copied().flatten()
1013 };
1014 let Some((position, width)) = entry else {
1015 out.nulls += 1;
1016 continue;
1017 };
1018 out.bytes = out.bytes.saturating_add(width);
1019 out.widest = out.widest.max(width);
1020 if out.values == 0 {
1021 low = position;
1022 high = position;
1023 first = position;
1024 } else if position < last {
1025 out.descents += 1;
1026 } else if position > last {
1027 out.ascents += 1;
1028 }
1029 low = low.min(position);
1030 high = high.max(position);
1031 last = position;
1032 out.values += 1;
1033 }
1034 let at = |position: u32| self.bounds[position as usize].clone();
1035 out.ends = (out.values > 0).then(|| Ends {
1036 low: at(low),
1037 high: at(high),
1038 first: at(first),
1039 last: at(last),
1040 });
1041 out
1042 }
1043}
1044
1045/// What one vector came to, in the terms the pass folds rather than in the terms it was read in.
1046///
1047/// The two fast arms read a vector very differently and reduce it to the same nine numbers, so the
1048/// folding is written once. Everything here is about the vector alone: nothing in it depends on the
1049/// vector before, which is the half [`Pass::fold`] settles.
1050#[derive(Debug)]
1051struct Reduced {
1052 rows: u64,
1053 nulls: u64,
1054 /// Non-null values, which is what says whether `ends` means anything.
1055 values: u64,
1056 bytes: u64,
1057 widest: u64,
1058 /// Adjacent non-null pairs where the later value is the larger, which rules out a descending
1059 /// column, and where it is the smaller, which starts a run.
1060 ascents: u64,
1061 descents: u64,
1062 ends: Option<Ends>,
1063}
1064
1065impl Reduced {
1066 fn empty(rows: u64) -> Self {
1067 Self { rows, nulls: 0, values: 0, bytes: 0, widest: 0, ascents: 0, descents: 0, ends: None }
1068 }
1069}
1070
1071/// The four values of a vector the pass needs by name: its two ends, and its two edges.
1072#[derive(Debug)]
1073struct Ends {
1074 low: Bound,
1075 high: Bound,
1076 /// The first and last non-null values, for joining to the vectors either side.
1077 first: Bound,
1078 last: Bound,
1079}
1080
1081/// The bound of a dictionary entry that [`Pass::scan_dictionary`] has already found is not null.
1082fn bound_of(entries: &[Option<(Bound, u64)>], at: usize) -> &Bound {
1083 match &entries[at] {
1084 Some((bound, _)) => bound,
1085 // Unreachable: every index handed here came out of the filter that dropped the nulls. The
1086 // low bound is the answer that costs a wider range rather than a wrong one, if it ever is.
1087 None => &Bound::Int(i128::MIN),
1088 }
1089}
1090
1091/// What one vector of a flat signed column came to, computed without building a single [`Bound`].
1092///
1093/// Everything a [`Pass`] needs from a vector that is not about the vector before it. The two ends,
1094/// the two rows at the edges so that the joining comparison can be made, and the counts.
1095#[derive(Debug)]
1096struct Spread<T = i128> {
1097 /// Rows in the vector, nulls included.
1098 rows: u64,
1099 nulls: u64,
1100 /// The ends, meaningless when `values` is zero.
1101 low: T,
1102 high: T,
1103 /// The first and last non-null values, for joining to the vectors either side.
1104 first: T,
1105 last: T,
1106 /// Adjacent non-null pairs where the later value is the smaller, which is what starts a run.
1107 descents: u64,
1108 /// And where it is the larger, which is what rules out a descending column.
1109 ascents: u64,
1110 /// Non-null values, which is what the byte total is a multiple of.
1111 values: u64,
1112}
1113
1114/// One pass over a vector's non-null values, reading them through `get`.
1115///
1116/// Generic over the reader rather than over the element type, so that the caller can widen a layout
1117/// into an `i128` at the call site and this gets compiled once per layout with the widening inlined.
1118fn spread(rows: usize, validity: &Validity, get: impl Fn(usize) -> i128) -> Spread {
1119 let mut out = Spread {
1120 rows: rows as u64,
1121 nulls: 0,
1122 low: 0,
1123 high: 0,
1124 first: 0,
1125 last: 0,
1126 descents: 0,
1127 ascents: 0,
1128 values: 0,
1129 };
1130 let nullable = validity.has_nulls(rows);
1131 // row at a time: this is the loop the whole fast path is, and it is a row at a time because the
1132 // ascents and the descents are about adjacent rows. No `Value` is built here and none can be:
1133 // `get` hands back an `i128` read out of a typed slice.
1134 for row in 0..rows {
1135 if nullable && !validity.is_valid(row) {
1136 out.nulls += 1;
1137 continue;
1138 }
1139 let value = get(row);
1140 if out.values == 0 {
1141 out.low = value;
1142 out.high = value;
1143 out.first = value;
1144 } else {
1145 if value < out.last {
1146 out.descents += 1;
1147 } else if value > out.last {
1148 out.ascents += 1;
1149 }
1150 out.low = out.low.min(value);
1151 out.high = out.high.max(value);
1152 }
1153 out.last = value;
1154 out.values += 1;
1155 }
1156 out
1157}
1158
1159/// [`spread`] for a float column, or `None` when a value is NaN.
1160///
1161/// NaN is the one float with no order, and the row at a time path has its own answers for it, so a
1162/// vector holding one goes back there rather than this path making up another. Every other pair of
1163/// floats compares the way [`Bound::order`] compares them. The ends move only on a strictly smaller
1164/// or larger value, which keeps the first of `0.0` and `-0.0` the way the row path does.
1165fn real_spread(
1166 rows: usize,
1167 validity: &Validity,
1168 get: impl Fn(usize) -> f64,
1169) -> Option<Spread<f64>> {
1170 let mut out = Spread {
1171 rows: rows as u64,
1172 nulls: 0,
1173 low: 0.0,
1174 high: 0.0,
1175 first: 0.0,
1176 last: 0.0,
1177 descents: 0,
1178 ascents: 0,
1179 values: 0,
1180 };
1181 let nullable = validity.has_nulls(rows);
1182 // row at a time: the same loop as `spread`, for the same reason, over an `f64` read out of a
1183 // typed slice.
1184 for row in 0..rows {
1185 if nullable && !validity.is_valid(row) {
1186 out.nulls += 1;
1187 continue;
1188 }
1189 let value = get(row);
1190 if value.is_nan() {
1191 return None;
1192 }
1193 if out.values == 0 {
1194 out.low = value;
1195 out.high = value;
1196 out.first = value;
1197 } else {
1198 if value < out.last {
1199 out.descents += 1;
1200 } else if value > out.last {
1201 out.ascents += 1;
1202 }
1203 if value < out.low {
1204 out.low = value;
1205 }
1206 if value > out.high {
1207 out.high = value;
1208 }
1209 }
1210 out.last = value;
1211 out.values += 1;
1212 }
1213 Some(out)
1214}
1215
1216fn takes(held: &Option<Bound>, bound: &Bound, want: Ordering) -> bool {
1217 match held {
1218 None => true,
1219 Some(held) => bound.order(held) == Some(want),
1220 }
1221}
1222
1223/// The same question asked of a slice, so that nothing is built to ask it.
1224fn takes_bytes(held: &Option<Bound>, bytes: &[u8], want: Ordering) -> bool {
1225 match held {
1226 None => true,
1227 Some(Bound::Bytes(held)) => bytes.cmp(held.as_slice()) == want,
1228 Some(_) => false,
1229 }
1230}
1231
1232/// Puts these bytes in an end, reusing the buffer that is already there.
1233///
1234/// The whole of the byte path's advantage. A `Vec` that is cleared and refilled does not allocate
1235/// once it is wide enough, and these ends plus the previous row are where every allocation of the
1236/// value path went.
1237fn fill(held: &mut Option<Bound>, bytes: &[u8]) {
1238 match held {
1239 Some(Bound::Bytes(held)) => {
1240 held.clear();
1241 held.extend_from_slice(bytes);
1242 }
1243 held => *held = Some(Bound::Bytes(bytes.to_vec())),
1244 }
1245}
1246
1247/// Whether any two of these stripe ranges overlap.
1248///
1249/// Sorted by low end and then walked, so this is one sort rather than the square. A pair this cannot
1250/// order counts as overlapping, which is the answer that costs a skipped stripe rather than a wrong
1251/// one.
1252fn overlapping(stripes: &[(Bound, Bound)]) -> bool {
1253 let mut order = (0..stripes.len()).collect::<Vec<_>>();
1254 order
1255 .sort_by(|&one, &other| stripes[one].0.order(&stripes[other].0).unwrap_or(Ordering::Equal));
1256 order.windows(2).any(|pair| {
1257 let before = &stripes[pair[0]].1;
1258 let after = &stripes[pair[1]].0;
1259 before.order(after) != Some(Ordering::Less)
1260 })
1261}
1262
1263/// What every value of this type takes, when they all take the same.
1264///
1265/// `None` for the variable width types, which is the two string ones and nothing else. Read off the
1266/// type once by `Pass::new` rather than off each value.
1267fn fixed_width(ty: &LogicalType) -> Option<u64> {
1268 Some(match ty {
1269 LogicalType::Boolean | LogicalType::TinyInt | LogicalType::UTinyInt => 1,
1270 LogicalType::SmallInt | LogicalType::USmallInt => 2,
1271 LogicalType::Integer | LogicalType::UInteger | LogicalType::Float | LogicalType::Date => 4,
1272 LogicalType::HugeInt | LogicalType::UHugeInt | LogicalType::Decimal { .. } => 16,
1273 LogicalType::Varchar | LogicalType::Blob => return None,
1274 // The eight byte types: the two big integers, the double, and the four time ones. Anything
1275 // else that reaches here is refused a summary by `countable` long before this.
1276 _ => 8,
1277 })
1278}
1279
1280/// What one value takes, for the byte total and the widest value.
1281///
1282/// The logical width and not the stored one. The stored width is what the column's encoding chose
1283/// and is already in the layout; this is what the value costs a plan that has to materialize it,
1284/// which is the number a hash table sizing decision wants.
1285fn width(value: &Value) -> u64 {
1286 match value {
1287 Value::Null => 0,
1288 Value::Boolean(_) | Value::TinyInt(_) | Value::UTinyInt(_) => 1,
1289 Value::SmallInt(_) | Value::USmallInt(_) => 2,
1290 Value::Integer(_) | Value::UInteger(_) | Value::Float(_) | Value::Date(_) => 4,
1291 Value::HugeInt(_) | Value::UHugeInt(_) | Value::Decimal { .. } => 16,
1292 Value::Varchar(text) => text.len() as u64,
1293 Value::Blob(bytes) => bytes.len() as u64,
1294 // The eight byte types and anything else, which is every remaining scalar. A nested value
1295 // reaching here would be counted at eight and is refused a summary long before this by
1296 // `countable`.
1297 _ => 8,
1298 }
1299}
1300
1301/// One column's statistics built as the rows go past on their way into the file.
1302///
1303/// # Why this exists beside [`build_summary`]
1304///
1305/// Section 3.7 gives the build ten percent of the native write time, and [`build_summary`] cannot
1306/// fit inside that however tight its inner loop gets, because it starts by reading the file back. A
1307/// second full read of a committed table, decode included, is not ten percent of the first one. It
1308/// is most of it: on a TPC-H SF1 `lineitem` the standalone build is 11.4 seconds against a write of
1309/// 20.0 seconds of processor time, and the read is the bulk of the 11.4.
1310///
1311/// The writer has the vectors already. It buffers a stripe as chunks and hands one column of all of
1312/// them to each encode worker, so every value is in memory, in `rid` order, on a thread that is
1313/// about to walk it anyway. What is left of the build once the read is taken out is the hashing and
1314/// the comparisons, and those do fit. So this is the same [`Pass`] and the same [`Counts`] driven
1315/// from the write rather than from a reader, and [`build_summary`] stays as the path for a file
1316/// that was written before any of this existed.
1317///
1318/// # No per stripe sketches here
1319///
1320/// Section 3.8 promotes a column when something has declared a relationship or a key over it, and
1321/// [`read_columns`] reads that off the file. A table being written for the first time has no
1322/// sections at all, so the promoted set is empty by construction and there is nothing for this to
1323/// decide. A later checkpoint that declares a key is what promotes the column, and that goes through
1324/// [`build_stats_for`] with the file in front of it.
1325#[derive(Debug)]
1326pub(crate) struct Gather {
1327 pass: Pass,
1328 counts: Counts,
1329}
1330
1331impl Gather {
1332 /// One for a column that can be summarized, and nothing for one that cannot.
1333 ///
1334 /// `None` rather than an error, because a table with an interval column in it still gets
1335 /// summaries for its other fifteen and section 3.1 says the interval column plans the way it
1336 /// planned before.
1337 pub(crate) fn new(ty: &LogicalType, generation: u64) -> Option<Self> {
1338 countable(ty).then(|| Self { pass: Pass::new(ty, generation), counts: Counts::new(1) })
1339 }
1340
1341 /// Opens a stripe of this column, whose parts come to [`Gather::part`] in part order until
1342 /// [`Gather::close_stripe`].
1343 ///
1344 /// The stripe is the unit the pass opens and closes its ends over, so a caller says where one
1345 /// starts and stops. The key is where the stripe goes once the writer sorts its stripes, which
1346 /// need not be the order they reach this in.
1347 pub(crate) fn open_stripe(&mut self, key: (u64, u64)) {
1348 self.pass.open_stripe(key);
1349 }
1350
1351 /// Takes one part of the stripe that is open, in order.
1352 pub(crate) fn part(&mut self, vector: &Vector) {
1353 self.counts.add_column(0, vector);
1354 self.pass.scan(vector);
1355 }
1356
1357 /// Ends the stripe that is open.
1358 pub(crate) fn close_stripe(&mut self) {
1359 self.pass.close_stripe();
1360 }
1361
1362 /// Takes in a gather that folded stripes of the same column on its own, as though this had
1363 /// folded them.
1364 ///
1365 /// This is what lets a stripe be summarized on the thread that encodes it, before the writer's
1366 /// lock is taken. Nothing a stripe adds depends on the stripes before it: the order fields are
1367 /// kept a stripe at a time and put together by key at the end, the ends and the totals are a
1368 /// minimum, a maximum and sums, and the counts union. The one thing that does depend on order
1369 /// is the tally's list, which comes out in the order the stripes are absorbed in, and that is
1370 /// the order they reached the writer in, which is what it was before.
1371 pub(crate) fn absorb(&mut self, later: Gather) {
1372 self.pass.absorb(later.pass);
1373 self.counts.absorb(later.counts);
1374 }
1375
1376 /// How many rows went past, which is what the caller checks against the table's own count.
1377 pub(crate) fn rows(&self) -> u64 {
1378 self.pass.rows
1379 }
1380
1381 /// The sketch's estimate of the column's distinct values so far, or nothing for a blind one.
1382 pub(crate) fn distinct(&self) -> Option<f64> {
1383 self.counts.sketch(0).map(|sketch| sketch.distinct())
1384 }
1385
1386 /// Every non-null value of the column with the rows holding it, and the rows holding a null,
1387 /// while the tally still holds the whole column.
1388 ///
1389 /// Nothing once the column has passed the tally's cap or turned out to be blind. The counts are
1390 /// exact, which is what lets the close take a narrow column's frequencies from here rather than
1391 /// read its pages back and count them a second time.
1392 pub(crate) fn frequencies(&self) -> Option<(Vec<(Value, u64)>, u64)> {
1393 Some((self.counts.frequencies(0)?, self.pass.nulls))
1394 }
1395
1396 /// The summary and the merged sketch, or nothing if the column turned out to be blind.
1397 ///
1398 /// Blind means a form `rudb_storage::count` has no arm for turned up, so the sketch is missing
1399 /// rows and cannot say which. A distinct count that is too low is the one error an estimator has
1400 /// no defence against, so the column gets no sections rather than sections with a number in them
1401 /// nothing can check.
1402 pub(crate) fn finish(self) -> Option<Stats> {
1403 let sketch = self.counts.sketch(0)?;
1404 Some(self.pass.finish(sketch, Vec::new()))
1405 }
1406}
1407
1408/// Everything the columns of a table being written cost so far, which is what the budget is a share
1409/// of.
1410///
1411/// The same sum [`crate::Layout::columns_total`] takes, off the table rather than off a reader,
1412/// because the writer has no reader and the file it would open is not committed yet. Every stripe's
1413/// pages are written by the time this is asked and so are the dictionaries, so the two agree.
1414pub(crate) fn column_bytes(table: &crate::Table) -> u64 {
1415 (0..table.fields.len())
1416 .map(|at| {
1417 crate::sum(table.stripes.iter().map(|stripe| crate::span_bytes(&stripe.pages, at)))
1418 .saturating_add(crate::sum(
1419 table.stripes.iter().map(|stripe| stripe.memberships.bytes(at)),
1420 ))
1421 .saturating_add(crate::sum(
1422 table.stripes.iter().map(|stripe| stripe.sieves.bytes(at)),
1423 ))
1424 .saturating_add(crate::sum(
1425 table.stripes.iter().map(|stripe| stripe.part_ranges.bytes(at)),
1426 ))
1427 .saturating_add(crate::dictionary_bytes(table, at))
1428 })
1429 .fold(0, u64::saturating_add)
1430}
1431
1432/// Which of these payloads fit the allowance, smallest first.
1433///
1434/// Smallest first so that a budget that cannot hold everything holds as many columns as it can. The
1435/// alternative is column order, which would give the summaries to whichever columns the schema
1436/// happened to list early, and there is nothing about being the first column that makes a summary
1437/// worth more.
1438pub(crate) fn within(costs: &[usize], allowance: u64, spent: u64) -> Vec<bool> {
1439 let mut order = (0..costs.len()).collect::<Vec<_>>();
1440 order.sort_by_key(|&at| costs[at]);
1441 let mut spent = spent;
1442 let mut keep = vec![false; costs.len()];
1443 for at in order {
1444 let cost = costs[at] as u64;
1445 if spent.saturating_add(cost) <= allowance {
1446 spent += cost;
1447 keep[at] = true;
1448 }
1449 }
1450 keep
1451}
1452
1453/// What the allowance is for a table whose columns come to this many bytes.
1454pub(crate) fn allowance(column_bytes: u64, share: u64) -> u64 {
1455 (column_bytes.saturating_mul(share) / 100).max(BUDGET_FLOOR)
1456}
1457
1458/// Builds the statistics for each of these columns and attaches them all in one commit.
1459///
1460/// One commit and not one each, for the reason `graph::build_key_maps` gives: a checkpoint that
1461/// published one generation per column would be one chance per column of being interrupted halfway.
1462///
1463/// # Errors
1464///
1465/// If the file cannot be opened, a column cannot be summarized, or the attach fails.
1466pub fn build_stats(path: &Path, table: &str, columns: &[usize]) -> Result<Vec<Built>> {
1467 build_stats_within(path, table, columns, BUDGET_SHARE)
1468}
1469
1470/// The columns of this table the per stripe rule promotes, in column order.
1471///
1472/// Section 3.8's default set: the columns something has declared a relationship or a key over. What
1473/// this build has to go on for that is the file itself, so the answer is the columns that already
1474/// carry a graph section, which is a key map or a forward link. That is not a proxy for the
1475/// question, it is the same question asked of the only party that has been told the answer: a key
1476/// map exists on a column because something declared it a key.
1477///
1478/// Empty is the ordinary answer and it is the right one. A table nothing has declared anything over
1479/// gets table level summaries and no per stripe sketches, which is what section 3.8 says and what
1480/// keeps SF100 inside two percent.
1481///
1482/// The other source the spec names is document 06's observation log, which promotes a column that
1483/// queries turned out to read at the next checkpoint. It is not built yet. When it is, it adds
1484/// columns here and nothing else in this file changes.
1485#[must_use]
1486pub fn read_columns(reader: &Reader) -> Vec<usize> {
1487 let generation = reader.table().generation();
1488 let mut promoted = reader
1489 .table()
1490 .sections()
1491 .iter()
1492 .filter(|held| held.among(section::GRAPH_KINDS) && held.usable(generation))
1493 .filter_map(|held| usize::try_from(held.id).ok())
1494 .collect::<Vec<_>>();
1495 promoted.sort_unstable();
1496 promoted.dedup();
1497 promoted
1498}
1499
1500/// The same, against a budget of `share` percent of the table's stored column bytes.
1501///
1502/// The budget is over the table rather than over a column, and when it binds the cheapest columns
1503/// are admitted first. That is the same degenerate case section 3.7's expected value ordering has
1504/// for a key map with no relationship over it: nothing has said which column a plan will ask about,
1505/// so no summary is worth more than another and the ordering falls back to the denominator. Cheapest
1506/// first is also the order that fits the most summaries in the room there is.
1507///
1508/// A column is all or nothing. Its summary and its sketches are admitted together or neither is,
1509/// because a summary whose distinct count came from a sketch that was then dropped is a number with
1510/// nothing behind it to check it against.
1511///
1512/// # Errors
1513///
1514/// If the file cannot be opened, a column cannot be summarized, or the attach fails.
1515pub fn build_stats_within(
1516 path: &Path,
1517 table: &str,
1518 columns: &[usize],
1519 share: u64,
1520) -> Result<Vec<Built>> {
1521 let promoted = read_columns(&Catalog::open(path)?.table(table)?);
1522 build_stats_for(path, table, columns, &promoted, share)
1523}
1524
1525/// The same, with the per stripe set named rather than read off the file.
1526///
1527/// For a caller that knows something this build does not, which today is the measurement harness and
1528/// tomorrow is whatever reads document 06's observation log. [`build_stats_within`] is the ordinary
1529/// entry point and it asks [`read_columns`].
1530///
1531/// A column in `per_stripe` that is not in `columns` is ignored rather than refused, because the two
1532/// lists answer different questions and a caller that names a promoted column it is not building is
1533/// not making a mistake worth stopping for.
1534///
1535/// # Errors
1536///
1537/// If the file cannot be opened, a column cannot be summarized, or the attach fails.
1538pub fn build_stats_for(
1539 path: &Path,
1540 table: &str,
1541 columns: &[usize],
1542 per_stripe: &[usize],
1543 share: u64,
1544) -> Result<Vec<Built>> {
1545 let reader = Catalog::open(path)?.table(table)?;
1546 let column_bytes = reader.layout().columns_total();
1547 let allowance = allowance(column_bytes, share);
1548 let spent = held_bytes(&reader, columns)?;
1549 let mut report = Vec::with_capacity(columns.len());
1550 let mut payloads = Vec::with_capacity(columns.len());
1551 for &column in columns {
1552 let start = Instant::now();
1553 let stats = build_summary_for(&reader, column, per_stripe.contains(&column))?;
1554 let mut summary = Vec::new();
1555 stats.summary.encode(&mut summary)?;
1556 let mut sketches = Vec::new();
1557 stats.sketches.encode(&mut sketches)?;
1558 report.push(Built {
1559 column,
1560 rows: stats.summary.rows,
1561 distinct: stats.summary.distinct,
1562 exact: stats.summary.distinct_class == Class::Exact,
1563 order: stats.summary.order,
1564 summary_bytes: summary.len(),
1565 sketch_bytes: sketches.len(),
1566 stripes: stats.sketches.stripes.len(),
1567 column_bytes,
1568 built: false,
1569 build: start.elapsed(),
1570 });
1571 payloads.push((column, summary, sketches));
1572 }
1573 let costs = report.iter().map(Built::bytes).collect::<Vec<_>>();
1574 let keep = within(&costs, allowance, spent);
1575 for (one, &keep) in report.iter_mut().zip(&keep) {
1576 one.built = keep;
1577 }
1578 // The reader holds the file open and the attach opens it again to write, so it is dropped first
1579 // for the reason `graph` drops it: the moment the file is written is a moment nothing else in
1580 // this function is reading it.
1581 drop(reader);
1582 let mut attachments = Vec::with_capacity(payloads.len() * 2);
1583 for ((column, summary, sketches), _) in payloads.iter().zip(&keep).filter(|&(_, &keep)| keep) {
1584 let id = u64::try_from(*column).map_err(|_| invalid("column index overflow"))?;
1585 attachments.push(Attachment {
1586 kind: *section::SUMMARY,
1587 id,
1588 flags: 0,
1589 // A summary is a header the whole way down. There is nothing behind it that a reader
1590 // could decide not to read, which is the shape section 3.2's field is for and not a
1591 // misuse of it: the answer to "how much do I read to know what this says" is all of it.
1592 header_bytes: u32::try_from(summary.len())
1593 .map_err(|_| invalid("a summary longer than a u32 can count"))?,
1594 bytes: summary,
1595 });
1596 attachments.push(Attachment {
1597 kind: *section::SKETCHES,
1598 id,
1599 flags: 0,
1600 header_bytes: SKETCH_HEADER,
1601 bytes: sketches,
1602 });
1603 }
1604 crate::attach(path, table, &attachments)?;
1605 Ok(report)
1606}
1607
1608/// What the table's existing statistics sections cost, leaving out the ones this build is replacing.
1609///
1610/// Statistics sections only. The two percent of section 3.8 and the graph layer's ten percent are
1611/// separate shares of the same column bytes, and separate means each counts only what it owns. A
1612/// TPC-H SF10 file's key maps are 7.7 MB against a two percent allowance of 54 MB, so counting them
1613/// here would hand a seventh of the statistics budget to sections that already have one of their
1614/// own, and a table would lose summaries for a reason that has nothing to do with summaries.
1615///
1616/// Reading the extent tables is what this costs, which is one small read per section and not a read
1617/// of a payload. A section whose extent table does not checksum is counted as nothing, because it
1618/// is a section that is already not there.
1619fn held_bytes(reader: &Reader, replacing: &[usize]) -> Result<u64> {
1620 let mut total = 0;
1621 for held in reader.table().sections() {
1622 if !held.among(section::STATISTICS_KINDS) {
1623 continue;
1624 }
1625 let replaced = replacing.iter().any(|&column| u64::try_from(column) == Ok(held.id));
1626 if replaced || !held.usable(reader.table().generation()) {
1627 continue;
1628 }
1629 let Ok(extents) = reader.extents(held) else { continue };
1630 total += extents.iter().map(|extent| u64::from(extent.length)).sum::<u64>();
1631 }
1632 Ok(total)
1633}
1634
1635/// The summary this table carries for a column, when it carries one this build can use.
1636///
1637/// `None` covers every reason there is not one and covering them all is the point. Section 3.1 says
1638/// deleting every statistics section changes no answer, so there is no reason to distinguish *no
1639/// summary was built* from *the summary is stale*, *the payload does not checksum*, or *the layout
1640/// is one a later build invented*. The answer to all four is to plan the query the way it was
1641/// planned before summaries existed.
1642#[must_use]
1643pub fn summary(reader: &Reader, column: usize) -> Option<Summary> {
1644 let bytes = payload(reader, column, section::SUMMARY)?;
1645 Summary::decode(&bytes).ok()
1646}
1647
1648/// The sketches this table carries for a column, same.
1649///
1650/// One more reason for `None` here than above: a sketch built by a hash this build does not use is
1651/// declined by [`Sketches::decode`] rather than merged into anything, which costs a rebuild where
1652/// merging would cost an answer.
1653#[must_use]
1654pub fn sketches(reader: &Reader, column: usize) -> Option<Sketches> {
1655 let bytes = payload(reader, column, section::SKETCHES)?;
1656 Sketches::decode(&bytes).ok()
1657}
1658
1659fn payload(reader: &Reader, column: usize, kind: &[u8; 8]) -> Option<Vec<u8>> {
1660 let table = reader.table();
1661 let id = u64::try_from(column).ok()?;
1662 let held = table.sections().iter().find(|section| section.kind == *kind && section.id == id)?;
1663 if !held.usable(table.generation()) {
1664 return None;
1665 }
1666 reader.payload(held).ok()
1667}
1668
1669/// Whether a type can be summarized at all, which is whether it has a hash rule.
1670#[must_use]
1671pub fn summarizable(ty: &LogicalType) -> bool {
1672 countable(ty)
1673}
1674
1675#[cfg(test)]
1676mod tests {
1677 use std::fs;
1678 use std::path::PathBuf;
1679 use std::sync::Arc;
1680 use std::time::{SystemTime, UNIX_EPOCH};
1681
1682 use rudb_common::Field;
1683 use rudb_encoding::sketch::hash64;
1684 use rudb_storage::count::hash_value;
1685 use rudb_vector::{Chunk, Vector};
1686
1687 use super::*;
1688 use crate::Writer;
1689
1690 fn path(label: &str) -> PathBuf {
1691 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
1692 std::env::temp_dir().join(format!("rudb-stats-{label}-{}-{stamp}.rdb", std::process::id()))
1693 }
1694
1695 /// A one column table of these values, written a thousand rows to a part.
1696 fn table_of(label: &str, values: &[Option<i64>]) -> PathBuf {
1697 let path = path(label);
1698 let mut writer =
1699 Writer::create(&path, "t", vec![Field::new("v", LogicalType::BigInt)]).expect("new");
1700 for part in values.chunks(1000) {
1701 let held =
1702 part.iter().map(|v| v.map_or(Value::Null, Value::BigInt)).collect::<Vec<_>>();
1703 let chunk =
1704 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("values")])
1705 .expect("one column");
1706 writer.append(&chunk).expect("a part");
1707 }
1708 writer.finish().expect("commit");
1709 path
1710 }
1711
1712 /// The same, with the part size named, for a test that needs more than one stripe.
1713 ///
1714 /// A stripe is up to `STRIPE_PARTS` parts, so small parts are how a test crosses a stripe
1715 /// boundary without writing a hundred and thirty thousand rows to do it.
1716 fn table_of_parts(label: &str, values: &[Option<i64>], per_part: usize) -> PathBuf {
1717 let path = path(label);
1718 let mut writer =
1719 Writer::create(&path, "t", vec![Field::new("v", LogicalType::BigInt)]).expect("new");
1720 for part in values.chunks(per_part) {
1721 let held =
1722 part.iter().map(|v| v.map_or(Value::Null, Value::BigInt)).collect::<Vec<_>>();
1723 let chunk =
1724 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("values")])
1725 .expect("one column");
1726 writer.append(&chunk).expect("a part");
1727 }
1728 writer.finish().expect("commit");
1729 path
1730 }
1731
1732 #[test]
1733 fn stripes_that_arrive_out_of_order_are_summarized_in_the_order_they_are_read() {
1734 // What a parallel load does: three pipeline instances each hand the writer a contiguous
1735 // run of the source as its own stripe, and they finish in whatever order they finish. The
1736 // table reads back sorted by source position, so that is the order the summary is about.
1737 // Every key repeats across a seam, the way an order's line items straddle two stripes.
1738 let path = path("late-stripes");
1739 let mut writer =
1740 Writer::create(&path, "t", vec![Field::new("v", LogicalType::BigInt)]).expect("new");
1741 let part = |from: i64| {
1742 let held = (from..from + 10).map(|v| Value::BigInt(v / 2)).collect::<Vec<_>>();
1743 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("values")])
1744 .expect("one column")
1745 };
1746 for stripe in [2_u64, 0, 1] {
1747 let parts = (0..3)
1748 .map(|at| {
1749 (
1750 (stripe * 3 + at, 0),
1751 part(i64::try_from(stripe * 30 + at * 10).expect("small")),
1752 )
1753 })
1754 .collect();
1755 writer.append_stripe(parts).expect("a stripe");
1756 }
1757 writer.finish().expect("commit");
1758
1759 let reader = reopen(&path);
1760 let summary = summary(&reader, 0).expect("the summary is in the file");
1761 assert_eq!(summary.rows, 90);
1762 assert_eq!(summary.order, Order::Ascending, "the stripes are in order once sorted");
1763 assert_eq!(summary.runs, 1, "and the seams between them are not descents");
1764 assert_eq!(crate::ascending(&reader), vec!["v".to_owned()]);
1765 }
1766
1767 /// The vector at a time pass says exactly what the row at a time pass says.
1768 ///
1769 /// [`Pass::scan_flat`] and [`Pass::scan_dictionary`] took the ordinary columns off
1770 /// [`Pass::scan_rows`] and they are why the build fits inside its share of the write. What they
1771 /// have to be is not fast but identical, so each is driven here over the same vectors in the
1772 /// same stripes as the row at a time pass and the two summaries are compared whole.
1773 ///
1774 /// Six shapes and five types. The shapes, because the fields that differ between them are the
1775 /// order flags and the run count, and those are what a vector at a time pass has to rejoin by
1776 /// hand. The types, because the two arms read a value in three different ways between them and
1777 /// a bound that came out of one has to be the bound that came out of another.
1778 #[test]
1779 fn the_vector_at_a_time_pass_says_what_the_row_at_a_time_pass_says() {
1780 // Coprime with the length, so this visits every value once and every part spans the range.
1781 let shuffled = (0..500_i64).map(|at| Some(1 + at * 307 % 500)).collect::<Vec<_>>();
1782 let shapes: [(&str, Vec<Option<i64>>); 7] = [
1783 ("ascending", (1..=500_i64).map(Some).collect()),
1784 ("descending", (1..=500_i64).rev().map(Some).collect()),
1785 ("constant", vec![Some(7); 500]),
1786 ("shuffled", shuffled),
1787 ("every third null", (1..=500_i64).map(|at| (at % 3 != 0).then_some(at)).collect()),
1788 ("all nulls", vec![None; 500]),
1789 ("twenty values over and over", (0..500_i64).map(|at| Some(at * 7 % 20)).collect()),
1790 ];
1791 let types = [
1792 LogicalType::SmallInt,
1793 LogicalType::Integer,
1794 LogicalType::BigInt,
1795 LogicalType::Decimal { width: 18, scale: 2 },
1796 LogicalType::Varchar,
1797 ];
1798 for (label, values) in &shapes {
1799 for ty in &types {
1800 // Sixty rows to a vector and five vectors to a stripe, so the stripe ends and the
1801 // overlap answer are in the comparison rather than left where they started.
1802 let held = values
1803 .chunks(60)
1804 .map(|part| {
1805 let values = part.iter().map(|value| one(ty, *value)).collect::<Vec<_>>();
1806 Vector::from_values(ty.clone(), &values).expect("values")
1807 })
1808 .collect::<Vec<_>>();
1809 // The same rows again as a dictionary of the twenty distinct values a vector holds,
1810 // in an order that is not the sorted one, so that the positions the arm hands out
1811 // are doing work rather than agreeing with the codes by accident.
1812 let coded = values
1813 .chunks(60)
1814 .map(|part| {
1815 let mut distinct = part.to_vec();
1816 distinct.sort_unstable();
1817 distinct.dedup();
1818 distinct.reverse();
1819 let values =
1820 distinct.iter().map(|value| one(ty, *value)).collect::<Vec<_>>();
1821 let codes = part
1822 .iter()
1823 .map(|value| {
1824 distinct.iter().position(|held| held == value).expect("a code")
1825 as u32
1826 })
1827 .collect::<Vec<_>>();
1828 Vector::dictionary(
1829 codes,
1830 Vector::from_values(ty.clone(), &values).expect("values"),
1831 )
1832 .expect("a dictionary")
1833 })
1834 .collect::<Vec<_>>();
1835 let flat = drive(ty, &held, |pass, vector| {
1836 assert!(pass.scan_flat(vector), "{label} {ty}");
1837 });
1838 let dictionary = drive(ty, &coded, |pass, vector| {
1839 assert!(pass.scan_dictionary(vector), "{label} {ty} is dictionary coded");
1840 });
1841 let rows = drive(ty, &held, Pass::scan_rows);
1842 assert_eq!(flat.summary, rows.summary, "flat: {label} {ty}");
1843 assert_eq!(dictionary.summary, rows.summary, "dictionary: {label} {ty}");
1844 // A signed column's dictionaries read through their codes, per vector and then as
1845 // one dictionary of every distinct value, which is what a Parquet row group hands
1846 // over and is too wide for the arm above. A dictionary holding a null is turned
1847 // down and read a row at a time.
1848 if *ty != LogicalType::Varchar {
1849 let mut distinct = values.clone();
1850 distinct.sort_unstable();
1851 distinct.dedup();
1852 distinct.reverse();
1853 let every = Arc::new(
1854 Vector::from_values(
1855 ty.clone(),
1856 &distinct.iter().map(|value| one(ty, *value)).collect::<Vec<_>>(),
1857 )
1858 .expect("values"),
1859 );
1860 let wide = values
1861 .chunks(60)
1862 .map(|part| {
1863 let codes = part.iter().map(|value| code(values, *value)).collect();
1864 Vector::dictionary_over(codes, Arc::clone(&every))
1865 .expect("a dictionary")
1866 })
1867 .collect::<Vec<_>>();
1868 for vectors in [&coded, &wide] {
1869 let gathered = drive(ty, vectors, |pass, vector| {
1870 if !pass.scan_gathered(vector) {
1871 assert!(values.contains(&None), "{label} {ty} is gathered");
1872 pass.scan_rows(vector);
1873 }
1874 });
1875 assert_eq!(gathered.summary, rows.summary, "gathered: {label} {ty}");
1876 }
1877 }
1878 // And again over one dictionary that every vector shares, which is what a Parquet
1879 // load hands over and what the pass keeps its last dictionary for. Only for the
1880 // shapes narrow enough to have one, since a dictionary wider than the vector it
1881 // codes is one this arm turns down.
1882 let Some(shared) = shared(ty, values) else { continue };
1883 let coded = values
1884 .chunks(60)
1885 .map(|part| {
1886 let codes = part.iter().map(|value| code(values, *value)).collect();
1887 Vector::dictionary_over(codes, Arc::clone(&shared)).expect("a dictionary")
1888 })
1889 .collect::<Vec<_>>();
1890 let held = drive(ty, &coded, |pass, vector| {
1891 assert!(pass.scan_dictionary(vector), "{label} {ty} is dictionary coded");
1892 });
1893 assert_eq!(held.summary, rows.summary, "one dictionary: {label} {ty}");
1894 }
1895 }
1896 }
1897
1898 /// A float column read a vector at a time says what it says read a row at a time, and a vector
1899 /// holding a NaN goes back to the row at a time pass.
1900 #[test]
1901 fn a_float_column_a_vector_at_a_time_says_what_it_says_a_row_at_a_time() {
1902 let shuffled = (0..500_i64).map(|at| Some((1 + at * 307 % 500) as f64 / 4.0)).collect();
1903 let shapes: [(&str, Vec<Option<f64>>); 6] = [
1904 ("ascending", (1..=500).map(|at| Some(f64::from(at) * 0.5)).collect()),
1905 ("descending", (1..=500).rev().map(|at| Some(f64::from(at) * 0.5)).collect()),
1906 ("shuffled", shuffled),
1907 (
1908 "signed zeros",
1909 (0..500).map(|at| Some(if at % 2 == 0 { 0.0 } else { -0.0 })).collect(),
1910 ),
1911 (
1912 "every third null",
1913 (1..=500).map(|at| (at % 3 != 0).then_some(f64::from(at))).collect(),
1914 ),
1915 (
1916 "a NaN in one vector",
1917 (0..500).map(|at| Some(if at == 130 { f64::NAN } else { f64::from(at) })).collect(),
1918 ),
1919 ];
1920 for ty in [LogicalType::Double, LogicalType::Float] {
1921 for (label, values) in &shapes {
1922 let held = values
1923 .chunks(60)
1924 .map(|part| {
1925 let values = part
1926 .iter()
1927 .map(|value| match (value, &ty) {
1928 (None, _) => Value::Null,
1929 (Some(value), LogicalType::Float) => Value::Float(*value as f32),
1930 (Some(value), _) => Value::Double(*value),
1931 })
1932 .collect::<Vec<_>>();
1933 Vector::from_values(ty.clone(), &values).expect("values")
1934 })
1935 .collect::<Vec<_>>();
1936 let mut fell_back = 0;
1937 let flat = drive(&ty, &held, |pass, vector| {
1938 if !pass.scan_flat(vector) {
1939 fell_back += 1;
1940 pass.scan_rows(vector);
1941 }
1942 });
1943 let rows = drive(&ty, &held, Pass::scan_rows);
1944 assert_eq!(
1945 format!("{:?}", flat.summary),
1946 format!("{:?}", rows.summary),
1947 "{label} {ty}"
1948 );
1949 assert_eq!(fell_back, usize::from(label.contains("NaN")), "{label} {ty}");
1950 }
1951 }
1952 }
1953
1954 /// A string column read a vector at a time says what it says a row at a time.
1955 #[test]
1956 fn a_string_column_a_vector_at_a_time_says_what_it_says_a_row_at_a_time() {
1957 let word = |at: i64| format!("w{:03}", at);
1958 let shapes: [(&str, Vec<Option<String>>); 6] = [
1959 ("ascending", (0..500).map(|at| Some(word(at))).collect()),
1960 ("descending", (0..500).rev().map(|at| Some(word(at))).collect()),
1961 ("shuffled", (0..500).map(|at| Some(word(at * 307 % 500))).collect()),
1962 ("repeated", (0..500).map(|at| Some(word(at / 7 % 5))).collect()),
1963 ("prefixes", (0..500).map(|at| Some("ab".repeat(1 + at % 9))).collect()),
1964 (
1965 "every third null",
1966 (0..500).map(|at| (at % 3 != 0).then(|| word(at * 13 % 500))).collect(),
1967 ),
1968 ];
1969 for ty in [LogicalType::Varchar] {
1970 for (label, values) in &shapes {
1971 let held = values
1972 .chunks(60)
1973 .map(|part| {
1974 let values = part
1975 .iter()
1976 .map(|value| value.clone().map_or(Value::Null, Value::Varchar))
1977 .collect::<Vec<_>>();
1978 Vector::from_values(ty.clone(), &values).expect("values")
1979 })
1980 .collect::<Vec<_>>();
1981 let mut fell_back = 0;
1982 let flat = drive(&ty, &held, |pass, vector| {
1983 if !pass.scan_flat(vector) {
1984 fell_back += 1;
1985 pass.scan_rows(vector);
1986 }
1987 });
1988 let rows = drive(&ty, &held, Pass::scan_rows);
1989 assert_eq!(
1990 format!("{:?}", flat.summary),
1991 format!("{:?}", rows.summary),
1992 "{label} {ty}"
1993 );
1994 assert_eq!(fell_back, 0, "{label} {ty}");
1995 }
1996 }
1997 }
1998
1999 /// The distinct values of a column as one dictionary, or nothing if there are too many of them
2000 /// for [`Pass::scan_dictionary`] to take it.
2001 fn shared(ty: &LogicalType, values: &[Option<i64>]) -> Option<Arc<Vector>> {
2002 let mut distinct = values.to_vec();
2003 distinct.sort_unstable();
2004 distinct.dedup();
2005 // Wider than the sixty rows a vector holds is what the arm turns down, and a test that fed
2006 // it one would be asserting over the row at a time pass twice.
2007 if distinct.len() > 60 {
2008 return None;
2009 }
2010 // Reversed, so the positions the arm hands out are doing work rather than agreeing with the
2011 // codes by accident.
2012 distinct.reverse();
2013 let held = distinct.iter().map(|value| one(ty, *value)).collect::<Vec<_>>();
2014 Some(Arc::new(Vector::from_values(ty.clone(), &held).expect("values")))
2015 }
2016
2017 /// Where a value sits in the dictionary [`shared`] builds.
2018 fn code(values: &[Option<i64>], value: Option<i64>) -> u32 {
2019 let mut distinct = values.to_vec();
2020 distinct.sort_unstable();
2021 distinct.dedup();
2022 distinct.reverse();
2023 distinct.iter().position(|held| *held == value).expect("a code") as u32
2024 }
2025
2026 /// One value of this type, or a null, for the equivalence test above.
2027 fn one(ty: &LogicalType, value: Option<i64>) -> Value {
2028 let Some(value) = value else { return Value::Null };
2029 match ty {
2030 LogicalType::SmallInt => Value::SmallInt(value as i16),
2031 LogicalType::Integer => Value::Integer(value as i32),
2032 LogicalType::BigInt => Value::BigInt(value),
2033 LogicalType::Varchar => Value::Varchar(format!("v{value:04}")),
2034 _ => Value::Decimal { unscaled: i128::from(value), width: 18, scale: 2 },
2035 }
2036 }
2037
2038 /// A whole pass over these vectors, five to a stripe, read by whichever arm the caller names.
2039 fn drive(ty: &LogicalType, held: &[Vector], mut scan: impl FnMut(&mut Pass, &Vector)) -> Stats {
2040 let mut pass = Pass::new(ty, 1);
2041 for (at, stripe) in held.chunks(5).enumerate() {
2042 pass.open_stripe((at as u64, 0));
2043 for vector in stripe {
2044 scan(&mut pass, vector);
2045 }
2046 pass.close_stripe();
2047 }
2048 pass.finish(Sketch::of(&[]), Vec::new())
2049 }
2050
2051 /// A one column table of intervals, which is a type with no hash rule and so a table this
2052 /// build writes no statistics section for.
2053 ///
2054 /// The only way left to make a file whose table names no sections, now that an ordinary write
2055 /// writes them. See the criterion 3 test for why stamping the version back onto a file that has
2056 /// them does not do it.
2057 fn table_of_intervals(label: &str, months: &[i32]) -> PathBuf {
2058 let path = path(label);
2059 let mut writer =
2060 Writer::create(&path, "t", vec![Field::new("v", LogicalType::Interval)]).expect("new");
2061 for part in months.chunks(1000) {
2062 let held = part
2063 .iter()
2064 .map(|months| Value::Interval { months: *months, days: 0, micros: 0 })
2065 .collect::<Vec<_>>();
2066 let chunk = Chunk::new(vec![
2067 Vector::from_values(LogicalType::Interval, &held).expect("values"),
2068 ])
2069 .expect("one column");
2070 writer.append(&chunk).expect("a part");
2071 }
2072 writer.finish().expect("commit");
2073 path
2074 }
2075
2076 /// Every value of the one column, in rid order, which is what a scan of this table answers.
2077 fn rows_of(reader: &Reader) -> Vec<Value> {
2078 let mut out = Vec::new();
2079 for part in 0..reader.parts() {
2080 let chunk = reader.read(part, &[0]).expect("a part reads back");
2081 for row in 0..chunk.len() {
2082 out.push(chunk.value_at(0, row));
2083 }
2084 }
2085 out
2086 }
2087
2088 fn reopen(path: &PathBuf) -> Reader {
2089 Catalog::open(path).expect("reopen").table("t").expect("the table")
2090 }
2091
2092 #[test]
2093 fn a_summary_built_over_a_file_says_what_the_column_holds() {
2094 // End to end: the column goes to disk, comes back through the reader, and every field of
2095 // the summary is the truth about it. Three thousand rows so the scan crosses parts, because
2096 // a pass that read them in the wrong order would be right about one part and wrong about
2097 // the order fields for the rest.
2098 let values = (1..=3000_i64).map(Some).collect::<Vec<_>>();
2099 let path = table_of("sorted", &values);
2100 let built = build_stats(&path, "t", &[0]).expect("build");
2101 assert_eq!(built.len(), 1);
2102 assert!(built[0].built, "a one column table is nowhere near the budget");
2103 assert_eq!(built[0].rows, 3000);
2104 assert_eq!(built[0].distinct, 3000);
2105 assert!(built[0].exact, "three thousand values is under the default k");
2106 assert_eq!(built[0].order, Order::Ascending);
2107
2108 let reader = reopen(&path);
2109 let summary = summary(&reader, 0).expect("the summary is in the file");
2110 assert_eq!(summary.rows, 3000);
2111 assert_eq!(summary.nulls, 0);
2112 assert_eq!(summary.low, Some(Bound::Int(1)));
2113 assert_eq!(summary.high, Some(Bound::Int(3000)));
2114 assert!(summary.ends_exact);
2115 assert!(summary.unique, "a sorted run of distinct values is a key candidate");
2116 assert_eq!(summary.runs, 1, "one ascending run");
2117 assert_eq!(summary.distinct_class, Class::Exact);
2118 assert_eq!(summary.newest, reader.table().generation());
2119
2120 let sketches = sketches(&reader, 0).expect("the sketches are in the file");
2121 assert!(sketches.merged.is_exact());
2122 assert!(sketches.stripes.is_empty(), "the per stripe rule gives this column none");
2123
2124 fs::remove_file(&path).expect("clean up");
2125 }
2126
2127 #[test]
2128 fn nulls_are_counted_and_do_not_reach_the_ends_or_the_sketch() {
2129 // The distinction that costs an answer if it is got wrong. A null is a row and is not a
2130 // value, so it moves `rows` and `nulls` and moves nothing else.
2131 let values: Vec<Option<i64>> =
2132 (0..2000).map(|at| if at % 3 == 0 { None } else { Some(at) }).collect();
2133 let path = table_of("nulls", &values);
2134 build_stats(&path, "t", &[0]).expect("build");
2135
2136 let reader = reopen(&path);
2137 let summary = summary(&reader, 0).expect("the summary");
2138 let nulls = values.iter().filter(|v| v.is_none()).count() as u64;
2139 assert_eq!(summary.rows, 2000);
2140 assert_eq!(summary.nulls, nulls);
2141 assert_eq!(summary.present(), 2000 - nulls);
2142 assert_eq!(summary.distinct, 2000 - nulls, "a null is not a distinct value");
2143 assert_eq!(summary.low, Some(Bound::Int(1)), "zero is null here");
2144 assert!(summary.unique);
2145
2146 fs::remove_file(&path).expect("clean up");
2147 }
2148
2149 #[test]
2150 fn a_column_that_repeats_is_not_reported_unique_and_a_descending_one_is_seen() {
2151 let values = (0..2000_i64).map(|at| Some(-(at / 2))).collect::<Vec<_>>();
2152 let path = table_of("repeats", &values);
2153 build_stats(&path, "t", &[0]).expect("build");
2154
2155 let reader = reopen(&path);
2156 let summary = summary(&reader, 0).expect("the summary");
2157 assert_eq!(summary.distinct, 1000);
2158 assert!(!summary.unique, "every value appears twice");
2159 assert_eq!(summary.order, Order::Descending);
2160 assert_eq!(summary.runs, 1000, "a descending column is a run per distinct value");
2161
2162 fs::remove_file(&path).expect("clean up");
2163 }
2164
2165 #[test]
2166 fn a_column_past_the_default_k_is_estimated_and_says_so() {
2167 // The rule the module doc names, at the point where it bites. Past k the sketch threw values
2168 // away, so the count is an estimate, and the class has to say so or a COUNT(DISTINCT) is
2169 // answered out of metadata with a number that is close and wrong.
2170 let values = (0..20_000_i64).map(Some).collect::<Vec<_>>();
2171 let path = table_of("estimated", &values);
2172 let built = build_stats(&path, "t", &[0]).expect("build");
2173 assert!(!built[0].exact, "twenty thousand values is past the default k");
2174
2175 let reader = reopen(&path);
2176 let summary = summary(&reader, 0).expect("the summary");
2177 assert_eq!(summary.distinct_class, Class::Estimated);
2178 assert!(!summary.unique, "uniqueness is never claimed off an estimate");
2179 assert!(summary.distinct > 17_000 && summary.distinct <= 20_000, "{}", summary.distinct);
2180 assert!(summary.distinct <= summary.present(), "more distinct values than rows");
2181
2182 fs::remove_file(&path).expect("clean up");
2183 }
2184
2185 #[test]
2186 fn a_shuffled_column_is_neither_ordered_nor_one_run() {
2187 let values = (0..2000_i64).map(|at| Some((at * 7919) % 2000)).collect::<Vec<_>>();
2188 let path = table_of("shuffled", &values);
2189 build_stats(&path, "t", &[0]).expect("build");
2190
2191 let reader = reopen(&path);
2192 let summary = summary(&reader, 0).expect("the summary");
2193 assert_eq!(summary.order, Order::Neither);
2194 assert!(summary.runs > 100, "a shuffle is many runs, not one: {}", summary.runs);
2195 assert_eq!(summary.low, Some(Bound::Int(0)));
2196 assert_eq!(summary.high, Some(Bound::Int(1999)));
2197
2198 fs::remove_file(&path).expect("clean up");
2199 }
2200
2201 #[test]
2202 fn a_column_something_declared_a_key_over_is_sketched_per_stripe_and_a_plain_one_is_not() {
2203 // Section 3.8's rule, both halves of it. Nothing has declared anything over this column, so
2204 // the first build gives it the table level summary and no per stripe sketches, which is the
2205 // state most columns are in and is what keeps SF100 inside two percent. A key map is then
2206 // built over it, which is something declaring it a key, and the next build promotes it.
2207 let values = (1..=19_200_i64).map(Some).collect::<Vec<_>>();
2208 let path = table_of_parts("promoted", &values, 100);
2209
2210 let plain = build_stats(&path, "t", &[0]).expect("build");
2211 assert_eq!(plain[0].stripes, 0, "nothing has declared anything over this column yet");
2212
2213 crate::graph::build_key_maps(&path, "t", &[0]).expect("a key map declares it a key");
2214 let promoted = build_stats(&path, "t", &[0]).expect("rebuild");
2215 assert!(promoted[0].stripes > 1, "{} stripes, wanted more than one", promoted[0].stripes);
2216 assert!(promoted[0].built, "and they fit");
2217 // The equality rather than a tolerance. The merged sketch of a promoted column is the union
2218 // of its stripe sketches at the column's own k, and a union of bottom-k sketches at one k
2219 // is the bottom-k of everything they saw, so it holds the same hashes as the single sketch
2220 // the plain build made. Promotion changes where the counting is reset and nothing else.
2221 assert_eq!(promoted[0].distinct, plain[0].distinct, "the merged count did not move");
2222
2223 let reader = reopen(&path);
2224 let sketches = sketches(&reader, 0).expect("the sketches came back");
2225 assert_eq!(sketches.stripes.len(), promoted[0].stripes);
2226 assert!(
2227 sketches.stripes.iter().all(|stripe| stripe.k() == STRIPE_K),
2228 "a stripe sketch is written down at the smaller k"
2229 );
2230 let floor = sketches.floor(0, sketches.stripes.len()).expect("a floor over every stripe");
2231 let actual = 19_200.0;
2232 assert!(
2233 (floor - actual).abs() / actual < 0.25,
2234 "{floor:.0} over every stripe against {actual:.0}"
2235 );
2236
2237 drop(reader);
2238 fs::remove_file(&path).expect("clean up");
2239 }
2240
2241 #[test]
2242 fn a_file_from_before_the_section_table_opens_and_every_statistic_is_unknown() {
2243 // Exit criterion 3 of #762, the statistics half of it. A build that knows about summaries
2244 // opens a file written by a build that did not, with no rewrite and no repair, states
2245 // nothing about that file's columns, and reads back exactly what the same rows read back
2246 // out of a file this build wrote.
2247 //
2248 // `None` is what `Unknown` is at this layer, and the two readers answer it for every reason
2249 // there is rather than distinguishing them, which is section 3.1: there is nothing a caller
2250 // could do differently on hearing *the file predates statistics* rather than *the section
2251 // does not checksum*, because both are answered by planning the query the way it was
2252 // planned before statistics existed.
2253 //
2254 // The older file is a table of a type with no hash rule, with its version stamped back. A
2255 // build before section 3.8 wrote no section block at all, and a table this build writes no
2256 // sections for is that file on disk, so there is no fixture to go stale and no second
2257 // encoder to drift.
2258 //
2259 // The obvious construction, stamping the version back onto a file that does carry
2260 // summaries, does not work and is worth saying why. The section block is found by a magic
2261 // at the end of the directory rather than by the number in the header, so a stamped file
2262 // with sections in it is a file with sections in it, and the test would be asserting
2263 // nothing.
2264 let months = (1..=3000_i32).collect::<Vec<_>>();
2265 let older = table_of_intervals("before_sections", &months);
2266 let current = table_of("with_sections", &(1..=3000_i64).map(Some).collect::<Vec<_>>());
2267
2268 let file = fs::OpenOptions::new().write(true).open(&older).expect("reopen to patch");
2269 crate::write_at(&file, 8, &22_u32.to_le_bytes()).expect("stamp the older format");
2270 drop(file);
2271
2272 let new = reopen(¤t);
2273 assert!(summary(&new, 0).is_some(), "the file this build wrote says what it holds");
2274
2275 let old = reopen(&older);
2276 assert!(old.table().sections().is_empty(), "an older file names no sections");
2277 assert!(summary(&old, 0).is_none(), "and so says nothing about its columns");
2278 assert!(sketches(&old, 0).is_none());
2279 assert!(read_columns(&old).is_empty(), "nor promotes any of them");
2280 assert_eq!(old.table().rows(), 3000, "and reads every row it holds");
2281 assert_eq!(
2282 rows_of(&old).first(),
2283 Some(&Value::Interval { months: 1, days: 0, micros: 0 }),
2284 "with the values it was written with"
2285 );
2286
2287 drop(new);
2288 drop(old);
2289 fs::remove_file(¤t).expect("clean up");
2290 fs::remove_file(&older).expect("clean up");
2291 }
2292
2293 #[test]
2294 fn the_stripe_ends_say_whether_a_scan_can_skip_and_a_shuffle_says_it_cannot() {
2295 // The per stripe ends, which is the one thing the pass tracks that nothing else checks and
2296 // which a scan reads to skip a whole stripe. A sorted column's stripes do not overlap and a
2297 // shuffled column's every stripe spans the column, so the same rows in a different order
2298 // give the opposite answer. Three stripes, so that the ends are opened and closed more than
2299 // once and a pass that never reset them would be caught.
2300 let sorted = (1..=19_200_i64).map(Some).collect::<Vec<_>>();
2301 let ordered = table_of_parts("stripes_sorted", &sorted, 100);
2302 build_stats(&ordered, "t", &[0]).expect("build");
2303 let reader = reopen(&ordered);
2304 let ordered_summary = summary(&reader, 0).expect("the summary");
2305 assert!(!ordered_summary.overlapping, "a sorted column's stripes are disjoint");
2306 assert_eq!(ordered_summary.low, Some(Bound::Int(1)));
2307 assert_eq!(ordered_summary.high, Some(Bound::Int(19_200)));
2308 drop(reader);
2309
2310 // A fixed stride rather than a random shuffle, so a failure is the same failure twice. The
2311 // stride and the row count share no factor, so this visits every value exactly once and
2312 // every stripe ends up holding values from very nearly the whole range.
2313 let shuffled = (0..19_200_i64).map(|at| Some(1 + at * 7919 % 19_200)).collect::<Vec<_>>();
2314 let mixed = table_of_parts("stripes_shuffled", &shuffled, 100);
2315 build_stats(&mixed, "t", &[0]).expect("build");
2316 let reader = reopen(&mixed);
2317 let mixed_summary = summary(&reader, 0).expect("the summary");
2318 assert!(mixed_summary.overlapping, "a shuffled column's stripes all span it");
2319 assert_eq!(mixed_summary.low, Some(Bound::Int(1)), "the same values in a different order");
2320 assert_eq!(mixed_summary.high, Some(Bound::Int(19_200)));
2321 drop(reader);
2322
2323 fs::remove_file(&ordered).expect("clean up");
2324 fs::remove_file(&mixed).expect("clean up");
2325 }
2326
2327 #[test]
2328 fn the_graph_sections_do_not_count_against_the_statistics_budget() {
2329 // The direction of box 4 that costs more, because the two percent is the smaller share. A
2330 // TPC-H SF10 file's key maps are 7.7 MB against an allowance of 54 MB, so a statistics
2331 // build that counted them would start a seventh of the way through a budget it was given
2332 // all of, and columns at the far end of a wide table would go unsummarized for a reason
2333 // that has nothing to do with summaries.
2334 let values = (1..=3000_i64).map(Some).collect::<Vec<_>>();
2335 let path = table_of("apart", &values);
2336 crate::graph::build_key_maps(&path, "t", &[0]).expect("a key map first");
2337
2338 let reader = reopen(&path);
2339 let graph = reader
2340 .table()
2341 .sections()
2342 .iter()
2343 .filter(|held| held.among(section::GRAPH_KINDS))
2344 .count();
2345 assert_eq!(graph, 1, "the key map is in the file");
2346 assert_eq!(held_bytes(&reader, &[0]).expect("held"), 0, "and it is not the statistics'");
2347
2348 drop(reader);
2349 fs::remove_file(&path).expect("clean up");
2350 }
2351
2352 #[test]
2353 fn deleting_the_sections_changes_nothing_but_whether_they_are_there() {
2354 // Section 3.1, as close to directly as a test can put it. The same file, read once with the
2355 // sections and once with the generation moved past them, and the reader opens and scans the
2356 // same either way.
2357 let values = (1..=1500_i64).map(Some).collect::<Vec<_>>();
2358 let path = table_of("invariant", &values);
2359 build_stats(&path, "t", &[0]).expect("build");
2360
2361 let reader = reopen(&path);
2362 assert!(summary(&reader, 0).is_some());
2363 let generation = reader.table().generation();
2364 let held: Vec<_> = reader
2365 .table()
2366 .sections()
2367 .iter()
2368 .filter(|s| s.kind == *section::SUMMARY || s.kind == *section::SKETCHES)
2369 .copied()
2370 .collect();
2371 assert_eq!(held.len(), 2, "a summary and a sketch section");
2372 for section in &held {
2373 assert!(section.usable(generation));
2374 assert!(!section.usable(generation + 1), "a rewrite invalidates rather than corrupts");
2375 }
2376 let rows: usize =
2377 (0..reader.parts()).map(|part| reader.read(part, &[0]).expect("a part").len()).sum();
2378 assert_eq!(rows, 1500, "the scan is the scan whether the sections are read or not");
2379
2380 fs::remove_file(&path).expect("clean up");
2381 }
2382
2383 #[test]
2384 fn a_string_column_is_read_through_the_typed_path_and_measured_by_its_bytes() {
2385 // The other fast path. A varchar has no fixed width, so the byte total and the widest value
2386 // are measured per value, and the ends are the string ends rather than the hash ends.
2387 let path = path("strings");
2388 let mut writer =
2389 Writer::create(&path, "t", vec![Field::new("v", LogicalType::Varchar)]).expect("new");
2390 let words = ["alpha", "bravo", "charlie", "delta", "alpha"];
2391 let held = words.iter().map(|w| Value::Varchar((*w).into())).collect::<Vec<_>>();
2392 let chunk =
2393 Chunk::new(vec![Vector::from_values(LogicalType::Varchar, &held).expect("words")])
2394 .expect("one column");
2395 writer.append(&chunk).expect("a part");
2396 writer.finish().expect("commit");
2397 build_stats(&path, "t", &[0]).expect("build");
2398
2399 let reader = reopen(&path);
2400 let summary = summary(&reader, 0).expect("the summary");
2401 assert_eq!(summary.rows, 5);
2402 assert_eq!(summary.distinct, 4, "alpha twice");
2403 assert!(!summary.unique);
2404 assert_eq!(summary.bytes, words.iter().map(|w| w.len() as u64).sum::<u64>());
2405 assert_eq!(summary.widest, 7, "charlie");
2406 assert_eq!(summary.low, Some(Bound::Bytes(b"alpha".to_vec())));
2407 assert_eq!(summary.high, Some(Bound::Bytes(b"delta".to_vec())));
2408
2409 drop(reader);
2410 fs::remove_file(&path).expect("clean up");
2411 }
2412
2413 #[test]
2414 fn a_type_with_no_hash_rule_is_refused_by_name_rather_than_summarized_as_empty() {
2415 let path = table_of("refused", &[Some(1)]);
2416 let reader = reopen(&path);
2417 assert!(summarizable(&LogicalType::BigInt));
2418 assert!(!summarizable(&LogicalType::Interval));
2419 assert!(build_summary(&reader, 1).is_err(), "a column past the end");
2420 drop(reader);
2421 fs::remove_file(&path).expect("clean up");
2422 }
2423
2424 #[test]
2425 fn the_stored_sketch_depends_on_the_value_rule_and_not_only_on_the_hash() {
2426 // HASH_IDENTITY pins `hash64`, which is half of what a stored sketch depends on. The other
2427 // half is the rule that turns a value into the bytes `hash64` sees, and that rule lives in
2428 // `rudb_storage::count`. Changing it without bumping HASH_IDENTITY would leave every stored
2429 // sketch readable, accepted, and built over a different universe than the one a new sketch
2430 // is built over, which is exactly the merge the identity exists to prevent.
2431 //
2432 // So the rule is pinned here. If this fails because `hash_value` changed on purpose, the fix
2433 // is to bump HASH_IDENTITY and then update these numbers, in that order.
2434 assert_eq!(hash_value(&Value::BigInt(1)), Some(hash64(&1_u128.to_le_bytes())));
2435 assert_eq!(hash_value(&Value::Integer(1)), hash_value(&Value::BigInt(1)));
2436 assert_eq!(hash_value(&Value::Varchar("a".into())), Some(hash64(b"a")));
2437 assert_eq!(hash_value(&Value::Null), None);
2438 }
2439}