use std::sync::Arc;
use rudb_common::{Error, Field, LogicalType, Result};
use rudb_io::File;
use rudb_vector::{Chunk, VECTOR_SIZE};
use crate::convert::{self, Cells};
use crate::dialect::{self, Dialect, Given};
use crate::infer;
use crate::scan::{Records, Span};
const BLOCK: usize = 1 << 20;
const TAIL: usize = 64 << 10;
const BLOCK_ROWS: usize = 1024;
#[derive(Debug)]
pub struct Reader {
file: Arc<dyn File>,
path: String,
given: Given,
dialect: Dialect,
fields: Vec<Field>,
projection: Vec<usize>,
buffer: Vec<u8>,
at: usize,
offset: u64,
drained: bool,
line: u64,
scratch: Vec<String>,
records: Records,
block: usize,
origin: u64,
end: u64,
cap: u64,
}
impl Reader {
pub fn open(file: Box<dyn File>, path: &str) -> Result<Self> {
Self::open_with(file, path, Given::default())
}
pub fn open_with(file: Box<dyn File>, path: &str, given: Given) -> Result<Self> {
Self::open_sized(file, path, given, BLOCK)
}
pub(crate) fn open_sized(
file: Box<dyn File>,
path: &str,
given: Given,
block: usize,
) -> Result<Self> {
let mut reader = Self {
file: Arc::from(file),
path: path.to_string(),
given,
dialect: Dialect::comma_separated(),
fields: Vec::new(),
projection: Vec::new(),
buffer: Vec::new(),
at: 0,
offset: 0,
drained: false,
line: 1,
scratch: Vec::new(),
records: Records::default(),
block,
origin: 0,
end: u64::MAX,
cap: u64::MAX,
};
reader.fill(0)?;
let sample = reader.buffer.clone();
let quote = given.quote.or_else(|| dialect::quote(&sample));
let delimiter = match given.delimiter {
Some(byte) => byte,
None => dialect::delimiter(&sample, quote)?,
};
let escape = given.escape.or(quote);
reader.dialect = Dialect { delimiter, quote, escape, header: false };
let rows = reader.sample_rows(&sample)?;
let (header, fields) = describe(&rows, given.header);
reader.dialect.header = header;
reader.fields = fields;
reader.projection = (0..reader.fields.len()).collect();
if header {
reader.skip_record()?;
}
Ok(reader)
}
#[must_use]
pub fn fields(&self) -> Vec<Field> {
self.projection.iter().map(|&at| self.fields[at].clone()).collect()
}
pub fn project(&mut self, columns: &[usize]) -> Result<()> {
for &column in columns {
if column >= self.fields.len() {
return Err(Error::io(format!(
"column {column} is past the {} the file has",
self.fields.len()
)));
}
}
self.projection = columns.to_vec();
Ok(())
}
pub fn retype(&mut self, types: &[LogicalType]) -> Result<()> {
if types.len() != self.projection.len() {
return Err(Error::io(format!(
"{} types for a projection of {} columns",
types.len(),
self.projection.len()
)));
}
for (&at, ty) in self.projection.iter().zip(types) {
self.fields[at].ty = ty.clone();
}
Ok(())
}
#[must_use]
pub const fn dialect(&self) -> Dialect {
self.dialect
}
pub fn next_chunk(&mut self) -> Result<Option<Chunk>> {
let rows = self.next_records()?;
if rows == 0 {
return Ok(None);
}
let first = self.line;
self.line += rows as u64;
let cells = Cells { bytes: &self.buffer, records: &self.records, dialect: self.dialect };
let projected: Vec<_> =
self.projection.iter().map(|&at| (at, &self.fields[at].ty)).collect();
let mut builders = convert::builders(&cells, &projected);
let mut start = 0;
while start < rows {
let end = rows.min(start + BLOCK_ROWS);
for (build, &at) in builders.iter_mut().zip(&self.projection) {
let field = &self.fields[at];
let refuse = |text: &str, row: usize| {
Error::conversion(self.conversion_error(text, field, first + row as u64))
};
if let Err(error) = build.rows(&cells, at, start..end, &refuse) {
return Err(self.first_bad_value(&cells, first).unwrap_or(error));
}
}
start = end;
}
let columns = builders.into_iter().map(|build| build.finish()).collect::<Result<_>>()?;
Ok(Some(Chunk::with_rows(columns, rows)?))
}
fn first_bad_value(&self, cells: &Cells<'_>, first: u64) -> Option<Error> {
self.projection.iter().find_map(|&at| {
let field = &self.fields[at];
let refuse = |text: &str, row: usize| {
Error::conversion(self.conversion_error(text, field, first + row as u64))
};
convert::column(cells, at, &field.ty, &refuse).err()
})
}
fn next_records(&mut self) -> Result<usize> {
self.records.clear();
let mut start = self.at;
let mut careful = false;
loop {
if self.here() >= self.end {
break;
}
let limit = if careful { self.records.len() + 1 } else { VECTOR_SIZE };
self.at = crate::scan::records(
&self.buffer,
self.at,
self.dialect,
self.drained,
limit,
&mut self.records,
)?;
if !careful && self.here() > self.end {
self.records.clear();
self.at = start;
careful = true;
continue;
}
if self.records.len() == VECTOR_SIZE {
break;
}
if careful && self.records.len() == limit {
continue;
}
if self.drained {
break;
}
if self.buffer.len() - start + self.block > Span::MOST {
if self.records.is_empty() {
return Err(Error::io("a record is longer than two gigabytes"));
}
break;
}
self.fill(start)?;
self.records.shift(start);
start = 0;
}
Ok(self.records.len())
}
pub(crate) fn here(&self) -> u64 {
self.offset - self.buffer.len() as u64 + self.at as u64
}
#[must_use]
pub fn bytes_read(&self) -> u64 {
self.offset - self.origin
}
pub(crate) fn stretch(&self, from: u64, end: u64, line: u64) -> Self {
Self {
file: Arc::clone(&self.file),
path: self.path.clone(),
given: self.given,
dialect: self.dialect,
fields: self.fields.clone(),
projection: self.projection.clone(),
buffer: Vec::new(),
at: 0,
offset: from,
drained: false,
line,
scratch: Vec::new(),
records: Records::default(),
block: self.block,
origin: from,
end,
cap: u64::MAX,
}
}
pub(crate) fn give_up_at(&mut self, cap: u64) {
self.cap = cap;
}
pub(crate) fn skim(&mut self) -> Result<u64> {
while self.next_records()? > 0 {}
Ok(self.here())
}
pub(crate) fn skip_to(&mut self, target: u64) -> Result<()> {
loop {
let from = self.here();
let rows = self.next_records()?;
if rows == 0 {
return Ok(());
}
if self.here() > target {
let front = self.offset - self.buffer.len() as u64;
self.at = usize::try_from(from - front)
.map_err(|_| Error::internal("a chunk start outside the buffer"))?;
return Ok(());
}
self.line += rows as u64;
}
}
pub(crate) const fn line(&self) -> u64 {
self.line
}
pub(crate) fn file(&self) -> &dyn File {
self.file.as_ref()
}
fn conversion_error(&self, text: &str, field: &Field, line: u64) -> String {
format!(
"CSV Error on Line: {line}\nOriginal Line: {text}\nError when converting column \
\"{}\". Could not convert string \"{text}\" to '{}'\n\nColumn {} is being converted \
as type {}\nThis type was auto-detected from the CSV file.\nPossible solutions:\n* \
Override the type for this column manually by setting the type explicitly, e.g., \
types={{'{}': 'VARCHAR'}}\n* Set the sample size to a larger value to enable the \
auto-detection to scan more values, e.g., sample_size=-1\n* Use a COPY statement to \
automatically derive types from an existing table.\n* Check whether the null string \
value is set correctly (e.g., nullstr = 'N/A')\n\n file = {}\n delimiter = {}\n \
quote = {}\n escape = {}\n header = {} {}\n sample_size = {}\n",
field.name,
field.ty,
field.name,
field.ty,
field.name,
self.path,
Given::shown(self.given.delimiter, Some(self.dialect.delimiter)),
Given::shown(self.given.quote, self.dialect.quote),
Given::shown(self.given.escape, self.dialect.escape),
self.dialect.header,
Given::source(self.given.header.is_some()),
infer::SAMPLE,
)
}
#[cfg(test)]
fn next_chunk_by_record(&mut self) -> Result<Option<Chunk>> {
let mut rows: Vec<Vec<Option<String>>> = Vec::new();
while rows.len() < VECTOR_SIZE {
match self.next_record()? {
Some(fields) => rows.push(fields),
None => break,
}
}
if rows.is_empty() {
return Ok(None);
}
let mut columns = Vec::with_capacity(self.projection.len());
for &at in &self.projection {
let field = &self.fields[at];
let mut values = Vec::with_capacity(rows.len());
for (row, held) in rows.iter().enumerate() {
let text = held.get(at).and_then(Option::as_deref);
values.push(self.convert(
text,
field,
self.line - rows.len() as u64 + row as u64,
)?);
}
columns.push(rudb_vector::Vector::from_values(field.ty.clone(), &values)?);
}
Ok(Some(Chunk::with_rows(columns, rows.len())?))
}
#[cfg(test)]
fn convert(&self, text: Option<&str>, field: &Field, line: u64) -> Result<rudb_common::Value> {
let Some(text) = text else { return Ok(rudb_common::Value::Null) };
if field.ty == LogicalType::Varchar {
return Ok(rudb_common::Value::Varchar(text.to_string()));
}
let value = rudb_common::Value::Varchar(text.to_string());
match rudb_kernels::cast_value(&value, &field.ty, false) {
Ok(converted) => Ok(converted),
Err(_) => Err(Error::conversion(self.conversion_error(text, field, line))),
}
}
#[cfg(test)]
fn next_record(&mut self) -> Result<Option<Vec<Option<String>>>> {
let Some(()) = self.advance()? else { return Ok(None) };
Ok(Some(
self.scratch
.iter()
.map(|text| if text.is_empty() { None } else { Some(text.clone()) })
.collect(),
))
}
fn advance(&mut self) -> Result<Option<()>> {
loop {
let mut scratch = std::mem::take(&mut self.scratch);
let outcome = crate::scan::record(
&self.buffer,
self.at,
self.dialect,
self.drained,
&mut scratch,
);
self.scratch = scratch;
match outcome? {
Some(next) => {
self.at = next;
self.line += 1;
return Ok(Some(()));
}
None if self.drained => return Ok(None),
None => self.fill(self.at)?,
}
}
}
fn skip_record(&mut self) -> Result<()> {
self.advance()?;
Ok(())
}
fn fill(&mut self, keep: usize) -> Result<()> {
if self.offset >= self.cap {
return Err(Error::io("a record runs further than a guessed start is followed"));
}
self.buffer.drain(..keep);
self.at -= keep;
let held = self.buffer.len();
let want = match self.end.checked_sub(self.offset) {
Some(left) if left > 0 => self.block.min(usize::try_from(left).unwrap_or(usize::MAX)),
_ => self.block.min(TAIL.max(held)),
};
self.buffer.resize(held + want, 0);
let read = self.file.read_at(self.offset, &mut self.buffer[held..])?;
self.buffer.truncate(held + read);
self.offset += read as u64;
if read == 0 {
self.drained = true;
}
Ok(())
}
fn sample_rows(&self, sample: &[u8]) -> Result<Vec<Vec<Option<String>>>> {
let mut rows = Vec::new();
let mut fields = Vec::new();
let mut at = 0;
while rows.len() <= infer::SAMPLE {
let Some(next) = crate::scan::record(sample, at, self.dialect, false, &mut fields)?
else {
break;
};
at = next;
rows.push(
fields
.iter()
.map(|text| if text.is_empty() { None } else { Some(text.clone()) })
.collect(),
);
}
Ok(rows)
}
}
fn describe(rows: &[Vec<Option<String>>], told: Option<bool>) -> (bool, Vec<Field>) {
let width = rows.iter().map(Vec::len).max().unwrap_or(0);
let body = types(&rows[1.min(rows.len())..], width);
let all_text = body.iter().all(|ty| *ty == LogicalType::Varchar);
let first_fits = rows.first().is_some_and(|first| {
first.iter().zip(&body).all(|(text, ty)| match text {
None => true,
Some(text) => infer::fits(text, ty),
})
});
let header = !rows.is_empty() && told.unwrap_or(rows.len() > 1 && (all_text || !first_fits));
if !header {
let types = types(rows, width);
let fields = types
.into_iter()
.enumerate()
.map(|(at, ty)| Field::new(format!("column{at}"), ty))
.collect();
return (false, fields);
}
let names = unique(&rows[0], width);
let fields = body.into_iter().zip(names).map(|(ty, name)| Field::new(name, ty)).collect();
(true, fields)
}
fn unique(header: &[Option<String>], width: usize) -> Vec<String> {
let mut taken: Vec<String> = Vec::with_capacity(width);
for at in 0..width {
let base = match header.get(at).and_then(Option::as_deref) {
Some(written) => written.to_string(),
None => format!("column{at}"),
};
let mut name = base.clone();
let mut next = 1;
while taken.iter().any(|held| held.eq_ignore_ascii_case(&name)) {
name = format!("{base}_{next}");
next += 1;
}
taken.push(name);
}
taken
}
fn types(rows: &[Vec<Option<String>>], width: usize) -> Vec<LogicalType> {
(0..width)
.map(|at| {
let values: Vec<Option<&str>> =
rows.iter().map(|row| row.get(at).and_then(Option::as_deref)).collect();
infer::column(&values)
})
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
use rudb_common::Value;
use rudb_io::{Filesystem, OpenMode, SimFilesystem};
use std::path::Path;
fn read(text: &str) -> Reader {
let filesystem = SimFilesystem::new();
let path = Path::new("/t.csv");
let file = filesystem.open(path, OpenMode::Create).expect("creates");
file.write_at(0, text.as_bytes()).expect("writes");
drop(file);
let file = filesystem.open(path, OpenMode::Read).expect("opens");
Reader::open(file, "/t.csv").expect("sniffs")
}
fn read_with(text: &str, given: Given) -> Reader {
let filesystem = SimFilesystem::new();
let path = Path::new("/t.csv");
let file = filesystem.open(path, OpenMode::Create).expect("creates");
file.write_at(0, text.as_bytes()).expect("writes");
drop(file);
let file = filesystem.open(path, OpenMode::Read).expect("opens");
Reader::open_with(file, "/t.csv", given).expect("reads")
}
fn names_and_types(reader: &Reader) -> Vec<(String, String)> {
reader.fields().iter().map(|f| (f.name.clone(), f.ty.to_string())).collect()
}
fn all(reader: &mut Reader) -> Vec<Vec<Value>> {
let mut rows = Vec::new();
while let Some(chunk) = reader.next_chunk().expect("reads") {
for row in 0..chunk.len() {
rows.push((0..chunk.width()).map(|at| chunk.value_at(row, at)).collect());
}
}
rows
}
#[test]
fn a_header_that_names_two_columns_the_same_thing_counts_the_second_one_up() {
let names: Vec<String> =
read("a,a,A,a_1\n1,2,3,4\nx,y,z,w\n").fields().into_iter().map(|f| f.name).collect();
assert_eq!(names, ["a", "a_1", "A_2", "a_1_1"]);
}
#[test]
fn a_header_and_three_types_are_what_duckdb_sniffs_for_the_same_bytes() {
let reader = read("a,b,c\n1,x,2.5\n2,y,3.5\n");
assert_eq!(
names_and_types(&reader),
[
("a".to_string(), "BIGINT".to_string()),
("b".to_string(), "VARCHAR".to_string()),
("c".to_string(), "DOUBLE".to_string()),
]
);
}
#[test]
fn a_file_with_no_header_gets_the_names_duckdb_gives_it() {
let reader = read("1,x\n2,y\n");
assert_eq!(
names_and_types(&reader),
[
("column0".to_string(), "BIGINT".to_string()),
("column1".to_string(), "VARCHAR".to_string()),
]
);
}
#[test]
fn two_rows_of_words_are_a_header_and_a_row() {
let reader = read("a,b\nc,d\n");
assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["a", "b"]);
}
#[test]
fn one_column_of_words_under_a_row_of_numbers_is_still_a_header() {
let reader = read("a,2\n3,4\n");
assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["a", "2"]);
}
#[test]
fn the_rows_are_the_rows_of_the_file() {
let mut reader = read("a,b\n1,x\n2,y\n");
assert_eq!(
all(&mut reader),
[
vec![Value::BigInt(1), Value::Varchar("x".into())],
vec![Value::BigInt(2), Value::Varchar("y".into())],
]
);
}
#[test]
fn an_empty_field_is_a_null_whether_it_was_quoted_or_not() {
let mut reader = read("a,b\n1,\n\"\",y\n");
assert_eq!(
all(&mut reader),
[vec![Value::BigInt(1), Value::Null], vec![Value::Null, Value::Varchar("y".into())],]
);
}
#[test]
fn a_projection_picks_columns_out_by_position_and_can_reorder_them() {
let mut reader = read("a,b,c\n1,x,2.5\n");
reader.project(&[2, 0]).expect("projects");
assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["c", "a"]);
assert_eq!(all(&mut reader), [vec![Value::Double(2.5), Value::BigInt(1)]]);
}
#[test]
fn a_projection_of_nothing_still_counts_the_rows() {
let mut reader = read("a,b\n1,x\n2,y\n3,z\n");
reader.project(&[]).expect("projects");
let chunk = reader.next_chunk().expect("reads").expect("a chunk");
assert_eq!(chunk.len(), 3);
assert_eq!(chunk.width(), 0);
}
#[test]
fn a_pipe_separated_file_reads_as_one() {
let mut reader = read("a|b\n1|x\n");
assert_eq!(reader.dialect().delimiter, b'|');
assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x".into())]]);
}
#[test]
fn a_quoted_field_with_a_delimiter_in_it_is_one_value() {
let mut reader = read("a,b\n1,\"x,y\"\n");
assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x,y".into())]]);
}
#[test]
fn more_rows_than_fit_one_chunk_arrive_as_more_than_one_chunk() {
let mut text = String::from("a\n");
for row in 0..VECTOR_SIZE + 5 {
text.push_str(&format!("{row}\n"));
}
let mut reader = read(&text);
let first = reader.next_chunk().expect("reads").expect("a chunk");
assert_eq!(first.len(), VECTOR_SIZE);
let second = reader.next_chunk().expect("reads").expect("a second chunk");
assert_eq!(second.len(), 5);
assert!(reader.next_chunk().expect("reads").is_none());
}
#[test]
fn a_value_the_sniffer_never_saw_is_an_error_rather_than_a_wider_column() {
let mut text = String::from("c\n");
for row in 0..infer::SAMPLE {
text.push_str(&format!("{row}\n"));
}
text.push_str("oops\n");
let mut reader = read(&text);
assert_eq!(reader.fields()[0].ty, LogicalType::BigInt);
let error = all_or_error(&mut reader).unwrap_err();
let line = infer::SAMPLE + 2;
assert!(error.message().starts_with(&format!("CSV Error on Line: {line}")), "{error}");
assert!(
error.message().contains("Could not convert string \"oops\" to 'BIGINT'"),
"{error}"
);
assert!(error.message().contains("sample_size = 20480"), "{error}");
}
#[test]
fn the_error_names_the_first_columns_bad_value_even_when_a_later_column_is_bad_sooner() {
let mut text = String::from("a,b\n");
for row in 0..infer::SAMPLE {
text.push_str(&format!("{row},{row}\n"));
}
text.push_str("1,late\n");
for row in 0..BLOCK_ROWS * 2 {
text.push_str(&format!("{row},{row}\n"));
}
text.push_str("early,1\n");
let mut reader = read(&text);
let error = all_or_error(&mut reader).unwrap_err();
assert!(error.message().contains("Could not convert string \"early\""), "{error}");
}
#[test]
fn a_file_told_it_has_no_header_reads_its_first_line_as_a_row() {
let mut reader =
read_with("a,b\n1,x\n2,y\n", Given { header: Some(false), ..Given::default() });
assert_eq!(
names_and_types(&reader),
[
("column0".to_string(), "VARCHAR".to_string()),
("column1".to_string(), "VARCHAR".to_string()),
]
);
assert_eq!(all(&mut reader).len(), 3);
}
#[test]
fn a_file_told_it_has_a_header_takes_its_first_line_as_the_names() {
let reader = read_with("1,2\n3,4\n", Given { header: Some(true), ..Given::default() });
assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["1", "2"]);
}
#[test]
fn a_given_delimiter_is_the_delimiter_whatever_the_file_looks_like() {
let reader = read_with("a,b\n1,x\n", Given { delimiter: Some(b';'), ..Given::default() });
assert_eq!(reader.dialect().delimiter, b';');
assert_eq!(reader.fields().len(), 1);
}
#[test]
fn a_given_quote_makes_a_field_that_holds_the_delimiter_one_value() {
let mut reader =
read_with("a,b\n1,'x,y'\n", Given { quote: Some(b'\''), ..Given::default() });
assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x,y".into())]]);
}
#[test]
fn the_block_says_set_by_user_for_what_the_call_gave_it() {
let mut text = String::from("c;d\n");
for row in 0..infer::SAMPLE {
text.push_str(&format!("{row};x\n"));
}
text.push_str("oops;x\n");
let given = Given { delimiter: Some(b';'), ..Given::default() };
let mut reader = read_with(&text, given);
let error = all_or_error(&mut reader).unwrap_err();
assert!(error.message().contains("delimiter = ; (Set By User)"), "{error}");
assert!(error.message().contains("header = true (Auto-Detected)"), "{error}");
}
#[test]
fn a_file_told_a_wider_type_than_it_sniffed_reads_its_whole_numbers_as_that_type() {
let mut reader = read("a\n1\n2\n");
assert_eq!(reader.fields()[0].ty, LogicalType::BigInt);
reader.retype(&[LogicalType::Double]).expect("one type for one column");
assert_eq!(reader.fields()[0].ty, LogicalType::Double);
assert_eq!(all(&mut reader), [[Value::Double(1.0)], [Value::Double(2.0)]]);
}
#[test]
fn a_type_list_that_is_not_as_long_as_the_projection_is_refused() {
let mut reader = read("a,b\n1,two\n");
let error = reader.retype(&[LogicalType::Double]).unwrap_err();
assert!(error.message().contains("1 types for a projection of 2 columns"), "{error}");
}
fn drained(reader: &mut Reader, old: bool) -> (Vec<(usize, Vec<String>)>, Option<String>) {
let mut chunks = Vec::new();
loop {
let next = if old { reader.next_chunk_by_record() } else { reader.next_chunk() };
match next {
Ok(Some(chunk)) => {
let mut values = Vec::new();
for row in 0..chunk.len() {
for at in 0..chunk.width() {
values.push(format!("{:?}", chunk.value_at(row, at)));
}
}
chunks.push((chunk.len(), values));
}
Ok(None) => return (chunks, None),
Err(error) => return (chunks, Some(error.to_string())),
}
}
}
struct Rng(u64);
impl Rng {
fn next(&mut self) -> u64 {
self.0 ^= self.0 << 13;
self.0 ^= self.0 >> 7;
self.0 ^= self.0 << 17;
self.0
}
fn below(&mut self, n: usize) -> usize {
(self.next() % n as u64) as usize
}
}
fn typed_file(rng: &mut Rng, rows: usize) -> String {
const ODD: [&str; 27] = [
"",
" 1",
"1 ",
"1e3",
"0x10",
"inf",
"-nan",
"abc",
"\"12\"",
"\"a\"\"b\"",
"\"x,y\"",
"\"x\ny\"",
"h\u{e9}llo",
"99999999999999999999",
"9999999999999999999",
"-",
"+5",
"007",
"2020-02-30",
"2020-02-29",
"0000-01-01",
"TRUE",
"no",
"1_000",
"1.5e-3",
"-0",
"\"\"",
];
let width = 1 + rng.below(6);
let kinds: Vec<usize> = (0..width).map(|_| rng.below(5)).collect();
let mut text: String = (0..width).map(|at| format!("c{at}")).collect::<Vec<_>>().join(",");
text.push('\n');
for _ in 0..rows {
let mut fields = Vec::with_capacity(width);
for &kind in &kinds {
let odd = rng.below(60) == 0;
fields.push(if odd {
ODD[rng.below(ODD.len())].to_string()
} else {
let n = rng.next();
match kind {
0 => format!("{}", (n % 2_000_001) as i64 - 1_000_000),
1 => format!("{}.{:02}", n % 100_000, n % 100),
2 => format!("{}-{:02}-{:02}", 1990 + n % 20, 1 + n % 12, 1 + n % 28),
3 => ["true", "false", "t", "F"][(n % 4) as usize].to_string(),
_ => ["x", "hello world", "a longer piece of text", "\"q,\"\"q\""]
[(n % 4) as usize]
.to_string(),
}
});
}
if rng.below(200) == 0 {
fields.pop();
}
if rng.below(200) == 0 {
fields.push("extra".to_string());
}
text.push_str(&fields.join(","));
text.push_str(["\n", "\n", "\n", "\r\n", "\r"][rng.below(5)]);
}
if rng.below(4) == 0 {
text.pop();
}
text
}
fn open_sized(text: &[u8], block: usize) -> Result<Reader> {
let filesystem = SimFilesystem::new();
let path = Path::new("/t.csv");
let file = filesystem.open(path, OpenMode::Create).expect("creates");
file.write_at(0, text).expect("writes");
drop(file);
let file = filesystem.open(path, OpenMode::Read).expect("opens");
Reader::open_sized(file, "/t.csv", Given::default(), block)
}
#[test]
fn generated_files_read_the_same_a_chunk_at_a_time_as_a_record_at_a_time() {
let types = [
LogicalType::BigInt,
LogicalType::Integer,
LogicalType::SmallInt,
LogicalType::TinyInt,
LogicalType::UBigInt,
LogicalType::UInteger,
LogicalType::USmallInt,
LogicalType::UTinyInt,
LogicalType::Double,
LogicalType::Float,
LogicalType::Date,
LogicalType::Boolean,
LogicalType::Varchar,
LogicalType::Timestamp,
LogicalType::Decimal { width: 18, scale: 3 },
];
let mut rng = Rng(0x2545_f491_4f6c_dd1d);
for case in 0..200 {
let rows = if case % 100 == 0 { 8192 + rng.below(1000) } else { rng.below(200) };
let mut text = typed_file(&mut rng, rows).into_bytes();
if rng.below(20) == 0 {
let at = rng.below(text.len() + 1);
text.splice(at..at, *b",\"x\"y,");
}
for block in [1 << 20, 32 + rng.below(400)] {
let (Ok(mut new), Ok(mut old)) =
(open_sized(&text, block), open_sized(&text, block))
else {
continue;
};
assert_eq!(new.fields(), old.fields());
let width = new.fields().len();
if width > 0 && rng.below(2) == 0 {
let columns: Vec<usize> =
(0..rng.below(width + 2)).map(|_| rng.below(width)).collect();
new.project(&columns).expect("projects");
old.project(&columns).expect("projects");
let wanted: Vec<LogicalType> =
columns.iter().map(|_| types[rng.below(types.len())].clone()).collect();
new.retype(&wanted).expect("retypes");
old.retype(&wanted).expect("retypes");
}
let expected = drained(&mut old, true);
let found = drained(&mut new, false);
assert_eq!(found.1, expected.1, "case {case}, block {block}");
assert_eq!(found.0, expected.0, "case {case}, block {block}");
}
}
}
#[test]
#[ignore = "a measurement, run by hand"]
fn reads_lineitem_faster_a_chunk_at_a_time() {
use rudb_io::RealFilesystem;
use std::time::Instant;
const ROWS: usize = 200_000;
let mut rng = Rng(0x1234_5678_9abc_def1);
let mut text = String::from(
"l_orderkey,l_partkey,l_suppkey,l_linenumber,l_quantity,l_extendedprice,l_discount,\
l_tax,l_returnflag,l_linestatus,l_shipdate,l_commitdate,l_receiptdate,\
l_shipinstruct,l_shipmode,l_comment\n",
);
let words = ["carefully", "final", "deposits", "furiously", "regular", "ideas", "sleep"];
for row in 0..ROWS {
let n = rng.next();
let date = |shift: u64| {
format!(
"{}-{:02}-{:02}",
1992 + (n >> shift) % 7,
1 + (n >> shift) % 12,
1 + (n >> shift) % 28
)
};
let comment: Vec<&str> =
(0..3 + n % 4).map(|k| words[((n >> (k * 3)) % 7) as usize]).collect();
text.push_str(&format!(
"{},{},{},{},{}.00,{}.{:02},0.0{},0.0{},{},{},{},{},{},{},{},{}\n",
row / 4 + 1,
n % 200_000,
n % 10_000,
row % 4 + 1,
1 + n % 50,
900 + n % 100_000,
n % 100,
n % 10,
(n >> 7) % 9,
["A", "N", "R"][(n % 3) as usize],
["O", "F"][(n % 2) as usize],
date(3),
date(11),
date(19),
["DELIVER IN PERSON", "NONE", "TAKE BACK RETURN"][(n % 3) as usize],
["TRUCK", "MAIL", "AIR", "SHIP"][(n % 4) as usize],
comment.join(" "),
));
}
let path = std::env::temp_dir().join(format!("rudb-lineitem-{}.csv", std::process::id()));
std::fs::write(&path, &text).expect("writes");
let megabytes = text.len() as f64 / 1e6;
let filesystem = RealFilesystem::new();
let mut best = [f64::MAX; 2];
for _ in 0..5 {
for (slot, old) in [(0, true), (1, false)] {
let file = filesystem.open(&path, OpenMode::Read).expect("opens");
let mut reader = Reader::open(file, "lineitem.csv").expect("sniffs");
let started = Instant::now();
let mut rows = 0;
loop {
let next =
if old { reader.next_chunk_by_record() } else { reader.next_chunk() };
let Some(chunk) = next.expect("reads") else { break };
rows += chunk.len();
}
assert_eq!(rows, ROWS);
best[slot] = best[slot].min(started.elapsed().as_secs_f64());
}
}
std::fs::remove_file(&path).expect("removes");
let (old, new) = (megabytes / best[0], megabytes / best[1]);
println!(
"{ROWS} rows, {megabytes:.1} MB: a record at a time {old:.1} MB/s, a chunk at a time"
);
println!("{new:.1} MB/s, {:.2}x", best[0] / best[1]);
}
fn all_or_error(reader: &mut Reader) -> Result<Vec<Vec<Value>>> {
let mut rows = Vec::new();
while let Some(chunk) = reader.next_chunk()? {
for row in 0..chunk.len() {
rows.push((0..chunk.width()).map(|at| chunk.value_at(row, at)).collect());
}
}
Ok(rows)
}
}