use std::collections::HashMap;
use std::ops::Range;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use rudb_catalog::Table;
use rudb_common::{Error, Field, LogicalType, Result};
use rudb_csv::Reader as CsvReader;
use rudb_functions::{
FILE_ROW_NUMBER, Given, TableFunction, csv_given, open_csv, open_parquet, series_length,
};
use rudb_kernels::cast;
use rudb_metrics::Counters;
use rudb_parquet::{Bound, Op, Reader, Test, skips};
use rudb_pipeline::{Morsel, Progress, Source};
use rudb_plan::{ExprRef, Plan, Slice};
use rudb_storage::Probe;
use rudb_vector::{Chunk, Data, VECTOR_SIZE, Vector};
use crate::expr::evaluate_all;
use crate::schema::Schema;
#[derive(Debug)]
pub(crate) struct Handout {
next: AtomicU64,
total: u64,
}
impl Handout {
pub(crate) fn new(total: usize) -> Self {
Self { next: AtomicU64::new(0), total: u64::try_from(total).unwrap_or(u64::MAX) }
}
pub(crate) fn total(&self) -> usize {
usize::try_from(self.total).unwrap_or(usize::MAX)
}
pub(crate) fn take(&self) -> Option<Morsel> {
let at = self.next.fetch_add(1, Ordering::Relaxed);
(at < self.total).then(|| Morsel::new(at, at, at + 1))
}
}
pub(crate) fn position(morsel: &Morsel) -> usize {
usize::try_from(morsel.cursor()).unwrap_or(usize::MAX)
}
fn poisoned<T>(_: T) -> Error {
Error::internal("a thread panicked while reading a file")
}
#[derive(Debug)]
pub(crate) struct Scan<'a> {
table: &'a Table,
columns: Vec<usize>,
probes: Vec<Probe>,
schema: Schema,
chunks: Handout,
skipped: AtomicUsize,
}
impl<'a> Scan<'a> {
pub(crate) fn new(
plan: &Plan,
table: &'a Table,
index: u32,
projection: Slice,
tests: Vec<(usize, Op, Bound)>,
) -> Result<Self> {
let fields = plan.field_list(projection).to_vec();
let mut columns = Vec::with_capacity(fields.len());
for field in &fields {
let position = table.column_index(&field.name).ok_or_else(|| {
Error::catalog(format!(
"Table \"{}\" does not have a column named \"{}\"",
table.name().table,
field.name
))
})?;
columns.push(position);
}
let probes = tests
.into_iter()
.filter_map(|(at, op, value)| Some(Probe { column: *columns.get(at)?, op, value }))
.collect();
let schema = Schema::numbered(fields, index);
let chunks = Handout::new(table.rows().chunk_count());
Ok(Self { table, columns, probes, schema, chunks, skipped: AtomicUsize::new(0) })
}
pub(crate) fn schema(&self) -> &Schema {
&self.schema
}
}
impl Source for Scan<'_> {
fn morsel(&self) -> Option<Morsel> {
self.chunks.take()
}
fn morsels(&self, _threads: usize) -> Option<usize> {
Some(self.chunks.total())
}
fn read(&self, morsel: &mut Morsel, out: &mut Chunk) -> Result<Progress> {
let at = position(morsel);
if at >= self.table.rows().chunk_count() {
*out = Chunk::empty(&self.schema.types());
return Ok(Progress::Done);
}
morsel.advance(1);
if !self.probes.is_empty() && self.table.rows().skips(at, &self.probes) {
self.skipped.fetch_add(1, Ordering::Relaxed);
*out = Chunk::empty(&self.schema.types());
return Ok(Progress::Done);
}
*out = self.table.rows().read(at, &self.columns)?;
Ok(Progress::Done)
}
}
#[derive(Debug)]
pub(crate) struct Dummy {
schema: Schema,
one: Handout,
}
impl Dummy {
pub(crate) fn new() -> Self {
Self { schema: Schema::empty(), one: Handout::new(1) }
}
pub(crate) fn schema(&self) -> &Schema {
&self.schema
}
}
impl Source for Dummy {
fn morsel(&self) -> Option<Morsel> {
self.one.take()
}
fn morsels(&self, _threads: usize) -> Option<usize> {
Some(self.one.total())
}
fn read(&self, morsel: &mut Morsel, out: &mut Chunk) -> Result<Progress> {
*out = Chunk::with_rows(Vec::new(), 1)?;
morsel.advance(1);
Ok(Progress::Done)
}
}
#[derive(Debug)]
pub(crate) struct Values {
schema: Schema,
chunks: Vec<Chunk>,
handout: Handout,
}
impl Values {
pub(crate) fn new(plan: &Plan, index: u32, columns: Slice, rows: Slice) -> Result<Self> {
let fields = plan.field_list(columns).to_vec();
let schema = Schema::numbered(fields, index);
let types = schema.types();
let source = Schema::empty();
let one = Chunk::with_rows(Vec::new(), 1)?;
let mut down: Vec<Vec<rudb_common::Value>> = vec![Vec::new(); types.len()];
for row in plan.row_list(rows) {
let exprs: Vec<ExprRef> = plan.expr_list(*row).to_vec();
if exprs.len() != types.len() {
return Err(Error::internal(format!(
"a VALUES row of {} expressions in a {} column list",
exprs.len(),
types.len()
)));
}
let evaluated = evaluate_all(plan, &exprs, &source, &one)?;
for (position, vector) in evaluated.iter().enumerate() {
down[position].push(vector.value_at(0));
}
}
let total = down.first().map_or(0, Vec::len);
let mut chunks = Vec::new();
let mut start = 0;
while start < total {
let end = (start + VECTOR_SIZE).min(total);
let mut built = Vec::with_capacity(types.len());
for (position, ty) in types.iter().enumerate() {
built.push(Vector::from_values(ty.clone(), &down[position][start..end])?);
}
chunks.push(Chunk::with_rows(built, end - start)?);
start = end;
}
let handout = Handout::new(chunks.len());
Ok(Self { schema, chunks, handout })
}
pub(crate) fn schema(&self) -> &Schema {
&self.schema
}
}
#[derive(Debug)]
pub(crate) struct Series {
schema: Schema,
start: i64,
step: i64,
rows: u64,
morsels: AtomicU64,
}
const RUN: u64 = 16 * VECTOR_SIZE as u64;
impl Series {
pub(crate) fn new(plan: &Plan, index: u32, function: &str, args: Slice) -> Result<Self> {
let Some(function) = TableFunction::lookup(function) else {
return Err(Error::internal(format!("a plan with a table function called {function}")));
};
let fields = vec![Field::new(function.name(), LogicalType::BigInt)];
let schema = Schema::numbered(fields, index);
let exprs: Vec<ExprRef> = plan.expr_list(args).to_vec();
let source = Schema::empty();
let one = Chunk::with_rows(Vec::new(), 1)?;
let evaluated = evaluate_all(plan, &exprs, &source, &one)?;
let mut given = Vec::with_capacity(evaluated.len());
for vector in &evaluated {
match vector.value_at(0) {
rudb_common::Value::Null => return Ok(Self::empty(schema)),
rudb_common::Value::BigInt(n) => given.push(n),
other => {
return Err(Error::internal(format!(
"a table function argument bound as BIGINT arrived as {other}"
)));
}
}
}
let (start, stop, step) = match given.as_slice() {
[stop] => (0, *stop, 1),
[start, stop] => (*start, *stop, 1),
[start, stop, step] => (*start, *stop, *step),
_ => {
return Err(Error::internal(format!(
"{}() bound with {} arguments",
function.name(),
given.len()
)));
}
};
let rows = u64::try_from(series_length(function, start, stop, step)?).unwrap_or(u64::MAX);
Ok(Self { schema, start, step, rows, morsels: AtomicU64::new(0) })
}
fn empty(schema: Schema) -> Self {
Self { schema, start: 0, step: 1, rows: 0, morsels: AtomicU64::new(0) }
}
pub(crate) fn schema(&self) -> &Schema {
&self.schema
}
fn value_at(&self, position: u64) -> i64 {
let steps = i64::try_from(position).unwrap_or(i64::MAX);
self.start.saturating_add(self.step.saturating_mul(steps))
}
}
impl Source for Series {
fn morsel(&self) -> Option<Morsel> {
let index = self.morsels.fetch_add(1, Ordering::Relaxed);
let start = index.saturating_mul(RUN);
(start < self.rows)
.then(|| Morsel::new(index, start, self.rows.min(start.saturating_add(RUN))))
}
fn morsels(&self, _threads: usize) -> Option<usize> {
Some(usize::try_from(self.rows.div_ceil(RUN)).unwrap_or(usize::MAX))
}
fn read(&self, morsel: &mut Morsel, out: &mut Chunk) -> Result<Progress> {
let count = usize::try_from(morsel.remaining()).unwrap_or(usize::MAX).min(VECTOR_SIZE);
if count == 0 {
*out = Chunk::empty(&[LogicalType::BigInt]);
return Ok(Progress::Done);
}
let mut at = self.value_at(morsel.cursor());
let mut counted = Vec::with_capacity(count);
for _ in 0..count {
counted.push(at);
at = at.saturating_add(self.step);
}
morsel.advance(u64::try_from(count).unwrap_or(u64::MAX));
let vector = Vector::flat(LogicalType::BigInt, Data::Int64(counted.into()))?;
*out = Chunk::with_rows(vec![vector], count)?;
Ok(if morsel.is_drained() { Progress::Done } else { Progress::More })
}
}
#[derive(Debug)]
pub(crate) struct FileScan {
function: TableFunction,
paths: Vec<String>,
given: Given,
wanted: Vec<Field>,
numbered: bool,
schema: Schema,
tests: Vec<(usize, Op, Bound)>,
cutting: Mutex<Cutting>,
open: Mutex<HashMap<u64, Arc<Mutex<Piece>>>>,
counters: Option<Arc<Counters>>,
}
const MORSEL_ROWS: usize = 32_768;
const FINE: u64 = 4;
fn morsel_rows(reader: &FileReader, groups: usize, threads: usize) -> usize {
let FileReader::Parquet(parquet) = reader else { return 0 };
if groups == 0 || threads <= groups {
return 0;
}
let page = parquet.page_bytes().unwrap_or(u64::MAX);
cut_rows(page, parquet.chunk_bytes(), group_rows(parquet, 0), groups, threads)
}
fn cut_rows(page: u64, whole: u64, rows: usize, groups: usize, threads: usize) -> usize {
if groups == 0 || threads <= groups || rows == 0 || page == 0 {
return 0;
}
let cut = rows.div_ceil(threads.div_ceil(groups)).max(MORSEL_ROWS);
let taking = u64::try_from(cut.min(rows)).unwrap_or(u64::MAX);
let each = whole / u64::try_from(rows).unwrap_or(u64::MAX) * taking;
if page.saturating_mul(FINE) > each { 0 } else { cut }
}
fn next_piece(rows: usize, part: usize, target: usize) -> Range<usize> {
let each = rows.div_ceil(parts(rows, target));
let upto = part.saturating_add(each).min(rows);
part.min(upto)..upto
}
fn parts(rows: usize, target: usize) -> usize {
if target == 0 {
return 1;
}
rows.div_ceil(target).max(1)
}
fn group_rows(reader: &Reader, at: usize) -> usize {
reader
.metadata()
.row_groups
.get(at)
.map_or(0, |group| usize::try_from(group.rows).unwrap_or(usize::MAX))
}
fn aim(cutting: &mut Cutting) {
let Some(reader) = cutting.reader.as_ref() else { return };
cutting.cut = morsel_rows(reader, cutting.groups, cutting.threads);
cutting.pieces = pieces(reader, cutting.cut);
}
fn pieces(reader: &FileReader, target: usize) -> usize {
let FileReader::Parquet(reader) = reader else { return 1 };
reader
.metadata()
.row_groups
.iter()
.map(|group| parts(usize::try_from(group.rows).unwrap_or(usize::MAX), target))
.sum::<usize>()
.max(1)
}
#[derive(Debug)]
struct Cutting {
at: usize,
reader: Option<FileReader>,
group: usize,
groups: usize,
part: usize,
cut: usize,
threads: usize,
pieces: usize,
skipping: Vec<Test>,
skipped: usize,
row: i64,
given: u64,
}
#[derive(Debug)]
struct Piece {
file: usize,
reader: Option<FileReader>,
failure: Option<Error>,
row: i64,
}
impl FileScan {
#[allow(clippy::too_many_arguments)]
pub(crate) fn new(
plan: &Plan,
index: u32,
function: TableFunction,
args: Slice,
options: Slice,
settings: Slice,
columns: Slice,
tests: Vec<(usize, Op, Bound)>,
) -> Result<Self> {
let paths = file_arguments(plan, args, function)?;
let given = csv_options(plan, options, settings)?;
let produced = plan.field_list(columns).to_vec();
let numbered = produced.last().is_some_and(|field| field.name == FILE_ROW_NUMBER);
let wanted =
if numbered { produced[..produced.len() - 1].to_vec() } else { produced.clone() };
let scan = Self {
function,
paths,
given,
wanted,
numbered,
schema: Schema::numbered(produced, index),
tests,
cutting: Mutex::new(Cutting {
at: 0,
reader: None,
group: 0,
groups: 0,
part: 0,
cut: 0,
threads: 1,
pieces: 1,
skipping: Vec::new(),
skipped: 0,
row: 0,
given: 0,
}),
open: Mutex::new(HashMap::new()),
counters: None,
};
{
let mut cutting = scan.cutting.lock().map_err(poisoned)?;
scan.advance(&mut cutting)?;
}
Ok(scan)
}
pub(crate) fn watched(mut self, counters: Arc<Counters>) -> Self {
self.counters = Some(counters);
self
}
pub(crate) fn schema(&self) -> &Schema {
&self.schema
}
fn advance(&self, cutting: &mut Cutting) -> Result<()> {
cutting.reader = None;
cutting.row = 0;
cutting.group = 0;
cutting.groups = 0;
cutting.part = 0;
cutting.cut = 0;
cutting.pieces = 1;
cutting.skipping = Vec::new();
let Some(path) = self.paths.get(cutting.at) else { return Ok(()) };
let mut reader = FileReader::open(self.function, path, self.given)?;
let first = if cutting.at == 0 { None } else { self.paths.first().map(String::as_str) };
let held = positions(self.function, &self.wanted, &reader.fields(), path, first)?;
reader.project(&held)?;
reader.settle(&self.wanted)?;
cutting.groups = reader.row_groups();
cutting.skipping = self
.tests
.iter()
.filter_map(|(at, op, value)| {
Some(Test { column: *held.get(*at)?, op: *op, value: value.clone() })
})
.collect();
cutting.reader = Some(reader);
cutting.at += 1;
aim(cutting);
Ok(())
}
fn cut(&self, cutting: &mut Cutting) -> Result<Option<Piece>> {
let file = cutting.at.saturating_sub(1);
if let Some(FileReader::Parquet(reader)) = cutting.reader.as_ref() {
let metadata = reader.metadata();
while cutting.group < cutting.groups {
let Some(group) = metadata.row_groups.get(cutting.group) else { break };
if !skips(&cutting.skipping, group, &metadata.schema) {
break;
}
cutting.group += 1;
cutting.row = cutting.row.saturating_add(group.rows);
cutting.skipped = cutting.skipped.saturating_add(1);
}
}
let piece = match cutting.reader.as_ref() {
Some(FileReader::Parquet(reader)) if cutting.group < cutting.groups => {
let at = cutting.group;
let rows = group_rows(reader, at);
let piece = next_piece(rows, cutting.part, cutting.cut);
let split = reader.split_rows(at, piece.clone())?;
if piece.end >= rows {
cutting.group += 1;
cutting.part = 0;
} else {
cutting.part = piece.end;
}
let row = cutting.row;
cutting.row = cutting
.row
.saturating_add(i64::try_from(piece.end - piece.start).unwrap_or(i64::MAX));
Piece { file, reader: Some(FileReader::Parquet(split)), failure: None, row }
}
Some(FileReader::Csv(_)) => {
Piece { file, reader: cutting.reader.take(), failure: None, row: 0 }
}
_ => return Ok(None),
};
Ok(Some(piece))
}
fn hand(&self, cutting: &mut Cutting, piece: Piece) -> Option<Morsel> {
let index = cutting.given;
cutting.given += 1;
let covers = u64::from(piece.failure.is_none());
self.open.lock().ok()?.insert(index, Arc::new(Mutex::new(piece)));
Some(Morsel::new(index, 0, covers))
}
fn piece(&self, index: u64) -> Result<Arc<Mutex<Piece>>> {
let open = self.open.lock().map_err(poisoned)?;
open.get(&index)
.map(Arc::clone)
.ok_or_else(|| Error::internal(format!("a file scan was read at morsel {index}")))
}
fn number(&self, chunk: Chunk, piece: &mut Piece) -> Result<Chunk> {
let rows = chunk.len();
let mut columns = Vec::with_capacity(chunk.width() + 1);
for at in 0..chunk.width() {
columns.push(chunk.column(at)?.clone());
}
let first = piece.row;
piece.row = piece.row.saturating_add(i64::try_from(rows).unwrap_or(i64::MAX));
let mut data = Vec::with_capacity(rows);
for at in 0..rows {
data.push(first.saturating_add(i64::try_from(at).unwrap_or(i64::MAX)));
}
columns.push(Vector::flat(LogicalType::BigInt, Data::Int64(data.into()))?);
Chunk::with_rows(columns, rows)
}
fn conform(&self, chunk: Chunk, file: usize) -> Result<Chunk> {
let rows = chunk.len();
let settled = self.wanted.iter().enumerate().all(|(at, field)| {
chunk.column(at).is_ok_and(|column| column.logical_type() == &field.ty)
});
if settled {
return Ok(chunk);
}
let mut columns = Vec::with_capacity(self.wanted.len());
for (at, field) in self.wanted.iter().enumerate() {
let column = chunk.column(at)?;
if column.logical_type() == &field.ty {
columns.push(column.clone());
continue;
}
columns.push(cast(column, &field.ty, false).map_err(|error| {
let path = self.paths.get(file).map_or("", String::as_str);
Error::conversion(format!(
"Error while reading file \"{path}\": failed to cast column \"{}\" from type \
{} to {}: {}",
field.name,
column.logical_type(),
field.ty,
error.message()
))
})?);
}
Chunk::with_rows(columns, rows)
}
}
impl Source for FileScan {
fn morsel(&self) -> Option<Morsel> {
let mut cutting = self.cutting.lock().ok()?;
loop {
match self.cut(&mut cutting) {
Ok(Some(piece)) => return self.hand(&mut cutting, piece),
Ok(None) => {}
Err(error) => {
let file = cutting.at.saturating_sub(1);
let piece = Piece { file, reader: None, failure: Some(error), row: 0 };
return self.hand(&mut cutting, piece);
}
}
if cutting.at >= self.paths.len() {
return None;
}
if let Err(error) = self.advance(&mut cutting) {
let file = cutting.at.saturating_sub(1);
let piece = Piece { file, reader: None, failure: Some(error), row: 0 };
return self.hand(&mut cutting, piece);
}
}
}
fn morsels(&self, threads: usize) -> Option<usize> {
let mut cutting = self.cutting.lock().ok()?;
cutting.threads = threads;
aim(&mut cutting);
Some(cutting.pieces.max(1).saturating_mul(self.paths.len().max(1)))
}
fn read(&self, morsel: &mut Morsel, out: &mut Chunk) -> Result<Progress> {
let piece = self.piece(morsel.index())?;
let mut piece = piece.lock().map_err(poisoned)?;
if let Some(error) = piece.failure.take() {
return Err(error);
}
loop {
let file = piece.file;
let Some(reader) = piece.reader.as_mut() else {
morsel.advance(1);
*out = Chunk::empty(&self.schema.types());
return Ok(Progress::Done);
};
let before = reader.bytes_read();
let next = reader.next_chunk()?;
if let Some(counters) = &self.counters {
counters.read(reader.bytes_read().saturating_sub(before));
}
let Some(chunk) = next else {
piece.reader = None;
continue;
};
let mut chunk = self.conform(chunk, file)?;
if self.numbered {
chunk = self.number(chunk, &mut piece)?;
}
*out = chunk;
return Ok(Progress::More);
}
}
}
#[derive(Debug)]
enum FileReader {
Parquet(Reader),
Csv(CsvReader),
}
impl FileReader {
fn open(function: TableFunction, path: &str, given: Given) -> Result<Self> {
match function {
TableFunction::ReadCsv => Ok(Self::Csv(open_csv(path, given)?)),
_ => Ok(Self::Parquet(open_parquet(path)?)),
}
}
fn fields(&self) -> Vec<Field> {
match self {
Self::Parquet(reader) => reader.fields(),
Self::Csv(reader) => reader.fields(),
}
}
fn project(&mut self, columns: &[usize]) -> Result<()> {
match self {
Self::Parquet(reader) => reader.project(columns),
Self::Csv(reader) => reader.project(columns),
}
}
fn settle(&mut self, wanted: &[Field]) -> Result<()> {
match self {
Self::Parquet(reader) => {
let text: Vec<bool> =
wanted.iter().map(|field| field.ty == LogicalType::Varchar).collect();
reader.as_string(&text);
Ok(())
}
Self::Csv(reader) => {
let types: Vec<LogicalType> = wanted.iter().map(|field| field.ty.clone()).collect();
reader.retype(&types)
}
}
}
fn next_chunk(&mut self) -> Result<Option<Chunk>> {
match self {
Self::Parquet(reader) => reader.next_chunk(),
Self::Csv(reader) => reader.next_chunk(),
}
}
fn row_groups(&self) -> usize {
match self {
Self::Parquet(reader) => reader.metadata().row_groups.len(),
Self::Csv(_) => 0,
}
}
fn bytes_read(&self) -> u64 {
match self {
Self::Parquet(reader) => reader.bytes_read(),
Self::Csv(_) => 0,
}
}
}
fn csv_options(plan: &Plan, options: Slice, settings: Slice) -> Result<Given> {
if options.len == 0 {
return Ok(Given::default());
}
let exprs: Vec<ExprRef> = plan.expr_list(settings).to_vec();
let source = Schema::empty();
let one = Chunk::with_rows(Vec::new(), 1)?;
let evaluated = evaluate_all(plan, &exprs, &source, &one)?;
let names: Vec<&str> = plan.name_list(options).iter().map(|name| plan.string(*name)).collect();
let written: Vec<(&str, rudb_common::Value)> =
names.into_iter().zip(evaluated.iter().map(|vector| vector.value_at(0))).collect();
csv_given(&written)
}
pub(crate) fn file_arguments(
plan: &Plan,
args: Slice,
function: TableFunction,
) -> Result<Vec<String>> {
let exprs: Vec<ExprRef> = plan.expr_list(args).to_vec();
let source = Schema::empty();
let one = Chunk::with_rows(Vec::new(), 1)?;
let evaluated = evaluate_all(plan, &exprs, &source, &one)?;
let mut paths = Vec::with_capacity(evaluated.len());
for vector in &evaluated {
match vector.value_at(0) {
rudb_common::Value::Varchar(path) => paths.push(path),
other => {
return Err(Error::internal(format!(
"{}() bound with {other:?} rather than constant file names",
function.name()
)));
}
}
}
Ok(paths)
}
pub(crate) fn positions(
function: TableFunction,
wanted: &[Field],
held: &[Field],
path: &str,
first: Option<&str>,
) -> Result<Vec<usize>> {
let mut positions = Vec::with_capacity(wanted.len());
for field in wanted {
let at = held.iter().position(|column| column.name == field.name).ok_or_else(|| {
let Some(first) = first else {
return Error::io(format!(
"File \"{path}\" does not have a column named \"{}\"",
field.name
));
};
if matches!(function, TableFunction::ReadCsv) {
return rudb_csv::mismatch(first, path, &field.name);
}
let candidates: Vec<&str> = held.iter().map(|column| column.name.as_str()).collect();
Error::invalid_input(format!(
"Failed to read file \"{path}\": schema mismatch in glob: column \"{}\" was read \
from the original file \"{first}\", but could not be found in file \
\"{path}\".\nCandidate names: {}\nIf you are trying to read files with different \
schemas, try setting union_by_name=True",
field.name,
candidates.join(", ")
))
})?;
positions.push(at);
}
Ok(positions)
}
impl Source for Values {
fn morsel(&self) -> Option<Morsel> {
self.handout.take()
}
fn morsels(&self, _threads: usize) -> Option<usize> {
Some(self.handout.total())
}
fn read(&self, morsel: &mut Morsel, out: &mut Chunk) -> Result<Progress> {
*out = match self.chunks.get(position(morsel)) {
Some(chunk) => chunk.clone(),
None => Chunk::empty(&self.schema.types()),
};
morsel.advance(1);
Ok(Progress::Done)
}
}
#[cfg(test)]
mod tests {
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use rudb_catalog::{QualifiedName, Table};
use rudb_common::{Field, LogicalType, Value};
use rudb_functions::TableFunction;
use rudb_pipeline::{Progress, Source};
use rudb_plan::{Node, Plan};
use rudb_vector::Chunk;
use super::{
Bound, FileScan, Handout, Op, Probe, RUN, Scan, Schema, Series, VECTOR_SIZE, cut_rows,
next_piece, parts,
};
fn cutting(rows: usize, target: usize) -> Vec<std::ops::Range<usize>> {
let mut out = Vec::new();
let mut part = 0;
loop {
let piece = next_piece(rows, part, target);
if piece.is_empty() {
return out;
}
part = piece.end;
out.push(piece);
assert!(out.len() <= rows + 1, "cutting {rows} rows into {target} did not terminate");
}
}
fn series(start: i64, step: i64, rows: u64) -> Series {
Series { schema: Schema::empty(), start, step, rows, morsels: AtomicU64::new(0) }
}
fn drained(series: &Series) -> (Vec<i64>, usize) {
let mut values = Vec::new();
let mut morsels = 0;
while let Some(mut morsel) = series.morsel() {
morsels += 1;
loop {
let mut chunk = Chunk::empty(&[LogicalType::BigInt]);
let progress = series.read(&mut morsel, &mut chunk).expect("a series reads");
for row in 0..chunk.len() {
match chunk.value_at(row, 0) {
Value::BigInt(value) => values.push(value),
other => panic!("a series produced {other}"),
}
}
if progress == Progress::Done {
break;
}
}
}
(values, morsels)
}
#[test]
fn a_morsel_of_a_series_is_read_a_chunk_at_a_time() {
let rows = VECTOR_SIZE as u64 * 2 + 5;
let (values, morsels) = drained(&series(0, 1, rows));
assert_eq!(morsels, 1);
assert_eq!(values.len(), rows as usize);
assert_eq!(values[0], 0);
assert_eq!(values[values.len() - 1], rows as i64 - 1);
}
#[test]
fn a_series_longer_than_a_morsel_carries_on_where_the_last_one_stopped() {
let rows = RUN + 3;
let (values, morsels) = drained(&series(10, 3, rows));
assert_eq!(morsels, 2);
assert_eq!(values.len(), rows as usize);
assert_eq!(values[0], 10);
assert_eq!(values[RUN as usize], 10 + 3 * RUN as i64);
assert_eq!(values[values.len() - 1], 10 + 3 * (rows as i64 - 1));
}
#[test]
fn a_series_of_nothing_hands_out_no_work() {
let (values, morsels) = drained(&series(0, 1, 0));
assert!(values.is_empty());
assert_eq!(morsels, 0);
}
#[test]
fn a_handout_gives_each_position_to_one_caller_and_then_stops() {
let handout = Handout::new(3);
let taken: Vec<u64> = (0..3).map(|_| handout.take().expect("a position").start()).collect();
assert_eq!(taken, [0, 1, 2]);
assert!(handout.take().is_none());
assert!(handout.take().is_none());
}
fn fixture() -> FileScan {
pruned(Vec::new())
}
fn pruned(tests: Vec<(usize, Op, Bound)>) -> FileScan {
let path = format!("{}/../rudb-parquet/testdata/mixed.parquet", env!("CARGO_MANIFEST_DIR"));
let text = format!(
"TableFunction read_parquet args=['{path}'::VARCHAR] #0 [a::INTEGER, b::BIGINT]"
);
let plan = Plan::parse(&text).expect("the plan text round trips");
let Node::TableFunction { index, args, options, settings, columns, .. } =
*plan.node(plan.root())
else {
panic!("the plan is a table function");
};
FileScan::new(
&plan,
index,
TableFunction::ReadParquet,
args,
options,
settings,
columns,
tests,
)
.expect("the fixture is there")
}
fn morsels(scan: &FileScan) -> Vec<usize> {
let mut rows = Vec::new();
while let Some(mut morsel) = scan.morsel() {
let mut chunk = Chunk::empty(&[]);
let mut read = 0;
while let Progress::More = scan.read(&mut morsel, &mut chunk).expect("decodes") {
read += chunk.len();
}
rows.push(read);
}
rows
}
#[test]
fn a_parquet_file_is_one_morsel_per_row_group() {
let scan = fixture();
let rows = morsels(&scan);
assert_eq!(rows, [2048, 2048], "two row groups of 2048");
}
#[test]
fn a_row_group_whose_bounds_rule_out_the_filter_is_never_handed_out() {
let scan = pruned(vec![(0, Op::Greater, Bound::Int(1_000))]);
let rows = morsels(&scan);
assert_eq!(rows, Vec::<usize>::new(), "both row groups are ruled out");
let cutting = scan.cutting.lock().expect("the lock holds");
assert_eq!(cutting.skipped, 2, "and both were skipped rather than read");
}
#[test]
fn a_row_group_whose_bounds_overlap_the_filter_is_handed_out_as_usual() {
let scan = pruned(vec![(0, Op::Greater, Bound::Int(50))]);
let rows = morsels(&scan);
assert_eq!(rows, [2048, 2048], "nothing is ruled out");
let cutting = scan.cutting.lock().expect("the lock holds");
assert_eq!(cutting.skipped, 0);
}
#[test]
fn a_test_against_a_position_the_scan_does_not_produce_rules_nothing_out() {
let scan = pruned(vec![(7, Op::Greater, Bound::Int(1_000))]);
assert_eq!(morsels(&scan), [2048, 2048]);
}
#[test]
fn a_scan_that_has_handed_out_every_morsel_hands_out_no_more() {
let scan = fixture();
let _ = morsels(&scan);
assert!(scan.morsel().is_none());
assert!(scan.morsel().is_none());
}
#[test]
fn every_morsel_can_be_taken_before_any_of_them_is_read() {
let scan = fixture();
let mut taken = Vec::new();
while let Some(morsel) = scan.morsel() {
taken.push(morsel);
}
assert_eq!(taken.len(), 2);
let mut rows = Vec::new();
for morsel in &mut taken {
let mut chunk = Chunk::empty(&[]);
let mut read = 0;
while let Progress::More = scan.read(morsel, &mut chunk).expect("decodes") {
read += chunk.len();
}
rows.push(read);
}
assert_eq!(rows, [2048, 2048]);
}
fn counted(rows: usize) -> Table {
let mut table = Table::new(
QualifiedName::new("memory", "main", "t"),
vec![Field::new("n", LogicalType::Integer)],
)
.expect("one column");
let values: Vec<Vec<Value>> = (0..rows).map(|n| vec![Value::Integer(n as i32)]).collect();
table.append_rows(&values).expect("integers");
table
}
fn scanning(table: &Table, tests: Vec<(usize, Op, Bound)>) -> Scan<'_> {
let fields = vec![Field::new("n", LogicalType::Integer)];
let probes =
tests.into_iter().map(|(column, op, value)| Probe { column, op, value }).collect();
Scan {
table,
columns: vec![0],
probes,
schema: Schema::numbered(fields, 0),
chunks: Handout::new(table.rows().chunk_count()),
skipped: AtomicUsize::new(0),
}
}
fn counted_rows(scan: &Scan<'_>) -> usize {
let mut rows = 0;
while let Some(mut morsel) = scan.morsel() {
let mut chunk = Chunk::empty(&[LogicalType::Integer]);
loop {
let progress = scan.read(&mut morsel, &mut chunk).expect("a table scan reads");
rows += chunk.len();
if progress == Progress::Done {
break;
}
}
}
rows
}
#[test]
fn a_filter_on_a_table_reads_only_the_chunks_that_can_hold_a_match() {
let table = counted(VECTOR_SIZE * 5);
assert_eq!(table.rows().chunk_count(), 5);
let scan = scanning(&table, vec![(0, Op::Equal, Bound::Int(5_000))]);
assert_eq!(counted_rows(&scan), VECTOR_SIZE, "one chunk's worth");
assert_eq!(scan.skipped.load(Ordering::Relaxed), 4);
}
#[test]
fn a_scan_with_no_tests_reads_every_chunk() {
let table = counted(VECTOR_SIZE * 3);
let scan = scanning(&table, Vec::new());
assert_eq!(counted_rows(&scan), VECTOR_SIZE * 3);
assert_eq!(scan.skipped.load(Ordering::Relaxed), 0);
}
#[test]
fn conjuncts_that_rule_out_every_chunk_read_nothing() {
let table = counted(VECTOR_SIZE * 4);
let scan = scanning(
&table,
vec![(0, Op::GreaterOrEqual, Bound::Int(2_048)), (0, Op::Less, Bound::Int(2_048))],
);
assert_eq!(counted_rows(&scan), 0);
assert_eq!(scan.skipped.load(Ordering::Relaxed), 4);
}
#[test]
fn the_morsels_of_a_row_group_tile_it_once_each() {
for rows in [1, 2, 7, 9, 10, 100, 4_095, 4_096, 4_097, 122_880, 123_554] {
for target in [1, 2, 3, 300, 4_096, 32_768] {
let pieces = cutting(rows, target);
assert_eq!(pieces.len(), parts(rows, target), "{rows} rows in morsels of {target}");
let mut at = 0;
for piece in &pieces {
assert_eq!(piece.start, at, "{rows} rows in morsels of {target}");
assert!(piece.len() <= target, "{rows} rows in morsels of {target}");
at = piece.end;
}
assert_eq!(at, rows, "{rows} rows in morsels of {target}");
}
}
}
#[test]
fn a_row_group_is_cut_into_even_morsels() {
let pieces = cutting(123_554, 32_768);
assert_eq!(pieces.len(), 4);
let longest = pieces.iter().map(std::ops::Range::len).max().expect("four morsels");
let shortest = pieces.iter().map(std::ops::Range::len).min().expect("four morsels");
assert_eq!(longest, 30_889);
assert_eq!(shortest, 30_887);
assert!(longest - shortest < pieces.len(), "morsels of {longest} and {shortest} rows");
}
#[test]
fn a_row_group_no_larger_than_a_morsel_is_one_morsel() {
assert_eq!(cutting(2_048, 32_768), vec![0..2_048]);
assert_eq!(cutting(32_768, 32_768), vec![0..32_768]);
assert_eq!(parts(2_048, 32_768), 1);
}
#[test]
fn an_empty_row_group_is_no_morsels() {
assert!(next_piece(0, 0, 32_768).is_empty());
assert_eq!(cutting(0, 32_768), Vec::new());
assert_eq!(cutting(0, 0), Vec::new());
}
#[test]
fn a_row_group_that_is_not_worth_cutting_is_one_morsel() {
assert_eq!(cutting(123_554, 0), vec![0..123_554]);
assert_eq!(parts(123_554, 0), 1);
assert_eq!(parts(0, 0), 1);
}
#[test]
fn a_file_is_cut_only_as_finely_as_there_are_threads_to_want_it() {
let rows = 123_554;
let whole = 12_355_400;
for (threads, wanted) in [(1, 1), (8, 1), (9, 1), (16, 2), (32, 4), (64, 4)] {
let cut = cut_rows(1024, whole, rows, 9, threads);
assert_eq!(parts(rows, cut), wanted, "{threads} threads over nine row groups");
}
}
#[test]
fn a_file_of_coarse_pages_is_not_cut_however_many_threads_there_are() {
let rows = 123_554;
let whole = 12_355_400;
assert_eq!(cut_rows(whole, whole, rows, 9, 64), 0, "one page per chunk, as DuckDB writes");
assert_eq!(cut_rows(819_201, whole, rows, 9, 64), 0);
assert_eq!(cut_rows(0, whole, rows, 9, 64), 0);
assert_eq!(cut_rows(1024, 0, rows, 9, 64), 0, "a group of no bytes is never worth cutting");
assert!(cut_rows(819_200, whole, rows, 9, 64) > 0);
}
}