1use 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
27pub(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 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
52pub 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 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, ®ions).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, ®ions).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}