Skip to main content

rudb_native/
projection.rs

1//! Optional row-valued projection ordered by one signed integer column.
2//!
3//! The first format covers a second signed integer column with a small dictionary. It carries
4//! every row in order, never a grouped result. A table append changes the table generation and
5//! makes the attached section unusable until it is rebuilt.
6
7use std::collections::{HashMap, HashSet};
8use std::path::Path;
9
10use rudb_common::{LogicalType, Result};
11use rudb_vector::Vector;
12
13use crate::{Catalog, Reader, attach, invalid, section};
14
15const MAGIC: &[u8; 8] = b"RUDBSP1\0";
16const FIXED_HEADER: usize = 8 + 2 + 2 + 8 + 2;
17const ROW_BYTES: usize = 10;
18
19pub(crate) fn id(order: usize, covered: usize) -> Result<u64> {
20    let order =
21        u32::try_from(order).map_err(|_| invalid("projection order column is too large"))?;
22    let covered =
23        u32::try_from(covered).map_err(|_| invalid("projection covered column is too large"))?;
24    Ok((u64::from(order) << 32) | u64::from(covered))
25}
26
27/// Every value of a non-null signed integer column, as a block where the vector hands one over and
28/// through `signed_at` for the forms it does not.
29pub(crate) fn integers(vector: &Vector, out: &mut Vec<i64>) -> Result<()> {
30    if vector.signed_block(out) && out.len() == vector.len() {
31        return Ok(());
32    }
33    out.clear();
34    // row at a time: the run and compressed forms have no block to hand over.
35    for row in 0..vector.len() {
36        let value = vector
37            .signed_at(row)
38            .and_then(|value| i64::try_from(value).ok())
39            .ok_or_else(|| invalid("sorted projection requires non-null signed integers"))?;
40        out.push(value);
41    }
42    Ok(())
43}
44
45pub(crate) fn eligible(order: &LogicalType, covered: &LogicalType) -> bool {
46    matches!(
47        order,
48        LogicalType::TinyInt | LogicalType::SmallInt | LogicalType::Integer | LogicalType::BigInt
49    ) && matches!(covered, LogicalType::TinyInt | LogicalType::SmallInt | LogicalType::Integer)
50}
51
52/// Build a row-valued, sorted, covering projection inside an existing native file.
53///
54/// The caller chooses columns by name. This is an explicit storage operation whose build time,
55/// peak memory, and bytes must be charged to the workload that asks for it. No SQL aggregate is
56/// evaluated or stored. A later append makes the section stale until this is called again.
57///
58/// # Errors
59///
60/// If the table or columns are missing, nullable, unsupported, or the file cannot be read or
61/// updated. More than 65,535 covered values exceed this first format's dictionary.
62pub fn build_sorted_projection(
63    path: impl AsRef<Path>,
64    table: &str,
65    order_column: &str,
66    covered_column: &str,
67) -> Result<()> {
68    let path = path.as_ref();
69    let catalog = Catalog::open(path)?;
70    let reader = catalog.table(table)?;
71    let fields = reader.table().fields();
72    let order = fields
73        .iter()
74        .position(|field| field.name.eq_ignore_ascii_case(order_column))
75        .ok_or_else(|| invalid("projection order column is missing"))?;
76    let covered = fields
77        .iter()
78        .position(|field| field.name.eq_ignore_ascii_case(covered_column))
79        .ok_or_else(|| invalid("projection covered column is missing"))?;
80    if order == covered
81        || !fields[order].not_null
82        || !fields[covered].not_null
83        || !eligible(&fields[order].ty, &fields[covered].ty)
84    {
85        return Err(invalid("sorted projection needs two supported, non-null integer columns"));
86    }
87    let mut rows = Vec::<(i64, i32)>::with_capacity(reader.table().rows());
88    let mut values = HashSet::<i32>::new();
89    let (mut users, mut groups) = (Vec::new(), Vec::new());
90    for part in 0..reader.parts() {
91        let chunk = reader.read(part, &[order, covered])?;
92        integers(chunk.column(0)?, &mut users)?;
93        integers(chunk.column(1)?, &mut groups)?;
94        for (&user, &group) in users.iter().zip(&groups) {
95            let group = i32::try_from(group)
96                .map_err(|_| invalid("projection covered value exceeds INTEGER"))?;
97            rows.push((user, group));
98            values.insert(group);
99        }
100    }
101    if rows.len() != reader.table().rows() {
102        return Err(invalid("projection source row count differs from its table"));
103    }
104    let mut dictionary = values.into_iter().collect::<Vec<_>>();
105    dictionary.sort_unstable();
106    let dictionary_len = u16::try_from(dictionary.len())
107        .map_err(|_| invalid("projection covered column exceeds 65535 values"))?;
108    let codes = dictionary
109        .iter()
110        .enumerate()
111        .map(|(at, &value)| (value, at as u16))
112        .collect::<HashMap<_, _>>();
113    rows.sort_unstable_by_key(|&(user, _)| user);
114    let order_index =
115        u16::try_from(order).map_err(|_| invalid("projection column index overflow"))?;
116    let covered_index =
117        u16::try_from(covered).map_err(|_| invalid("projection column index overflow"))?;
118    let header = FIXED_HEADER + dictionary.len() * 4;
119    let mut bytes = Vec::with_capacity(header + rows.len() * ROW_BYTES);
120    bytes.extend_from_slice(MAGIC);
121    bytes.extend_from_slice(&order_index.to_le_bytes());
122    bytes.extend_from_slice(&covered_index.to_le_bytes());
123    bytes.extend_from_slice(&(rows.len() as u64).to_le_bytes());
124    bytes.extend_from_slice(&dictionary_len.to_le_bytes());
125    for value in &dictionary {
126        bytes.extend_from_slice(&value.to_le_bytes());
127    }
128    for (user, group) in rows {
129        bytes.extend_from_slice(&user.to_le_bytes());
130        bytes.extend_from_slice(&codes[&group].to_le_bytes());
131    }
132    drop(reader);
133    drop(catalog);
134    attach(
135        path,
136        table,
137        &[section::Attachment {
138            kind: *section::SORTED_PROJECTION,
139            id: id(order, covered)?,
140            flags: 0,
141            header_bytes: header as u32,
142            bytes: &bytes,
143        }],
144    )?;
145    Ok(())
146}
147
148impl Reader {
149    /// Count distinct order values by covered value from a current sorted projection.
150    ///
151    /// Returns `None` when no matching current projection exists. Workers split only between
152    /// order values, so each `(order, covered)` pair is counted once without a global hash set.
153    ///
154    /// # Errors
155    ///
156    /// If a matching projection or its source directory is damaged.
157    ///
158    /// # Panics
159    ///
160    /// Fixed-width header decoding assumes the header length checked immediately before it.
161    pub fn grouped_distinct_projection(
162        &self,
163        order: usize,
164        covered: usize,
165        limit: usize,
166    ) -> Result<Option<Vec<(i32, u64)>>> {
167        let wanted = id(order, covered)?;
168        let Some(section) = self.table().sections().iter().find(|section| {
169            section.kind == *section::SORTED_PROJECTION
170                && section.id == wanted
171                && section.usable(self.table().generation())
172        }) else {
173            return Ok(None);
174        };
175        let extents = self.extents(section)?;
176        let first = extents.first().ok_or_else(|| invalid("projection has no extent"))?;
177        let bytes = self.extent(first)?;
178        if bytes.len() < FIXED_HEADER || &bytes[..8] != MAGIC {
179            return Err(invalid("projection header differs"));
180        }
181        let stored_order = u16::from_le_bytes(bytes[8..10].try_into().unwrap());
182        let stored_covered = u16::from_le_bytes(bytes[10..12].try_into().unwrap());
183        if usize::from(stored_order) != order || usize::from(stored_covered) != covered {
184            return Err(invalid("projection columns differ from its section"));
185        }
186        let rows = u64::from_le_bytes(bytes[12..20].try_into().unwrap());
187        if rows != self.table().rows() as u64 {
188            return Err(invalid("projection row count differs from its table"));
189        }
190        let size = u16::from_le_bytes(bytes[20..22].try_into().unwrap()) as usize;
191        let header = FIXED_HEADER + size * 4;
192        if bytes.len() < header || section.header_bytes as usize != header {
193            return Err(invalid("projection dictionary exceeds its first extent"));
194        }
195        let dictionary = bytes[FIXED_HEADER..header]
196            .chunks_exact(4)
197            .map(|item| i32::from_le_bytes(item.try_into().unwrap()))
198            .collect::<Vec<_>>();
199        if dictionary.windows(2).any(|pair| pair[0] >= pair[1]) {
200            return Err(invalid("projection dictionary is not sorted and unique"));
201        }
202        let expected = (header as u64)
203            .checked_add(
204                rows.checked_mul(ROW_BYTES as u64)
205                    .ok_or_else(|| invalid("projection length overflow"))?,
206            )
207            .ok_or_else(|| invalid("projection length overflow"))?;
208        let base = first.offset;
209        let mut consumed = 0_u64;
210        for extent in &extents {
211            if extent.first != consumed || extent.offset != base + consumed {
212                return Err(invalid("projection extents are not contiguous"));
213            }
214            consumed += u64::from(extent.length);
215        }
216        if consumed != expected {
217            return Err(invalid("projection byte length differs from its row count"));
218        }
219        let rows = usize::try_from(rows).map_err(|_| invalid("projection rows exceed memory"))?;
220        let workers = if rows < 2_000_000 {
221            1
222        } else {
223            std::thread::available_parallelism().map_or(1, usize::from).min(8).min(rows)
224        };
225        let mut boundaries = Vec::with_capacity(workers + 1);
226        boundaries.push(0);
227        for worker in 1..workers {
228            let mut at = rows * worker / workers;
229            if at > 0 && at < rows {
230                let previous = projected_user(self, base, header, at - 1)?;
231                while at < rows && projected_user(self, base, header, at)? == previous {
232                    at += 1;
233                }
234            }
235            boundaries.push(at);
236        }
237        boundaries.push(rows);
238        let counts = std::thread::scope(|scope| -> Result<Vec<Vec<u64>>> {
239            let mut handles = Vec::with_capacity(workers);
240            for pair in boundaries.windows(2) {
241                let (start, end) = (pair[0], pair[1]);
242                let extents = &extents;
243                handles
244                    .push(scope.spawn(move || scan_range(self, extents, header, size, start, end)));
245            }
246            handles
247                .into_iter()
248                .map(|handle| handle.join().map_err(|_| invalid("projection worker panicked"))?)
249                .collect()
250        })?;
251        let mut totals = vec![0_u64; size];
252        for local in counts {
253            for (total, value) in totals.iter_mut().zip(local) {
254                *total =
255                    total.checked_add(value).ok_or_else(|| invalid("projection count overflow"))?;
256            }
257        }
258        let mut ranked =
259            dictionary.into_iter().zip(totals).filter(|(_, count)| *count != 0).collect::<Vec<_>>();
260        ranked.sort_unstable_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
261        ranked.truncate(limit);
262        Ok(Some(ranked))
263    }
264}
265
266fn projected_user(reader: &Reader, base: u64, header: usize, row: usize) -> Result<i64> {
267    let mut bytes = [0_u8; 8];
268    let offset = base + header as u64 + row as u64 * ROW_BYTES as u64;
269    crate::read_at(&reader.file, offset, &mut bytes)?;
270    Ok(i64::from_le_bytes(bytes))
271}
272
273fn scan_range(
274    reader: &Reader,
275    extents: &[section::Extent],
276    header: usize,
277    dictionary: usize,
278    start: usize,
279    end: usize,
280) -> Result<Vec<u64>> {
281    let low = header as u64 + start as u64 * ROW_BYTES as u64;
282    let high = header as u64 + end as u64 * ROW_BYTES as u64;
283    let mut marks = vec![0_u64; dictionary];
284    let mut counts = vec![0_u64; dictionary];
285    let mut current_user = None;
286    let mut epoch = 0_u64;
287    let mut seen = 0_usize;
288    let mut carry = [0_u8; ROW_BYTES];
289    let mut carry_len = 0_usize;
290    for extent in extents {
291        let extent_end = extent.first + u64::from(extent.length);
292        if extent_end <= low || extent.first >= high {
293            continue;
294        }
295        let bytes = reader.extent(extent)?;
296        let begin = low.saturating_sub(extent.first) as usize;
297        let finish = (high.min(extent_end) - extent.first) as usize;
298        let mut block = &bytes[begin..finish];
299        if carry_len != 0 {
300            let needed = ROW_BYTES - carry_len;
301            let taken = needed.min(block.len());
302            carry[carry_len..carry_len + taken].copy_from_slice(&block[..taken]);
303            carry_len += taken;
304            block = &block[taken..];
305            if carry_len < ROW_BYTES {
306                continue;
307            }
308            process_block(
309                &carry,
310                dictionary,
311                &mut current_user,
312                &mut epoch,
313                &mut marks,
314                &mut counts,
315            )?;
316            seen += 1;
317        }
318        let chunks = block.chunks_exact(ROW_BYTES);
319        let remainder = chunks.remainder();
320        for chunk in chunks {
321            process_block(
322                chunk,
323                dictionary,
324                &mut current_user,
325                &mut epoch,
326                &mut marks,
327                &mut counts,
328            )?;
329            seen += 1;
330        }
331        carry[..remainder.len()].copy_from_slice(remainder);
332        carry_len = remainder.len();
333    }
334    if carry_len != 0 || seen != end - start {
335        return Err(invalid("projection scan did not cover its range"));
336    }
337    Ok(counts)
338}
339
340#[inline(always)]
341fn process_block(
342    bytes: &[u8],
343    dictionary: usize,
344    current_user: &mut Option<i64>,
345    epoch: &mut u64,
346    marks: &mut [u64],
347    counts: &mut [u64],
348) -> Result<()> {
349    let user = i64::from_le_bytes(bytes[..8].try_into().unwrap());
350    let code = u16::from_le_bytes(bytes[8..10].try_into().unwrap()) as usize;
351    if code >= dictionary {
352        return Err(invalid("projection code is outside its dictionary"));
353    }
354    if current_user.is_some_and(|previous| user < previous) {
355        return Err(invalid("projection order is descending"));
356    }
357    if *current_user != Some(user) {
358        *epoch = epoch.checked_add(1).ok_or_else(|| invalid("projection epoch overflow"))?;
359        *current_user = Some(user);
360    }
361    if marks[code] != *epoch {
362        marks[code] = *epoch;
363        counts[code] += 1;
364    }
365    Ok(())
366}
367
368#[cfg(test)]
369mod tests {
370    use std::collections::{HashMap, HashSet};
371    use std::sync::atomic::{AtomicUsize, Ordering};
372
373    use rudb_common::{Field, LogicalType, Value};
374    use rudb_vector::{Chunk, Vector};
375
376    use crate::{Catalog, Writer};
377
378    use super::build_sorted_projection;
379
380    static NEXT: AtomicUsize = AtomicUsize::new(0);
381
382    #[test]
383    fn sorted_projection_counts_each_cover_value_once_per_order_value() {
384        let at = NEXT.fetch_add(1, Ordering::Relaxed);
385        let path = std::env::temp_dir()
386            .join(format!("rudb-sorted-projection-{}-{at}.rdb", std::process::id()));
387        let fields = vec![
388            Field::required("user", LogicalType::BigInt),
389            Field::required("region", LogicalType::Integer),
390        ];
391        let mut writer = Writer::create(&path, "events", fields).expect("create native file");
392        for (users, regions) in
393            [(vec![9_i64, 2, 9], vec![7_i32, 1, 7]), (vec![2_i64, 2, 5], vec![2_i32, 2, 1])]
394        {
395            let users = users.into_iter().map(Value::BigInt).collect::<Vec<_>>();
396            let regions = regions.into_iter().map(Value::Integer).collect::<Vec<_>>();
397            let chunk = Chunk::new(vec![
398                Vector::from_values(LogicalType::BigInt, &users).expect("users"),
399                Vector::from_values(LogicalType::Integer, &regions).expect("regions"),
400            ])
401            .expect("matching columns");
402            writer.append(&chunk).expect("append rows");
403        }
404        writer.finish().expect("commit native file");
405        let reader = Catalog::open(&path).expect("catalog").table("events").expect("table");
406        assert_eq!(reader.grouped_distinct_projection(0, 1, 10).expect("no projection"), None);
407        drop(reader);
408        build_sorted_projection(&path, "events", "user", "region").expect("build projection");
409        let reader = Catalog::open(&path).expect("catalog").table("events").expect("table");
410        assert_eq!(
411            reader.grouped_distinct_projection(0, 1, 10).expect("valid projection"),
412            Some(vec![(1, 2), (2, 1), (7, 1)]),
413        );
414        assert_eq!(reader.grouped_distinct_projection(1, 0, 10).expect("different order"), None);
415        std::fs::remove_file(path).expect("remove scratch file");
416    }
417
418    #[test]
419    fn sorted_projection_crosses_extent_and_worker_boundaries() {
420        let at = NEXT.fetch_add(1, Ordering::Relaxed);
421        let path = std::env::temp_dir()
422            .join(format!("rudb-sorted-projection-wide-{}-{at}.rdb", std::process::id()));
423        let fields = vec![
424            Field::required("user", LogicalType::BigInt),
425            Field::required("region", LogicalType::Integer),
426        ];
427        let mut writer = Writer::create(&path, "events", fields).expect("create native file");
428        let mut pairs = HashSet::new();
429        for first in (0..75_000_i64).step_by(1_000) {
430            let source = (first..first + 1_000)
431                .map(|row| (row * 17 % 20_000, (row % 7) as i32))
432                .collect::<Vec<_>>();
433            pairs.extend(source.iter().copied());
434            let users = source.iter().map(|(user, _)| Value::BigInt(*user)).collect::<Vec<_>>();
435            let regions =
436                source.iter().map(|(_, region)| Value::Integer(*region)).collect::<Vec<_>>();
437            let chunk = Chunk::new(vec![
438                Vector::from_values(LogicalType::BigInt, &users).expect("users"),
439                Vector::from_values(LogicalType::Integer, &regions).expect("regions"),
440            ])
441            .expect("matching columns");
442            writer.append(&chunk).expect("append rows");
443        }
444        writer.finish().expect("commit native file");
445        build_sorted_projection(&path, "events", "user", "region").expect("build projection");
446        let reader = Catalog::open(&path).expect("catalog").table("events").expect("table");
447        let mut expected = HashMap::<i32, u64>::new();
448        for (_, region) in pairs {
449            *expected.entry(region).or_default() += 1;
450        }
451        let mut expected = expected.into_iter().collect::<Vec<_>>();
452        expected.sort_unstable_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
453        assert_eq!(
454            reader.grouped_distinct_projection(0, 1, 10).expect("valid projection"),
455            Some(expected)
456        );
457        std::fs::remove_file(path).expect("remove scratch file");
458    }
459}