use std::collections::HashMap;
use std::ops::Deref;
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering as Atomic};
use std::sync::{Arc, Mutex};
use rudb_common::{Error, LogicalType, Result};
use rudb_metrics::{LoadProfile, Stage};
use rudb_storage::Range;
use rudb_vector::{Bitmap, Chunk, Data, StringColumn, Validity, Vector};
use super::{
ColumnStripe, DICTIONARY_CHECK_SEED, DICTIONARY_DECIDE_ROWS, DICTIONARY_DISTINCT_IN_TEN,
EncodedBlock, GlobalDictionary, MAX_ENCODE_WORKERS, MAX_PAGE, Part, PendingChunk, STRIPE_PARTS,
Settling, Spread, Unencoded, Writer, checksum, coded_page, invalid, push_validity,
seeded_checksum, stats, unique_codes, weight,
};
static BUSY: AtomicUsize = AtomicUsize::new(0);
pub const DICTIONARY_CAP_BYTES: u64 = 512 * 1024 * 1024;
#[derive(Debug)]
pub(crate) struct Coding {
flags: Box<[AtomicBool]>,
growth: Box<[AtomicU64]>,
held: AtomicU64,
cap: AtomicU64,
}
impl Coding {
pub(crate) fn new(flags: impl IntoIterator<Item = bool>) -> Self {
let flags = flags.into_iter().map(AtomicBool::new).collect::<Box<[_]>>();
let growth = flags.iter().map(|_| AtomicU64::new(0)).collect();
Self { flags, growth, held: AtomicU64::new(0), cap: AtomicU64::new(DICTIONARY_CAP_BYTES) }
}
pub(crate) fn cap(&self, bytes: u64) {
self.cap.store(bytes, Atomic::Relaxed);
}
fn recount(&self, before: u64, now: u64) -> u64 {
if now >= before {
self.held.fetch_add(now - before, Atomic::Relaxed) + (now - before)
} else {
self.held.fetch_sub(before - now, Atomic::Relaxed).saturating_sub(before - now)
}
}
fn grew_most(&self, index: usize) -> bool {
let mine = self.growth[index].load(Atomic::Relaxed);
mine > 0 && self.growth.iter().all(|other| other.load(Atomic::Relaxed) <= mine)
}
}
impl Deref for Coding {
type Target = [AtomicBool];
fn deref(&self) -> &[AtomicBool] {
&self.flags
}
}
struct Share(usize);
impl Share {
fn take(columns: usize, parts: usize) -> Self {
let busy = BUSY.fetch_add(1, Atomic::Relaxed) + 1;
let cores =
std::thread::available_parallelism().map_or(1, usize::from).min(MAX_ENCODE_WORKERS);
let workers = if parts <= 1 { 1 } else { (cores / busy).clamp(1, columns.max(1)) };
Self(workers)
}
}
impl Drop for Share {
fn drop(&mut self) {
BUSY.fetch_sub(1, Atomic::Relaxed);
}
}
#[derive(Debug, Clone)]
pub struct Preparer {
types: Vec<LogicalType>,
coded: Arc<Coding>,
profile: Option<Arc<LoadProfile>>,
}
#[derive(Debug)]
pub struct Prepared {
parts: Vec<Part>,
types: Vec<LogicalType>,
columns: Vec<Column>,
gathers: Vec<Option<stats::Gather>>,
profile: Option<Arc<LoadProfile>>,
}
#[derive(Debug)]
pub struct Merged {
parts: Vec<Part>,
columns: Vec<Merge>,
blocks: Vec<Unencoded>,
profile: Option<Arc<LoadProfile>>,
counted: bool,
}
#[derive(Debug)]
pub struct Paged {
parts: Vec<Part>,
columns: Vec<ColumnStripe>,
blocks: Vec<(usize, usize, EncodedBlock)>,
counted: bool,
}
enum Built {
Stripe(ColumnStripe),
Block(EncodedBlock),
}
#[derive(Debug)]
enum Column {
Pages(ColumnStripe),
Coded(Local),
}
#[derive(Debug)]
enum Merge {
Pages(ColumnStripe),
Codes {
parts: Vec<LocalPart>,
global: Vec<u32>,
},
Plain(Local),
}
const END: u32 = u32::MAX;
#[derive(Debug, Default)]
struct Local {
first: HashMap<u64, u32, Spread>,
next: Vec<u32>,
hashes: Vec<u64>,
checks: Vec<u64>,
bytes: Vec<u8>,
ends: Vec<usize>,
counts: Vec<u64>,
nulls: u64,
parts: Vec<LocalPart>,
blob: bool,
}
#[derive(Debug)]
struct LocalPart {
codes: Vec<u32>,
validity: Vec<u8>,
range: Range,
}
impl Local {
#[cfg(test)]
fn code_column(index: usize, held: &[PendingChunk]) -> Result<Self> {
let mut local = Self::default();
let mut mapped = None;
for pending in held {
local.code_part(pending.chunk.column(index)?, &mut mapped)?;
}
local.done();
Ok(local)
}
fn code_part(
&mut self,
column: &Vector,
mapped: &mut Option<(Arc<Vector>, Vec<u32>)>,
) -> Result<()> {
self.blob = column.logical_type() == &LogicalType::Blob;
if let Some(codes) = self.code_dictionary(column, mapped)? {
let mut validity = Vec::new();
push_validity(&mut validity, column);
self.parts.push(LocalPart { codes, validity, range: Range::of(column) });
return Ok(());
}
let flat = column.flatten()?;
let mut codes = Vec::with_capacity(flat.len());
let mut last = None;
for row in 0..flat.len() {
let text = flat.bytes_at(row).unwrap_or(b"");
let code = match last {
Some(code) if self.value(code) == text => code,
_ => self.code(text)?,
};
last = Some(code);
if flat.is_null_at(row) {
self.nulls += 1;
} else {
self.counts[code as usize] += 1;
}
codes.push(code);
}
let mut validity = Vec::new();
push_validity(&mut validity, &flat);
self.parts.push(LocalPart { codes, validity, range: Range::of(column) });
Ok(())
}
fn done(&mut self) {
self.first = HashMap::default();
self.next = Vec::new();
}
fn code_dictionary(
&mut self,
column: &Vector,
mapped: &mut Option<(Arc<Vector>, Vec<u32>)>,
) -> Result<Option<Vec<u32>>> {
let Some((codes, values)) = column.shared_dictionary_parts() else { return Ok(None) };
if !matches!(values.validity(), Validity::AllValid) {
return Ok(None);
}
let Some(codes) = codes.get(..column.len()) else { return Ok(None) };
let fresh = !matches!(mapped, Some((held, _)) if Arc::ptr_eq(held, values));
if fresh {
*mapped = Some((Arc::clone(values), vec![END; values.len()]));
}
let Some((_, map)) = mapped.as_mut() else { return Ok(None) };
let every = matches!(column.validity(), Validity::AllValid);
let mut coded = Vec::with_capacity(codes.len());
for (row, &code) in codes.iter().enumerate() {
if !every && !column.validity().is_valid(row) {
let code = self.code(b"")?;
self.nulls += 1;
coded.push(code);
continue;
}
let slot = map
.get_mut(code as usize)
.ok_or_else(|| invalid("a dictionary code is out of range"))?;
if *slot == END {
*slot = self.code(values.bytes_at(code as usize).unwrap_or(b""))?;
}
self.counts[*slot as usize] += 1;
coded.push(*slot);
}
Ok(Some(coded))
}
fn rows(&self) -> Result<Vec<Vector>> {
self.parts
.iter()
.map(|part| {
let len = part.codes.len();
let mut column = StringColumn::with_capacity(len);
for &code in &part.codes {
column.push_bytes(self.value(code));
}
let validity = match part.validity.split_first() {
Some((0, _)) => Validity::AllValid,
Some((1, _)) => Validity::AllInvalid,
Some((2, bits)) => {
let mut mask = Bitmap::all_valid(len);
for row in (0..len).filter(|row| bits[row / 8] & (1 << (row % 8)) == 0) {
mask.set(row, false);
}
Validity::Mask(mask)
}
_ => return Err(Error::internal("a coded part has no validity")),
};
let ty = if self.blob { LogicalType::Blob } else { LogicalType::Varchar };
Ok(Vector::flat(ty, Data::Varlen(column))?.with_validity(validity))
})
.collect()
}
fn values(&self) -> usize {
self.ends.len()
}
fn value(&self, code: u32) -> &[u8] {
let code = code as usize;
let from = if code == 0 { 0 } else { self.ends[code - 1] };
&self.bytes[from..self.ends[code]]
}
fn code(&mut self, text: &[u8]) -> Result<u32> {
let hash = checksum(text);
let Some(&first) = self.first.get(&hash) else {
let code = self.push(text, hash)?;
self.first.insert(hash, code);
return Ok(code);
};
let mut at = first;
loop {
if self.value(at) == text {
return Ok(at);
}
match self.next[at as usize] {
END => break,
next => at = next,
}
}
let code = self.push(text, hash)?;
self.next[at as usize] = code;
Ok(code)
}
fn push(&mut self, text: &[u8], hash: u64) -> Result<u32> {
let code = u32::try_from(self.ends.len())
.ok()
.filter(|&code| code != END)
.ok_or_else(|| invalid("a stripe has too many values in one column"))?;
self.bytes.extend_from_slice(text);
self.ends.push(self.bytes.len());
self.next.push(END);
self.hashes.push(hash);
self.checks.push(seeded_checksum(text, DICTIONARY_CHECK_SEED));
self.counts.push(0);
Ok(code)
}
fn merge_into(&self, dictionary: &mut GlobalDictionary) -> Result<Vec<u32>> {
let mut global = Vec::with_capacity(self.values());
for (code, (&hash, &check)) in self.hashes.iter().zip(&self.checks).enumerate() {
let text = self.value(code as u32);
let at = dictionary.code_hashed(text, hash, check)?;
let count = dictionary
.counts
.get_mut(at as usize)
.ok_or_else(|| invalid("global dictionary count code is out of range"))?;
*count = count.saturating_add(self.counts[code]);
global.push(at);
}
dictionary.nulls = dictionary.nulls.saturating_add(self.nulls);
Ok(global)
}
}
fn drops_dictionary(rows: usize, distinct: usize) -> bool {
rows >= DICTIONARY_DECIDE_ROWS
&& distinct.saturating_mul(10) > rows.saturating_mul(DICTIONARY_DISTINCT_IN_TEN)
}
fn demotes(rows: usize, new: usize, total: u64, coding: &Coding, index: usize) -> bool {
drops_dictionary(rows, new)
|| (total > coding.cap.load(Atomic::Relaxed) && coding.grew_most(index))
}
fn fan_out<T: Send>(
jobs: Vec<usize>,
workers: usize,
profile: Option<&LoadProfile>,
work: impl Fn(usize) -> Result<T> + Sync,
) -> Result<Vec<(usize, T)>> {
if workers <= 1 || jobs.len() <= 1 {
let _span = profile.map(|profile| profile.span(Stage::Pages));
return jobs.into_iter().map(|index| Ok((index, work(index)?))).collect();
}
let workers = workers.min(jobs.len());
let queue = Mutex::new(jobs);
let pieces = std::thread::scope(|scope| {
(0..workers)
.map(|_| {
scope.spawn(|| {
let _span = profile.map(|profile| profile.span(Stage::Pages));
let mut mine = Vec::new();
loop {
let taken = queue
.lock()
.map_err(|_| Error::internal("a native encode worker panicked"))?
.pop();
let Some(index) = taken else { break };
mine.push((index, work(index)?));
}
Ok(mine)
})
})
.collect::<Vec<_>>()
.into_iter()
.map(|handle| {
handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
})
.collect::<Result<Vec<Vec<_>>>>()
})?;
Ok(pieces.into_iter().flatten().collect())
}
impl Preparer {
pub fn prepare(&self, parts: Vec<((u64, u64), Chunk)>) -> Result<Prepared> {
if parts.len() > STRIPE_PARTS {
return Err(invalid("a stripe was handed more parts than it holds"));
}
let held = parts
.into_iter()
.filter(|(_, chunk)| !chunk.is_empty())
.map(|(order, chunk)| PendingChunk { order, chunk })
.collect::<Vec<_>>();
for pending in &held {
self.fits(&pending.chunk)?;
}
self.prepare_held(held)
}
fn fits(&self, chunk: &Chunk) -> Result<()> {
if chunk.width() != self.types.len() {
return Err(invalid("chunk width differs from table schema"));
}
for (index, ty) in self.types.iter().enumerate() {
if chunk.column(index)?.logical_type() != ty {
return Err(invalid("chunk type differs from table schema"));
}
}
Ok(())
}
pub(crate) fn prepare_held(&self, held: Vec<PendingChunk>) -> Result<Prepared> {
let mut building = self.start();
self.feed_held(&mut building, held)?;
self.finish(building)
}
#[must_use]
pub fn start(&self) -> Building {
let columns = (0..self.types.len())
.map(|index| {
let body = if self.coded[index].load(Atomic::Relaxed) {
Body::Coded(Local::default(), None)
} else {
Body::Pages(ColumnStripe::default(), Settling::default())
};
let gather = stats::Gather::new(&self.types[index], 0);
Mutex::new(Growing { body, gather })
})
.collect();
Building { parts: Vec::new(), columns }
}
pub fn feed(&self, building: &mut Building, parts: Vec<((u64, u64), Chunk)>) -> Result<()> {
if building.parts.len().saturating_add(parts.len()) > STRIPE_PARTS {
return Err(invalid("a stripe was handed more parts than it holds"));
}
let held = parts
.into_iter()
.filter(|(_, chunk)| !chunk.is_empty())
.map(|(order, chunk)| PendingChunk { order, chunk })
.collect::<Vec<_>>();
for pending in &held {
self.fits(&pending.chunk)?;
}
self.feed_held(building, held)
}
fn feed_held(&self, building: &mut Building, held: Vec<PendingChunk>) -> Result<()> {
if held.is_empty() {
return Ok(());
}
let width = self.types.len();
let opening = building.parts.is_empty();
let key = held.first().map_or((0, 0), |pending| pending.order);
let share = Share::take(width, held.len());
let mut jobs = (0..width).collect::<Vec<_>>();
jobs.sort_by_key(|&index| weight(&self.types[index]));
let columns = &building.columns;
fan_out(jobs, share.0, self.profile.as_deref(), |index| {
let mut growing = columns[index]
.lock()
.map_err(|_| Error::internal("a native encode worker panicked"))?;
let Growing { body, gather } = &mut *growing;
if matches!(body, Body::Coded(..)) && !self.coded[index].load(Atomic::Relaxed) {
body.plain()?;
}
if let Some(gather) = gather.as_mut() {
if opening {
gather.open_stripe(key);
}
for pending in &held {
gather.part(pending.chunk.column(index)?);
}
}
match body {
Body::Coded(local, mapped) => {
for pending in &held {
local.code_part(pending.chunk.column(index)?, mapped)?;
}
}
Body::Pages(stripe, settling) => {
for pending in &held {
Writer::encode_page(stripe, settling, pending.chunk.column(index)?)?;
}
}
}
Ok(())
})?;
drop(share);
building.parts.extend(held.iter().map(Part::of));
Ok(())
}
pub fn finish(&self, building: Building) -> Result<Prepared> {
let empty = building.parts.is_empty();
let (columns, gathers) = building
.columns
.into_iter()
.map(|growing| {
let Growing { body, gather } = growing
.into_inner()
.map_err(|_| Error::internal("a native encode worker panicked"))?;
let column = match body {
Body::Coded(mut local, _) => {
local.done();
Column::Coded(local)
}
Body::Pages(stripe, _) => Column::Pages(stripe),
};
let gather = gather.filter(|_| !empty).map(|mut gather| {
gather.close_stripe();
gather
});
Ok((column, gather))
})
.collect::<Result<Vec<_>>>()?
.into_iter()
.unzip();
Ok(Prepared {
parts: building.parts,
types: self.types.clone(),
columns,
gathers,
profile: self.profile.clone(),
})
}
}
#[derive(Debug)]
pub struct Building {
parts: Vec<Part>,
columns: Vec<Mutex<Growing>>,
}
impl Building {
#[must_use]
pub fn parts(&self) -> usize {
self.parts.len()
}
}
#[derive(Debug)]
struct Growing {
body: Body,
gather: Option<stats::Gather>,
}
#[derive(Debug)]
enum Body {
Pages(ColumnStripe, Settling),
Coded(Local, Option<(Arc<Vector>, Vec<u32>)>),
}
impl Body {
fn plain(&mut self) -> Result<()> {
let Self::Coded(local, _) = self else { return Ok(()) };
let mut stripe = ColumnStripe::default();
let mut settling = Settling::default();
for rows in local.rows()? {
Writer::encode_page(&mut stripe, &mut settling, &rows)?;
}
*self = Self::Pages(stripe, settling);
Ok(())
}
}
enum Slot<'a> {
Owned(&'a mut Option<GlobalDictionary>, &'a mut Option<stats::Gather>),
Lent(&'a Mutex<LentColumn>, &'a Lent),
}
struct Step<'a> {
index: usize,
column: Column,
slot: Slot<'a>,
gather: Option<stats::Gather>,
}
impl Step<'_> {
fn cost(&self, coded: &Coding) -> usize {
match &self.column {
Column::Coded(local) if coded[self.index].load(Atomic::Relaxed) => {
local.values().saturating_add(1)
}
_ => 0,
}
}
fn run(
self,
rows: usize,
coded: &Coding,
profile: Option<&LoadProfile>,
) -> Result<(usize, Merge, Vec<Unencoded>)> {
let Self { index, column, slot, gather } = self;
match slot {
Slot::Owned(dictionary, mine) => {
merge_column(index, column, gather, dictionary, mine, rows, coded, profile)
}
Slot::Lent(held, lent) => {
let mut held = held.lock().map_err(|_| Error::internal("a merge panicked"))?;
if lent.reclaimed.load(Atomic::Acquire) {
return Err(Error::internal("a stripe was merged after its table was closed"));
}
let LentColumn { dictionary, gather: mine } = &mut *held;
merge_column(index, column, gather, dictionary, mine, rows, coded, profile)
}
}
}
}
#[expect(clippy::too_many_arguments, reason = "one column's share of the stripe's merge state")]
fn merge_column(
index: usize,
column: Column,
stripe: Option<stats::Gather>,
dictionary: &mut Option<GlobalDictionary>,
gather: &mut Option<stats::Gather>,
rows: usize,
coded: &Coding,
profile: Option<&LoadProfile>,
) -> Result<(usize, Merge, Vec<Unencoded>)> {
if let (Some(mine), Some(stripe)) = (gather.as_mut(), stripe) {
mine.absorb(stripe);
}
let mut new = None;
let merge = match (column, dictionary.as_mut()) {
(Column::Pages(stripe), None) => Merge::Pages(stripe),
(Column::Pages(stripe), Some(global)) if global.demoted => Merge::Pages(stripe),
(Column::Pages(_), Some(_)) => {
return Err(Error::internal(
"a column with a global dictionary was prepared without one",
));
}
(Column::Coded(local), None) => Merge::Plain(local),
(Column::Coded(local), Some(global)) if global.demoted => Merge::Plain(local),
(Column::Coded(local), Some(global)) => {
if global.values() == 0 && drops_dictionary(rows, local.values()) {
if let Some(profile) = profile {
profile.release(global.charged);
}
coded.recount(global.charged, 0);
*dictionary = None;
coded[index].store(false, Atomic::Relaxed);
Merge::Plain(local)
} else {
let before = global.values();
let codes = local.merge_into(global)?;
new = Some(global.values() - before);
Merge::Codes { parts: local.parts, global: codes }
}
}
};
if let (Some(new), Some(global)) = (new, dictionary.as_mut()) {
let now = global.held_bytes();
coded.growth[index].store(now.saturating_sub(global.charged), Atomic::Relaxed);
let total =
coded.held.load(Atomic::Relaxed).saturating_add(now).saturating_sub(global.charged);
if demotes(rows, new, total, coded, index) {
global.demote();
coded[index].store(false, Atomic::Relaxed);
coded.growth[index].store(0, Atomic::Relaxed);
}
}
let blocks = match dictionary {
Some(dictionary) => {
dictionary.settle()?;
let blocks = dictionary.hand_out(index);
let (before, now) = dictionary.recharge(profile);
coded.recount(before, now);
blocks
}
None => Vec::new(),
};
Ok((index, merge, blocks))
}
fn merge_columns(prepared: Prepared, slots: Vec<Slot<'_>>, coded: &Coding) -> Result<Merged> {
let Prepared { parts, columns, gathers, profile, .. } = prepared;
let timing = profile.as_deref().map(|profile| profile.span(Stage::Dictionary));
let rows: usize = parts.iter().map(|part| part.rows).sum();
let width = columns.len();
if slots.len() != width || gathers.len() != width {
return Err(Error::internal("a stripe was merged into a table of another width"));
}
let mut steps = columns
.into_iter()
.zip(gathers)
.zip(slots)
.enumerate()
.map(|(index, ((column, gather), slot))| Step { index, column, slot, gather })
.collect::<Vec<_>>();
steps.sort_by_key(|step| step.cost(coded));
let workers = std::thread::available_parallelism()
.map_or(1, usize::from)
.min(MAX_ENCODE_WORKERS)
.min(steps.iter().filter(|step| step.cost(coded) > 0).count())
.max(1);
let done = if workers <= 1 {
steps
.into_iter()
.map(|step| step.run(rows, coded, profile.as_deref()))
.collect::<Result<Vec<_>>>()?
} else {
let queue = Mutex::new(steps);
let pieces = std::thread::scope(|scope| {
(0..workers)
.map(|_| {
scope.spawn(|| {
let mut mine = Vec::new();
loop {
let taken = queue
.lock()
.map_err(|_| Error::internal("a merge worker panicked"))?
.pop();
let Some(step) = taken else { break };
mine.push(step.run(rows, coded, profile.as_deref())?);
}
Ok(mine)
})
})
.collect::<Vec<_>>()
.into_iter()
.map(|handle| {
handle.join().map_err(|_| Error::internal("a merge worker panicked"))?
})
.collect::<Result<Vec<Vec<_>>>>()
})?;
pieces.into_iter().flatten().collect()
};
let mut slots: Vec<Option<(Merge, Vec<Unencoded>)>> = (0..width).map(|_| None).collect();
for (index, merge, blocks) in done {
slots[index] = Some((merge, blocks));
}
let mut merged = Vec::with_capacity(width);
let mut blocks = Vec::new();
for slot in slots {
let (merge, handed) = slot.ok_or_else(|| Error::internal("a column was never merged"))?;
merged.push(merge);
blocks.extend(handed);
}
drop(timing);
Ok(Merged { parts, columns: merged, blocks, profile, counted: false })
}
#[derive(Debug)]
pub(crate) struct Lent {
columns: Box<[Mutex<LentColumn>]>,
reclaimed: AtomicBool,
}
#[derive(Debug)]
pub(crate) struct LentColumn {
pub(crate) dictionary: Option<GlobalDictionary>,
gather: Option<stats::Gather>,
}
impl Lent {
pub(crate) fn columns(&self) -> &[Mutex<LentColumn>] {
&self.columns
}
fn take_back(&self, blocks: Vec<(usize, usize, EncodedBlock)>) -> Result<()> {
for (column, at, block) in blocks {
self.columns
.get(column)
.ok_or_else(|| Error::internal("a dictionary block came back to no column"))?
.lock()
.map_err(|_| Error::internal("a merge panicked"))?
.dictionary
.as_mut()
.ok_or_else(|| Error::internal("a dictionary block came back to no dictionary"))?
.take_back(at, block)?;
}
Ok(())
}
#[allow(clippy::type_complexity)]
pub(crate) fn reclaim(
&self,
) -> Result<(Vec<Option<GlobalDictionary>>, Vec<Option<stats::Gather>>)> {
self.reclaimed.store(true, Atomic::Release);
let mut dictionaries = Vec::with_capacity(self.columns.len());
let mut gathers = Vec::with_capacity(self.columns.len());
for column in &self.columns {
let mut held = column.lock().map_err(|_| Error::internal("a merge panicked"))?;
dictionaries.push(held.dictionary.take());
gathers.push(held.gather.take());
}
Ok((dictionaries, gathers))
}
}
#[derive(Debug, Clone)]
pub struct Merger {
lent: Arc<Lent>,
types: Vec<LogicalType>,
coded: Arc<Coding>,
}
impl Merger {
pub fn merge(&self, prepared: Prepared) -> Result<Merged> {
if prepared.types != self.types {
return Err(invalid("a stripe was prepared for a table of other columns"));
}
let slots = self.lent.columns.iter().map(|column| Slot::Lent(column, &self.lent)).collect();
merge_columns(prepared, slots, &self.coded)
}
pub fn give_back(&self, paged: &mut Paged) -> Result<()> {
self.lent.take_back(std::mem::take(&mut paged.blocks))
}
}
impl Merged {
pub fn pages(self) -> Result<Paged> {
let Self { parts, columns, blocks, profile, counted } = self;
let width = columns.len();
let mut jobs = (width..width + blocks.len())
.chain((0..width).filter(|&index| !matches!(columns[index], Merge::Pages(_))))
.collect::<Vec<_>>();
jobs.sort_by_key(|&index| index < width && matches!(columns[index], Merge::Plain(_)));
let share = Share::take(jobs.len(), parts.len());
let built = fan_out(jobs, share.0, profile.as_deref(), |index| {
let Some(column) = columns.get(index) else {
return Ok(Built::Block(blocks[index - width].encode()?));
};
Ok(Built::Stripe(match column {
Merge::Codes { parts, global } => code_pages(parts, global)?,
Merge::Plain(local) => {
Writer::encode_pages(&local.rows()?.iter().collect::<Vec<_>>())?
}
Merge::Pages(_) => {
return Err(Error::internal("a finished column was queued to be built"));
}
}))
})?;
drop(share);
let mut slots: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
let mut encoded = Vec::with_capacity(blocks.len());
for (index, one) in built {
match one {
Built::Stripe(stripe) => slots[index] = Some(stripe),
Built::Block(block) => {
let (column, at) = blocks[index - width].place();
encoded.push((column, at, block));
}
}
}
let columns = columns
.into_iter()
.zip(slots)
.map(|(column, slot)| match (column, slot) {
(Merge::Pages(stripe), _) | (_, Some(stripe)) => Ok(stripe),
_ => Err(Error::internal("a column was never encoded")),
})
.collect::<Result<Vec<_>>>()?;
Ok(Paged { parts, columns, blocks: encoded, counted })
}
}
fn code_pages(parts: &[LocalPart], global: &[u32]) -> Result<ColumnStripe> {
let mut stripe = ColumnStripe {
pages: Vec::with_capacity(parts.len()),
codes: Vec::with_capacity(parts.len()),
sieves: Vec::with_capacity(parts.len()),
ranges: Vec::with_capacity(parts.len()),
};
for part in parts {
let codes = part
.codes
.iter()
.map(|&code| global.get(code as usize).copied())
.collect::<Option<Vec<_>>>()
.ok_or_else(|| Error::internal("a stripe's code has no global code"))?;
let bytes = coded_page(&codes, &part.validity)?;
if bytes.len() > MAX_PAGE {
return Err(invalid("column page exceeds the configured bound"));
}
stripe.pages.push(bytes);
stripe.codes.push(Some(unique_codes(&codes)));
stripe.sieves.push(None);
stripe.ranges.push(part.range.clone());
}
Ok(stripe)
}
impl Writer {
#[must_use]
pub fn preparer(&self) -> Preparer {
Preparer {
types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
coded: Arc::clone(&self.coded),
profile: self.profile.clone(),
}
}
pub fn merge(&mut self, prepared: Prepared) -> Result<Merged> {
self.flush_pending()?;
if prepared.columns.len() != self.table.fields.len()
|| prepared.types.iter().ne(self.table.fields.iter().map(|field| &field.ty))
{
return Err(invalid("a stripe was prepared for a table of other columns"));
}
self.table.rows = prepared
.parts
.iter()
.try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
.ok_or_else(|| invalid("row count overflow"))?;
self.merge_held(prepared)
}
pub(crate) fn merge_held(&mut self, prepared: Prepared) -> Result<Merged> {
let slots = match &self.lent {
Some(lent) => lent.columns.iter().map(|column| Slot::Lent(column, lent)).collect(),
None => self
.dictionaries
.iter_mut()
.zip(self.gathers.iter_mut())
.map(|(dictionary, gather)| Slot::Owned(dictionary, gather))
.collect::<Vec<_>>(),
};
let mut merged = merge_columns(prepared, slots, &self.coded)?;
merged.counted = true;
Ok(merged)
}
pub fn merger(&mut self) -> Result<Merger> {
self.flush_pending()?;
let lent = match &self.lent {
Some(lent) => Arc::clone(lent),
None => {
let lent = Arc::new(Lent {
columns: std::mem::take(&mut self.dictionaries)
.into_iter()
.zip(std::mem::take(&mut self.gathers))
.map(|(dictionary, gather)| Mutex::new(LentColumn { dictionary, gather }))
.collect(),
reclaimed: AtomicBool::new(false),
});
self.lent = Some(Arc::clone(&lent));
lent
}
};
Ok(Merger {
lent,
types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
coded: Arc::clone(&self.coded),
})
}
pub fn write(&mut self, paged: Paged) -> Result<()> {
self.write_paged(paged)
}
pub(crate) fn write_paged(&mut self, paged: Paged) -> Result<()> {
let Paged { parts, columns, blocks, counted } = paged;
if !counted {
self.table.rows = parts
.iter()
.try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
.ok_or_else(|| invalid("row count overflow"))?;
}
if let Some(lent) = &self.lent {
lent.take_back(blocks)?;
} else {
for (column, at, block) in blocks {
self.dictionaries
.get_mut(column)
.and_then(Option::as_mut)
.ok_or_else(|| {
Error::internal("a dictionary block came back to no dictionary")
})?
.take_back(at, block)?;
}
}
if parts.is_empty() {
return self.place_blocks();
}
self.write_stripe(&parts, columns)
}
pub fn append_prepared(&mut self, prepared: Prepared) -> Result<()> {
let merged = self.merge(prepared)?;
let paged = merged.pages()?;
self.write(paged)
}
}
#[cfg(test)]
mod tests {
use std::fs;
use std::path::PathBuf;
use std::time::{SystemTime, UNIX_EPOCH};
use rudb_common::{Field, Value};
use rudb_vector::Vector;
use super::*;
use crate::Reader;
const PART: usize = 1_000;
#[test]
fn dictionary_parts_code_as_their_rows_would() {
let texts = |values: &[&str]| {
Arc::new(
Vector::from_values(
LogicalType::Varchar,
&values
.iter()
.map(|text| Value::Varchar((*text).to_string()))
.collect::<Vec<_>>(),
)
.expect("a dictionary"),
)
};
let shared = texts(&["b", "a", "", "c", "unused"]);
let other = texts(&["c", "d", "a"]);
let mut nulls = Bitmap::all_valid(6);
nulls.set(1, false);
nulls.set(4, false);
let parts = [
Vector::dictionary_over(vec![3, 3, 1, 0, 2, 1], Arc::clone(&shared)).expect("codes"),
Vector::dictionary_over(vec![0, 3, 1, 1, 3, 2], Arc::clone(&shared))
.expect("codes")
.with_validity(Validity::Mask(nulls)),
Vector::dictionary_over(vec![1, 2, 0, 1], other).expect("codes"),
];
let held = |flat: bool| {
parts
.iter()
.enumerate()
.map(|(at, part)| PendingChunk {
order: (at as u64, 0),
chunk: Chunk::new(vec![if flat {
part.flatten().expect("flat")
} else {
part.clone()
}])
.expect("a chunk"),
})
.collect::<Vec<_>>()
};
let parquet = held(false);
let before = rudb_common::slow::here();
let coded = Local::code_column(0, &parquet).expect("coded");
assert_eq!(
rudb_common::slow::here().since(before).get(rudb_common::slow::Cause::Flatten),
0,
"a part that came in as codes was flattened",
);
let flat = Local::code_column(0, &held(true)).expect("coded");
assert_eq!(coded.values(), flat.values());
for code in 0..flat.values() as u32 {
assert_eq!(coded.value(code), flat.value(code), "value {code}");
}
assert_eq!(coded.counts, flat.counts);
assert_eq!(coded.nulls, flat.nulls);
assert_eq!(coded.nulls, 2);
assert_eq!(coded.hashes, flat.hashes);
assert_eq!(coded.checks, flat.checks);
assert_eq!(coded.parts.len(), flat.parts.len());
for (coded, flat) in coded.parts.iter().zip(&flat.parts) {
assert_eq!(coded.codes, flat.codes);
assert_eq!(coded.validity, flat.validity);
}
}
fn path(label: &str) -> PathBuf {
let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
std::env::temp_dir()
.join(format!("rudb-prepare-{label}-{}-{stamp}.rdb", std::process::id()))
}
fn fields() -> Vec<Field> {
vec![
Field::required("id", LogicalType::BigInt),
Field::new("city", LogicalType::Varchar),
Field::new("note", LogicalType::Varchar),
]
}
fn row(id: usize) -> [Value; 3] {
let city = if id % 11 == 0 {
Value::Null
} else {
Value::Varchar(format!("city {}", (id / 7) % 13))
};
let note = if id % 17 == 0 { Value::Null } else { Value::Varchar(format!("note {id}")) };
[Value::BigInt(id as i64), city, note]
}
fn stripe(first: usize, parts: usize) -> Vec<((u64, u64), Chunk)> {
(first..first + parts)
.map(|part| {
let rows = (part * PART..(part + 1) * PART).map(row).collect::<Vec<_>>();
let column = |at: usize| {
let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
Vector::from_values(fields()[at].ty.clone(), &values).expect("a column")
};
let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
((part as u64, 0), chunk)
})
.collect()
}
fn runs() -> Vec<Vec<((u64, u64), Chunk)>> {
vec![stripe(5, 5), stripe(0, 5), stripe(10, 3)]
}
fn check(path: &PathBuf) {
let reader = Reader::open(path).expect("reopen");
assert_eq!(reader.parts(), 13);
for part in 0..13 {
let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
for at in [0, 17, PART - 1] {
let want = row(part * PART + at);
for (column, value) in want.iter().enumerate() {
assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
}
}
}
}
#[test]
fn stripes_prepared_before_any_is_merged_write_the_same_bytes_as_one_at_a_time() {
let alone = path("alone");
let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
for run in runs() {
writer.append_stripe(run).expect("a stripe");
}
writer.finish().expect("commit");
let split = path("split");
let mut writer = Writer::create(&split, "t", fields()).expect("a file");
let preparer = writer.preparer();
let prepared = runs()
.into_iter()
.map(|run| preparer.prepare(run).expect("prepared"))
.collect::<Vec<_>>();
for one in prepared {
writer.append_prepared(one).expect("a stripe");
}
assert!(!preparer.coded[2].load(Atomic::Relaxed), "note lost its dictionary");
assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
writer.finish().expect("commit");
assert_eq!(fs::read(&alone).expect("read"), fs::read(&split).expect("read"));
check(&split);
fs::remove_file(alone).expect("remove");
fs::remove_file(split).expect("remove");
}
#[test]
fn stripes_written_in_another_order_than_they_were_merged_read_back() {
let path = path("crossed");
let mut writer = Writer::create(&path, "t", fields()).expect("a file");
let preparer = writer.preparer();
let mut merged = runs()
.into_iter()
.map(|run| writer.merge(preparer.prepare(run).expect("prepared")).expect("merged"))
.map(|merged| merged.pages().expect("paged"))
.collect::<Vec<_>>();
merged.reverse();
for paged in merged {
writer.write(paged).expect("written");
}
writer.finish().expect("commit");
check(&path);
fs::remove_file(path).expect("remove");
}
#[test]
fn stripes_merged_through_a_merger_write_the_same_bytes_as_the_writer() {
let alone = path("alone-merger");
let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
for run in runs() {
writer.append_stripe(run).expect("a stripe");
}
writer.finish().expect("commit");
let lent = path("lent");
let mut writer = Writer::create(&lent, "t", fields()).expect("a file");
let preparer = writer.preparer();
let merger = writer.merger().expect("a merger");
for run in runs() {
let merged = merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
let mut paged = merged.pages().expect("paged");
merger.give_back(&mut paged).expect("given back");
writer.write(paged).expect("written");
}
assert_eq!(writer.table.rows, 13 * PART);
writer.finish().expect("commit");
assert_eq!(fs::read(&alone).expect("read"), fs::read(&lent).expect("read"));
check(&lent);
fs::remove_file(alone).expect("remove");
fs::remove_file(lent).expect("remove");
}
#[test]
fn stripes_fed_in_batches_write_the_same_bytes_as_prepared_whole() {
let whole = path("whole");
let mut writer = Writer::create(&whole, "t", fields()).expect("a file");
for run in runs() {
writer.append_stripe(run).expect("a stripe");
}
writer.finish().expect("commit");
let fed = path("fed");
let mut writer = Writer::create(&fed, "t", fields()).expect("a file");
let preparer = writer.preparer();
let merger = writer.merger().expect("a merger");
let nothing = Chunk::new(
fields()
.iter()
.map(|field| Vector::from_values(field.ty.clone(), &[]).expect("a column"))
.collect(),
)
.expect("a chunk");
for mut run in runs() {
let mut building = preparer.start();
preparer.feed(&mut building, Vec::new()).expect("fed nothing");
while !run.is_empty() {
let rest = run.split_off(2.min(run.len()));
let mut batch = std::mem::replace(&mut run, rest);
batch.push(((u64::MAX, 0), nothing.clone()));
preparer.feed(&mut building, batch).expect("fed");
}
let merged =
merger.merge(preparer.finish(building).expect("finished")).expect("merged");
let mut paged = merged.pages().expect("paged");
merger.give_back(&mut paged).expect("given back");
writer.write(paged).expect("written");
}
writer.finish().expect("commit");
assert_eq!(fs::read(&whole).expect("read"), fs::read(&fed).expect("read"));
check(&fed);
fs::remove_file(whole).expect("remove");
fs::remove_file(fed).expect("remove");
}
#[test]
fn a_column_dropped_while_its_stripe_is_built_turns_to_pages() {
let path = path("dropped-while-built");
let mut writer = Writer::create(&path, "t", fields()).expect("a file");
let preparer = writer.preparer();
let merger = writer.merger().expect("a merger");
let write = |writer: &mut Writer, building: Building| {
let merged =
merger.merge(preparer.finish(building).expect("finished")).expect("merged");
let mut paged = merged.pages().expect("paged");
merger.give_back(&mut paged).expect("given back");
writer.write(paged).expect("written");
};
let mut runs = runs().into_iter();
let mut first = preparer.start();
let mut second = preparer.start();
let mut later = runs.next().expect("a run");
preparer.feed(&mut second, later.drain(..2).collect()).expect("fed");
preparer.feed(&mut first, runs.next().expect("a run")).expect("fed");
write(&mut writer, first);
let note = |building: &Building| {
matches!(building.columns[2].lock().expect("unpoisoned").body, Body::Pages(..))
};
assert!(!note(&second), "still coded until it is fed again");
preparer.feed(&mut second, later).expect("fed");
assert!(note(&second), "turned to pages once fed after the drop");
write(&mut writer, second);
let mut last = preparer.start();
preparer.feed(&mut last, runs.next().expect("a run")).expect("fed");
write(&mut writer, last);
writer.finish().expect("commit");
check(&path);
fs::remove_file(path).expect("remove");
}
#[test]
fn a_stripe_fed_more_parts_than_it_holds_is_refused() {
let path = path("overfed");
let writer = Writer::create(&path, "t", fields()).expect("a file");
let preparer = writer.preparer();
let mut building = preparer.start();
preparer.feed(&mut building, stripe(0, STRIPE_PARTS - 1)).expect("fed");
assert_eq!(building.parts(), STRIPE_PARTS - 1);
assert!(preparer.feed(&mut building, stripe(STRIPE_PARTS, 2)).is_err());
drop(writer);
let _ = fs::remove_file(path);
}
#[test]
fn stripes_merged_on_several_threads_at_once_read_back() {
let path = path("merged-at-once");
let mut writer = Writer::create(&path, "t", fields()).expect("a file");
let preparer = writer.preparer();
let merger = writer.merger().expect("a merger");
let writer = Mutex::new(writer);
std::thread::scope(|scope| {
for run in runs() {
let (preparer, merger, writer) = (&preparer, &merger, &writer);
scope.spawn(move || {
let merged =
merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
let mut paged = merged.pages().expect("paged");
merger.give_back(&mut paged).expect("given back");
writer.lock().expect("the writer").write(paged).expect("written");
});
}
});
writer.into_inner().expect("the writer").finish().expect("commit");
check(&path);
fs::remove_file(path).expect("remove");
}
#[test]
fn a_merge_after_the_table_is_closed_is_refused() {
let path = path("late");
let mut writer = Writer::create(&path, "t", fields()).expect("a file");
let preparer = writer.preparer();
let merger = writer.merger().expect("a merger");
writer.finish().expect("commit");
let prepared = preparer.prepare(stripe(0, 2)).expect("prepared");
assert!(merger.merge(prepared).is_err());
fs::remove_file(path).expect("remove");
}
#[test]
fn a_stripe_of_another_table_is_refused_at_the_merge() {
let path = path("refused");
let mut writer = Writer::create(&path, "t", fields()).expect("a file");
let other = Writer::create(path.with_extension("other"), "u", vec![fields().remove(0)])
.expect("a file");
let prepared = other.preparer().prepare(vec![]).expect("nothing to prepare");
assert!(writer.merge(prepared).is_err());
assert_eq!(writer.table.rows, 0);
drop(other);
fs::remove_file(path.with_extension("other")).expect("remove");
fs::remove_file(path).expect("remove");
}
fn turning(id: usize) -> [Value; 3] {
let url = match id {
_ if id % 13 == 0 => Value::Null,
_ if id < 5 * PART => Value::Varchar(format!("https://example.com/{}", id % 20)),
_ => Value::Varchar(format!("https://example.com/page/{id}")),
};
[Value::BigInt(id as i64), Value::Varchar(format!("city {}", id % 13)), url]
}
fn turning_fields() -> Vec<Field> {
vec![
Field::required("id", LogicalType::BigInt),
Field::new("city", LogicalType::Varchar),
Field::new("url", LogicalType::Varchar),
]
}
fn turning_stripe(
rows: fn(usize) -> [Value; 3],
first: usize,
parts: usize,
) -> Vec<((u64, u64), Chunk)> {
(first..first + parts)
.map(|part| {
let rows = (part * PART..(part + 1) * PART).map(rows).collect::<Vec<_>>();
let column = |at: usize| {
let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
Vector::from_values(turning_fields()[at].ty.clone(), &values).expect("a column")
};
let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
((part as u64, 0), chunk)
})
.collect()
}
fn turning_runs() -> Vec<Vec<((u64, u64), Chunk)>> {
vec![
turning_stripe(turning, 0, 5),
turning_stripe(turning, 5, 5),
turning_stripe(turning, 10, 3),
]
}
fn check_turned(path: &PathBuf) {
let reader = Reader::open(path).expect("reopen");
assert_eq!(reader.parts(), 13);
assert_eq!(reader.table().demoted, [false, false, true], "only url is demoted");
for part in 0..13 {
let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
let url = chunk.column(2).expect("url");
assert!(url.stable_dictionary_parts().is_none(), "part {part} hands out no codes");
for at in 0..PART {
let want = turning(part * PART + at);
for (column, value) in want.iter().enumerate() {
assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
}
}
}
assert_eq!(reader.distinct_values(2).expect("asked"), None);
assert_eq!(reader.text_extremes(2).expect("asked"), None);
assert_eq!(reader.exact_frequencies(2).expect("asked"), None);
assert_eq!(reader.top_frequencies(2, 5).expect("asked"), None);
assert!(!reader.skips_codes(0, 2, &[0]).expect("asked"), "no code proves a value absent");
assert_eq!(reader.distinct_values(1).expect("asked"), Some(13));
assert!(reader.text_extremes(1).expect("asked").is_some());
}
#[test]
fn a_column_that_turns_unique_is_demoted_and_reads_back() {
let alone = path("demoted-alone");
let mut writer = Writer::create(&alone, "t", turning_fields()).expect("a file");
let preparer = writer.preparer();
let mut runs = turning_runs().into_iter();
writer.append_stripe(runs.next().expect("a run")).expect("a stripe");
assert!(preparer.coded[2].load(Atomic::Relaxed), "url repeats in its first stripe");
for run in runs {
writer.append_stripe(run).expect("a stripe");
}
assert!(!preparer.coded[2].load(Atomic::Relaxed), "url was demoted");
assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
writer.finish().expect("commit");
check_turned(&alone);
let split = path("demoted-split");
let mut writer = Writer::create(&split, "t", turning_fields()).expect("a file");
let preparer = writer.preparer();
let prepared = turning_runs()
.into_iter()
.map(|run| preparer.prepare(run).expect("prepared"))
.collect::<Vec<_>>();
for one in prepared {
writer.append_prepared(one).expect("a stripe");
}
writer.finish().expect("commit");
check_turned(&split);
fs::remove_file(alone).expect("remove");
fs::remove_file(split).expect("remove");
}
fn growing(id: usize) -> [Value; 3] {
let url = Value::Varchar(format!("https://example.com/{}", id / 4));
[Value::BigInt(id as i64), Value::Varchar(format!("city {}", id % 13)), url]
}
#[test]
fn the_dictionary_cap_demotes_the_column_that_grew_most() {
let path = path("capped");
let mut writer = Writer::create(&path, "t", turning_fields())
.expect("a file")
.with_dictionary_cap(1 << 30);
let preparer = writer.preparer();
writer.append_stripe(turning_stripe(growing, 0, 5)).expect("a stripe");
assert!(preparer.coded[2].load(Atomic::Relaxed), "url is under the cap");
assert!(preparer.coded[1].load(Atomic::Relaxed), "city is under the cap");
writer.coded.cap(1);
writer.append_stripe(turning_stripe(growing, 5, 5)).expect("a stripe");
assert!(!preparer.coded[2].load(Atomic::Relaxed), "url grew most and was demoted");
assert!(preparer.coded[1].load(Atomic::Relaxed), "city grew nothing and keeps it");
writer.append_stripe(turning_stripe(growing, 10, 3)).expect("a stripe");
assert!(preparer.coded[1].load(Atomic::Relaxed), "city still grows nothing");
writer.finish().expect("commit");
let reader = Reader::open(&path).expect("reopen");
assert_eq!(reader.table().demoted, [false, false, true]);
for part in 0..13 {
let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
for at in 0..PART {
let want = growing(part * PART + at);
for (column, value) in want.iter().enumerate() {
assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
}
}
}
assert_eq!(reader.distinct_values(1).expect("asked"), Some(13));
assert_eq!(reader.distinct_values(2).expect("asked"), None);
fs::remove_file(path).expect("remove");
}
}