1use std::sync::Arc;
17
18use rudb_common::{Error, Field, LogicalType, Result};
19use rudb_io::File;
20use rudb_vector::{Chunk, VECTOR_SIZE};
21
22use crate::convert::{self, Cells};
23use crate::dialect::{self, Dialect, Given};
24use crate::infer;
25use crate::scan::{Records, Span};
26
27const BLOCK: usize = 1 << 20;
33
34const TAIL: usize = 64 << 10;
40
41#[derive(Debug)]
43pub struct Reader {
44 file: Arc<dyn File>,
45 path: String,
46 given: Given,
47 dialect: Dialect,
48 fields: Vec<Field>,
49 projection: Vec<usize>,
50 buffer: Vec<u8>,
51 at: usize,
52 offset: u64,
53 drained: bool,
54 line: u64,
55 scratch: Vec<String>,
56 records: Records,
57 block: usize,
58 origin: u64,
60 end: u64,
62 cap: u64,
65}
66
67impl Reader {
68 pub fn open(file: Box<dyn File>, path: &str) -> Result<Self> {
77 Self::open_with(file, path, Given::default())
78 }
79
80 pub fn open_with(file: Box<dyn File>, path: &str, given: Given) -> Result<Self> {
90 Self::open_sized(file, path, given, BLOCK)
91 }
92
93 pub(crate) fn open_sized(
96 file: Box<dyn File>,
97 path: &str,
98 given: Given,
99 block: usize,
100 ) -> Result<Self> {
101 let mut reader = Self {
102 file: Arc::from(file),
103 path: path.to_string(),
104 given,
105 dialect: Dialect::comma_separated(),
106 fields: Vec::new(),
107 projection: Vec::new(),
108 buffer: Vec::new(),
109 at: 0,
110 offset: 0,
111 drained: false,
112 line: 1,
113 scratch: Vec::new(),
114 records: Records::default(),
115 block,
116 origin: 0,
117 end: u64::MAX,
118 cap: u64::MAX,
119 };
120 reader.fill(0)?;
121 let sample = reader.buffer.clone();
122 let quote = given.quote.or_else(|| dialect::quote(&sample));
123 let delimiter = match given.delimiter {
124 Some(byte) => byte,
125 None => dialect::delimiter(&sample, quote)?,
126 };
127 let escape = given.escape.or(quote);
128 reader.dialect = Dialect { delimiter, quote, escape, header: false };
129 let rows = reader.sample_rows(&sample)?;
130 let (header, fields) = describe(&rows, given.header);
131 reader.dialect.header = header;
132 reader.fields = fields;
133 reader.projection = (0..reader.fields.len()).collect();
134 if header {
135 reader.skip_record()?;
136 }
137 Ok(reader)
138 }
139
140 #[must_use]
142 pub fn fields(&self) -> Vec<Field> {
143 self.projection.iter().map(|&at| self.fields[at].clone()).collect()
144 }
145
146 pub fn project(&mut self, columns: &[usize]) -> Result<()> {
152 for &column in columns {
153 if column >= self.fields.len() {
154 return Err(Error::io(format!(
155 "column {column} is past the {} the file has",
156 self.fields.len()
157 )));
158 }
159 }
160 self.projection = columns.to_vec();
161 Ok(())
162 }
163
164 pub fn retype(&mut self, types: &[LogicalType]) -> Result<()> {
182 if types.len() != self.projection.len() {
183 return Err(Error::io(format!(
184 "{} types for a projection of {} columns",
185 types.len(),
186 self.projection.len()
187 )));
188 }
189 for (&at, ty) in self.projection.iter().zip(types) {
190 self.fields[at].ty = ty.clone();
191 }
192 Ok(())
193 }
194
195 #[must_use]
197 pub const fn dialect(&self) -> Dialect {
198 self.dialect
199 }
200
201 pub fn next_chunk(&mut self) -> Result<Option<Chunk>> {
213 let rows = self.next_records()?;
214 if rows == 0 {
215 return Ok(None);
216 }
217 let first = self.line;
218 self.line += rows as u64;
219 let cells = Cells { bytes: &self.buffer, records: &self.records, dialect: self.dialect };
220 let mut columns = Vec::with_capacity(self.projection.len());
221 for &at in &self.projection {
222 let field = &self.fields[at];
223 let refuse = |text: &str, row: usize| {
224 Error::conversion(self.conversion_error(text, field, first + row as u64))
225 };
226 columns.push(convert::column(&cells, at, &field.ty, &refuse)?);
227 }
228 Ok(Some(Chunk::with_rows(columns, rows)?))
229 }
230
231 fn next_records(&mut self) -> Result<usize> {
239 self.records.clear();
240 let mut start = self.at;
241 let mut careful = false;
242 loop {
243 if self.here() >= self.end {
244 break;
245 }
246 let limit = if careful { self.records.len() + 1 } else { VECTOR_SIZE };
247 self.at = crate::scan::records(
248 &self.buffer,
249 self.at,
250 self.dialect,
251 self.drained,
252 limit,
253 &mut self.records,
254 )?;
255 if !careful && self.here() > self.end {
256 self.records.clear();
257 self.at = start;
258 careful = true;
259 continue;
260 }
261 if self.records.len() == VECTOR_SIZE {
262 break;
263 }
264 if careful && self.records.len() == limit {
265 continue;
266 }
267 if self.drained {
268 break;
269 }
270 if self.buffer.len() - start + self.block > Span::MOST {
275 if self.records.is_empty() {
276 return Err(Error::io("a record is longer than two gigabytes"));
277 }
278 break;
279 }
280 self.fill(start)?;
281 self.records.shift(start);
282 start = 0;
283 }
284 Ok(self.records.len())
285 }
286
287 pub(crate) fn here(&self) -> u64 {
289 self.offset - self.buffer.len() as u64 + self.at as u64
290 }
291
292 #[must_use]
294 pub fn bytes_read(&self) -> u64 {
295 self.offset - self.origin
296 }
297
298 pub(crate) fn stretch(&self, from: u64, end: u64, line: u64) -> Self {
305 Self {
306 file: Arc::clone(&self.file),
307 path: self.path.clone(),
308 given: self.given,
309 dialect: self.dialect,
310 fields: self.fields.clone(),
311 projection: self.projection.clone(),
312 buffer: Vec::new(),
313 at: 0,
314 offset: from,
315 drained: false,
316 line,
317 scratch: Vec::new(),
318 records: Records::default(),
319 block: self.block,
320 origin: from,
321 end,
322 cap: u64::MAX,
323 }
324 }
325
326 pub(crate) fn give_up_at(&mut self, cap: u64) {
328 self.cap = cap;
329 }
330
331 pub(crate) fn skim(&mut self) -> Result<u64> {
334 while self.next_records()? > 0 {}
335 Ok(self.here())
336 }
337
338 pub(crate) fn skip_to(&mut self, target: u64) -> Result<()> {
344 loop {
345 let from = self.here();
346 let rows = self.next_records()?;
347 if rows == 0 {
348 return Ok(());
349 }
350 if self.here() > target {
351 let front = self.offset - self.buffer.len() as u64;
352 self.at = usize::try_from(from - front)
353 .map_err(|_| Error::internal("a chunk start outside the buffer"))?;
354 return Ok(());
355 }
356 self.line += rows as u64;
357 }
358 }
359
360 pub(crate) const fn line(&self) -> u64 {
362 self.line
363 }
364
365 pub(crate) fn file(&self) -> &dyn File {
367 self.file.as_ref()
368 }
369
370 fn conversion_error(&self, text: &str, field: &Field, line: u64) -> String {
382 format!(
383 "CSV Error on Line: {line}\nOriginal Line: {text}\nError when converting column \
384 \"{}\". Could not convert string \"{text}\" to '{}'\n\nColumn {} is being converted \
385 as type {}\nThis type was auto-detected from the CSV file.\nPossible solutions:\n* \
386 Override the type for this column manually by setting the type explicitly, e.g., \
387 types={{'{}': 'VARCHAR'}}\n* Set the sample size to a larger value to enable the \
388 auto-detection to scan more values, e.g., sample_size=-1\n* Use a COPY statement to \
389 automatically derive types from an existing table.\n* Check whether the null string \
390 value is set correctly (e.g., nullstr = 'N/A')\n\n file = {}\n delimiter = {}\n \
391 quote = {}\n escape = {}\n header = {} {}\n sample_size = {}\n",
392 field.name,
393 field.ty,
394 field.name,
395 field.ty,
396 field.name,
397 self.path,
398 Given::shown(self.given.delimiter, Some(self.dialect.delimiter)),
399 Given::shown(self.given.quote, self.dialect.quote),
400 Given::shown(self.given.escape, self.dialect.escape),
401 self.dialect.header,
402 Given::source(self.given.header.is_some()),
403 infer::SAMPLE,
404 )
405 }
406
407 #[cfg(test)]
410 fn next_chunk_by_record(&mut self) -> Result<Option<Chunk>> {
411 let mut rows: Vec<Vec<Option<String>>> = Vec::new();
412 while rows.len() < VECTOR_SIZE {
413 match self.next_record()? {
414 Some(fields) => rows.push(fields),
415 None => break,
416 }
417 }
418 if rows.is_empty() {
419 return Ok(None);
420 }
421 let mut columns = Vec::with_capacity(self.projection.len());
422 for &at in &self.projection {
423 let field = &self.fields[at];
424 let mut values = Vec::with_capacity(rows.len());
425 for (row, held) in rows.iter().enumerate() {
426 let text = held.get(at).and_then(Option::as_deref);
427 values.push(self.convert(
428 text,
429 field,
430 self.line - rows.len() as u64 + row as u64,
431 )?);
432 }
433 columns.push(rudb_vector::Vector::from_values(field.ty.clone(), &values)?);
434 }
435 Ok(Some(Chunk::with_rows(columns, rows.len())?))
436 }
437
438 #[cfg(test)]
440 fn convert(&self, text: Option<&str>, field: &Field, line: u64) -> Result<rudb_common::Value> {
441 let Some(text) = text else { return Ok(rudb_common::Value::Null) };
442 if field.ty == LogicalType::Varchar {
443 return Ok(rudb_common::Value::Varchar(text.to_string()));
444 }
445 let value = rudb_common::Value::Varchar(text.to_string());
446 match rudb_kernels::cast_value(&value, &field.ty, false) {
447 Ok(converted) => Ok(converted),
448 Err(_) => Err(Error::conversion(self.conversion_error(text, field, line))),
449 }
450 }
451
452 #[cfg(test)]
457 fn next_record(&mut self) -> Result<Option<Vec<Option<String>>>> {
458 let Some(()) = self.advance()? else { return Ok(None) };
459 Ok(Some(
460 self.scratch
461 .iter()
462 .map(|text| if text.is_empty() { None } else { Some(text.clone()) })
463 .collect(),
464 ))
465 }
466
467 fn advance(&mut self) -> Result<Option<()>> {
469 loop {
470 let mut scratch = std::mem::take(&mut self.scratch);
471 let outcome = crate::scan::record(
472 &self.buffer,
473 self.at,
474 self.dialect,
475 self.drained,
476 &mut scratch,
477 );
478 self.scratch = scratch;
479 match outcome? {
480 Some(next) => {
481 self.at = next;
482 self.line += 1;
483 return Ok(Some(()));
484 }
485 None if self.drained => return Ok(None),
486 None => self.fill(self.at)?,
487 }
488 }
489 }
490
491 fn skip_record(&mut self) -> Result<()> {
493 self.advance()?;
494 Ok(())
495 }
496
497 fn fill(&mut self, keep: usize) -> Result<()> {
502 if self.offset >= self.cap {
503 return Err(Error::io("a record runs further than a guessed start is followed"));
504 }
505 self.buffer.drain(..keep);
506 self.at -= keep;
507 let held = self.buffer.len();
508 let want = match self.end.checked_sub(self.offset) {
509 Some(left) if left > 0 => self.block.min(usize::try_from(left).unwrap_or(usize::MAX)),
510 _ => self.block.min(TAIL.max(held)),
511 };
512 self.buffer.resize(held + want, 0);
513 let read = self.file.read_at(self.offset, &mut self.buffer[held..])?;
514 self.buffer.truncate(held + read);
515 self.offset += read as u64;
516 if read == 0 {
517 self.drained = true;
518 }
519 Ok(())
520 }
521
522 fn sample_rows(&self, sample: &[u8]) -> Result<Vec<Vec<Option<String>>>> {
525 let mut rows = Vec::new();
526 let mut fields = Vec::new();
527 let mut at = 0;
528 while rows.len() <= infer::SAMPLE {
529 let Some(next) = crate::scan::record(sample, at, self.dialect, false, &mut fields)?
532 else {
533 break;
534 };
535 at = next;
536 rows.push(
537 fields
538 .iter()
539 .map(|text| if text.is_empty() { None } else { Some(text.clone()) })
540 .collect(),
541 );
542 }
543 Ok(rows)
544 }
545}
546
547fn describe(rows: &[Vec<Option<String>>], told: Option<bool>) -> (bool, Vec<Field>) {
558 let width = rows.iter().map(Vec::len).max().unwrap_or(0);
559 let body = types(&rows[1.min(rows.len())..], width);
560 let all_text = body.iter().all(|ty| *ty == LogicalType::Varchar);
561 let first_fits = rows.first().is_some_and(|first| {
562 first.iter().zip(&body).all(|(text, ty)| match text {
563 None => true,
564 Some(text) => infer::fits(text, ty),
565 })
566 });
567 let header = !rows.is_empty() && told.unwrap_or(rows.len() > 1 && (all_text || !first_fits));
569 if !header {
570 let types = types(rows, width);
571 let fields = types
572 .into_iter()
573 .enumerate()
574 .map(|(at, ty)| Field::new(format!("column{at}"), ty))
575 .collect();
576 return (false, fields);
577 }
578 let names = unique(&rows[0], width);
579 let fields = body.into_iter().zip(names).map(|(ty, name)| Field::new(name, ty)).collect();
580 (true, fields)
581}
582
583fn unique(header: &[Option<String>], width: usize) -> Vec<String> {
596 let mut taken: Vec<String> = Vec::with_capacity(width);
597 for at in 0..width {
598 let base = match header.get(at).and_then(Option::as_deref) {
599 Some(written) => written.to_string(),
600 None => format!("column{at}"),
601 };
602 let mut name = base.clone();
603 let mut next = 1;
604 while taken.iter().any(|held| held.eq_ignore_ascii_case(&name)) {
605 name = format!("{base}_{next}");
606 next += 1;
607 }
608 taken.push(name);
609 }
610 taken
611}
612
613fn types(rows: &[Vec<Option<String>>], width: usize) -> Vec<LogicalType> {
615 (0..width)
616 .map(|at| {
617 let values: Vec<Option<&str>> =
618 rows.iter().map(|row| row.get(at).and_then(Option::as_deref)).collect();
619 infer::column(&values)
620 })
621 .collect()
622}
623
624#[cfg(test)]
625mod tests {
626 use super::*;
627 use rudb_common::Value;
628 use rudb_io::{Filesystem, OpenMode, SimFilesystem};
629 use std::path::Path;
630
631 fn read(text: &str) -> Reader {
632 let filesystem = SimFilesystem::new();
633 let path = Path::new("/t.csv");
634 let file = filesystem.open(path, OpenMode::Create).expect("creates");
635 file.write_at(0, text.as_bytes()).expect("writes");
636 drop(file);
637 let file = filesystem.open(path, OpenMode::Read).expect("opens");
638 Reader::open(file, "/t.csv").expect("sniffs")
639 }
640
641 fn read_with(text: &str, given: Given) -> Reader {
643 let filesystem = SimFilesystem::new();
644 let path = Path::new("/t.csv");
645 let file = filesystem.open(path, OpenMode::Create).expect("creates");
646 file.write_at(0, text.as_bytes()).expect("writes");
647 drop(file);
648 let file = filesystem.open(path, OpenMode::Read).expect("opens");
649 Reader::open_with(file, "/t.csv", given).expect("reads")
650 }
651
652 fn names_and_types(reader: &Reader) -> Vec<(String, String)> {
653 reader.fields().iter().map(|f| (f.name.clone(), f.ty.to_string())).collect()
654 }
655
656 fn all(reader: &mut Reader) -> Vec<Vec<Value>> {
657 let mut rows = Vec::new();
658 while let Some(chunk) = reader.next_chunk().expect("reads") {
659 for row in 0..chunk.len() {
660 rows.push((0..chunk.width()).map(|at| chunk.value_at(row, at)).collect());
661 }
662 }
663 rows
664 }
665
666 #[test]
667 fn a_header_that_names_two_columns_the_same_thing_counts_the_second_one_up() {
668 let names: Vec<String> =
669 read("a,a,A,a_1\n1,2,3,4\nx,y,z,w\n").fields().into_iter().map(|f| f.name).collect();
670 assert_eq!(names, ["a", "a_1", "A_2", "a_1_1"]);
674 }
675
676 #[test]
677 fn a_header_and_three_types_are_what_duckdb_sniffs_for_the_same_bytes() {
678 let reader = read("a,b,c\n1,x,2.5\n2,y,3.5\n");
679 assert_eq!(
680 names_and_types(&reader),
681 [
682 ("a".to_string(), "BIGINT".to_string()),
683 ("b".to_string(), "VARCHAR".to_string()),
684 ("c".to_string(), "DOUBLE".to_string()),
685 ]
686 );
687 }
688
689 #[test]
690 fn a_file_with_no_header_gets_the_names_duckdb_gives_it() {
691 let reader = read("1,x\n2,y\n");
692 assert_eq!(
693 names_and_types(&reader),
694 [
695 ("column0".to_string(), "BIGINT".to_string()),
696 ("column1".to_string(), "VARCHAR".to_string()),
697 ]
698 );
699 }
700
701 #[test]
702 fn two_rows_of_words_are_a_header_and_a_row() {
703 let reader = read("a,b\nc,d\n");
704 assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["a", "b"]);
705 }
706
707 #[test]
708 fn one_column_of_words_under_a_row_of_numbers_is_still_a_header() {
709 let reader = read("a,2\n3,4\n");
712 assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["a", "2"]);
713 }
714
715 #[test]
716 fn the_rows_are_the_rows_of_the_file() {
717 let mut reader = read("a,b\n1,x\n2,y\n");
718 assert_eq!(
719 all(&mut reader),
720 [
721 vec![Value::BigInt(1), Value::Varchar("x".into())],
722 vec![Value::BigInt(2), Value::Varchar("y".into())],
723 ]
724 );
725 }
726
727 #[test]
728 fn an_empty_field_is_a_null_whether_it_was_quoted_or_not() {
729 let mut reader = read("a,b\n1,\n\"\",y\n");
732 assert_eq!(
733 all(&mut reader),
734 [vec![Value::BigInt(1), Value::Null], vec![Value::Null, Value::Varchar("y".into())],]
735 );
736 }
737
738 #[test]
739 fn a_projection_picks_columns_out_by_position_and_can_reorder_them() {
740 let mut reader = read("a,b,c\n1,x,2.5\n");
741 reader.project(&[2, 0]).expect("projects");
742 assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["c", "a"]);
743 assert_eq!(all(&mut reader), [vec![Value::Double(2.5), Value::BigInt(1)]]);
744 }
745
746 #[test]
747 fn a_projection_of_nothing_still_counts_the_rows() {
748 let mut reader = read("a,b\n1,x\n2,y\n3,z\n");
749 reader.project(&[]).expect("projects");
750 let chunk = reader.next_chunk().expect("reads").expect("a chunk");
751 assert_eq!(chunk.len(), 3);
752 assert_eq!(chunk.width(), 0);
753 }
754
755 #[test]
756 fn a_pipe_separated_file_reads_as_one() {
757 let mut reader = read("a|b\n1|x\n");
758 assert_eq!(reader.dialect().delimiter, b'|');
759 assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x".into())]]);
760 }
761
762 #[test]
763 fn a_quoted_field_with_a_delimiter_in_it_is_one_value() {
764 let mut reader = read("a,b\n1,\"x,y\"\n");
765 assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x,y".into())]]);
766 }
767
768 #[test]
769 fn more_rows_than_fit_one_chunk_arrive_as_more_than_one_chunk() {
770 let mut text = String::from("a\n");
771 for row in 0..VECTOR_SIZE + 5 {
772 text.push_str(&format!("{row}\n"));
773 }
774 let mut reader = read(&text);
775 let first = reader.next_chunk().expect("reads").expect("a chunk");
776 assert_eq!(first.len(), VECTOR_SIZE);
777 let second = reader.next_chunk().expect("reads").expect("a second chunk");
778 assert_eq!(second.len(), 5);
779 assert!(reader.next_chunk().expect("reads").is_none());
780 }
781
782 #[test]
783 fn a_value_the_sniffer_never_saw_is_an_error_rather_than_a_wider_column() {
784 let mut text = String::from("c\n");
788 for row in 0..infer::SAMPLE {
789 text.push_str(&format!("{row}\n"));
790 }
791 text.push_str("oops\n");
792 let mut reader = read(&text);
793 assert_eq!(reader.fields()[0].ty, LogicalType::BigInt);
794 let error = all_or_error(&mut reader).unwrap_err();
795 let line = infer::SAMPLE + 2;
796 assert!(error.message().starts_with(&format!("CSV Error on Line: {line}")), "{error}");
797 assert!(
798 error.message().contains("Could not convert string \"oops\" to 'BIGINT'"),
799 "{error}"
800 );
801 assert!(error.message().contains("sample_size = 20480"), "{error}");
802 }
803
804 #[test]
807 fn a_file_told_it_has_no_header_reads_its_first_line_as_a_row() {
808 let mut reader =
809 read_with("a,b\n1,x\n2,y\n", Given { header: Some(false), ..Given::default() });
810 assert_eq!(
811 names_and_types(&reader),
812 [
813 ("column0".to_string(), "VARCHAR".to_string()),
814 ("column1".to_string(), "VARCHAR".to_string()),
815 ]
816 );
817 assert_eq!(all(&mut reader).len(), 3);
818 }
819
820 #[test]
822 fn a_file_told_it_has_a_header_takes_its_first_line_as_the_names() {
823 let reader = read_with("1,2\n3,4\n", Given { header: Some(true), ..Given::default() });
824 assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["1", "2"]);
825 }
826
827 #[test]
829 fn a_given_delimiter_is_the_delimiter_whatever_the_file_looks_like() {
830 let reader = read_with("a,b\n1,x\n", Given { delimiter: Some(b';'), ..Given::default() });
831 assert_eq!(reader.dialect().delimiter, b';');
832 assert_eq!(reader.fields().len(), 1);
833 }
834
835 #[test]
837 fn a_given_quote_makes_a_field_that_holds_the_delimiter_one_value() {
838 let mut reader =
839 read_with("a,b\n1,'x,y'\n", Given { quote: Some(b'\''), ..Given::default() });
840 assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x,y".into())]]);
841 }
842
843 #[test]
845 fn the_block_says_set_by_user_for_what_the_call_gave_it() {
846 let mut text = String::from("c;d\n");
847 for row in 0..infer::SAMPLE {
848 text.push_str(&format!("{row};x\n"));
849 }
850 text.push_str("oops;x\n");
851 let given = Given { delimiter: Some(b';'), ..Given::default() };
852 let mut reader = read_with(&text, given);
853 let error = all_or_error(&mut reader).unwrap_err();
854 assert!(error.message().contains("delimiter = ; (Set By User)"), "{error}");
855 assert!(error.message().contains("header = true (Auto-Detected)"), "{error}");
856 }
857
858 #[test]
859 fn a_file_told_a_wider_type_than_it_sniffed_reads_its_whole_numbers_as_that_type() {
860 let mut reader = read("a\n1\n2\n");
864 assert_eq!(reader.fields()[0].ty, LogicalType::BigInt);
865 reader.retype(&[LogicalType::Double]).expect("one type for one column");
866 assert_eq!(reader.fields()[0].ty, LogicalType::Double);
867 assert_eq!(all(&mut reader), [[Value::Double(1.0)], [Value::Double(2.0)]]);
868 }
869
870 #[test]
871 fn a_type_list_that_is_not_as_long_as_the_projection_is_refused() {
872 let mut reader = read("a,b\n1,two\n");
873 let error = reader.retype(&[LogicalType::Double]).unwrap_err();
874 assert!(error.message().contains("1 types for a projection of 2 columns"), "{error}");
875 }
876
877 fn drained(reader: &mut Reader, old: bool) -> (Vec<(usize, Vec<String>)>, Option<String>) {
880 let mut chunks = Vec::new();
881 loop {
882 let next = if old { reader.next_chunk_by_record() } else { reader.next_chunk() };
883 match next {
884 Ok(Some(chunk)) => {
885 let mut values = Vec::new();
886 for row in 0..chunk.len() {
887 for at in 0..chunk.width() {
888 values.push(format!("{:?}", chunk.value_at(row, at)));
889 }
890 }
891 chunks.push((chunk.len(), values));
892 }
893 Ok(None) => return (chunks, None),
894 Err(error) => return (chunks, Some(error.to_string())),
895 }
896 }
897 }
898
899 struct Rng(u64);
901
902 impl Rng {
903 fn next(&mut self) -> u64 {
904 self.0 ^= self.0 << 13;
905 self.0 ^= self.0 >> 7;
906 self.0 ^= self.0 << 17;
907 self.0
908 }
909
910 fn below(&mut self, n: usize) -> usize {
911 (self.next() % n as u64) as usize
912 }
913 }
914
915 fn typed_file(rng: &mut Rng, rows: usize) -> String {
918 const ODD: [&str; 27] = [
919 "",
920 " 1",
921 "1 ",
922 "1e3",
923 "0x10",
924 "inf",
925 "-nan",
926 "abc",
927 "\"12\"",
928 "\"a\"\"b\"",
929 "\"x,y\"",
930 "\"x\ny\"",
931 "h\u{e9}llo",
932 "99999999999999999999",
933 "9999999999999999999",
934 "-",
935 "+5",
936 "007",
937 "2020-02-30",
938 "2020-02-29",
939 "0000-01-01",
940 "TRUE",
941 "no",
942 "1_000",
943 "1.5e-3",
944 "-0",
945 "\"\"",
946 ];
947 let width = 1 + rng.below(6);
948 let kinds: Vec<usize> = (0..width).map(|_| rng.below(5)).collect();
949 let mut text: String = (0..width).map(|at| format!("c{at}")).collect::<Vec<_>>().join(",");
950 text.push('\n');
951 for _ in 0..rows {
952 let mut fields = Vec::with_capacity(width);
953 for &kind in &kinds {
954 let odd = rng.below(60) == 0;
955 fields.push(if odd {
956 ODD[rng.below(ODD.len())].to_string()
957 } else {
958 let n = rng.next();
959 match kind {
960 0 => format!("{}", (n % 2_000_001) as i64 - 1_000_000),
961 1 => format!("{}.{:02}", n % 100_000, n % 100),
962 2 => format!("{}-{:02}-{:02}", 1990 + n % 20, 1 + n % 12, 1 + n % 28),
963 3 => ["true", "false", "t", "F"][(n % 4) as usize].to_string(),
964 _ => ["x", "hello world", "a longer piece of text", "\"q,\"\"q\""]
965 [(n % 4) as usize]
966 .to_string(),
967 }
968 });
969 }
970 if rng.below(200) == 0 {
971 fields.pop();
972 }
973 if rng.below(200) == 0 {
974 fields.push("extra".to_string());
975 }
976 text.push_str(&fields.join(","));
977 text.push_str(["\n", "\n", "\n", "\r\n", "\r"][rng.below(5)]);
978 }
979 if rng.below(4) == 0 {
980 text.pop();
981 }
982 text
983 }
984
985 fn open_sized(text: &[u8], block: usize) -> Result<Reader> {
986 let filesystem = SimFilesystem::new();
987 let path = Path::new("/t.csv");
988 let file = filesystem.open(path, OpenMode::Create).expect("creates");
989 file.write_at(0, text).expect("writes");
990 drop(file);
991 let file = filesystem.open(path, OpenMode::Read).expect("opens");
992 Reader::open_sized(file, "/t.csv", Given::default(), block)
993 }
994
995 #[test]
999 fn generated_files_read_the_same_a_chunk_at_a_time_as_a_record_at_a_time() {
1000 let types = [
1001 LogicalType::BigInt,
1002 LogicalType::Integer,
1003 LogicalType::SmallInt,
1004 LogicalType::TinyInt,
1005 LogicalType::UBigInt,
1006 LogicalType::UInteger,
1007 LogicalType::USmallInt,
1008 LogicalType::UTinyInt,
1009 LogicalType::Double,
1010 LogicalType::Float,
1011 LogicalType::Date,
1012 LogicalType::Boolean,
1013 LogicalType::Varchar,
1014 LogicalType::Timestamp,
1015 LogicalType::Decimal { width: 18, scale: 3 },
1016 ];
1017 let mut rng = Rng(0x2545_f491_4f6c_dd1d);
1018 for case in 0..200 {
1019 let rows = if case % 100 == 0 { 8192 + rng.below(1000) } else { rng.below(200) };
1020 let mut text = typed_file(&mut rng, rows).into_bytes();
1021 if rng.below(20) == 0 {
1022 let at = rng.below(text.len() + 1);
1024 text.splice(at..at, *b",\"x\"y,");
1025 }
1026 for block in [1 << 20, 32 + rng.below(400)] {
1027 let (Ok(mut new), Ok(mut old)) =
1028 (open_sized(&text, block), open_sized(&text, block))
1029 else {
1030 continue;
1031 };
1032 assert_eq!(new.fields(), old.fields());
1033 let width = new.fields().len();
1034 if width > 0 && rng.below(2) == 0 {
1035 let columns: Vec<usize> =
1036 (0..rng.below(width + 2)).map(|_| rng.below(width)).collect();
1037 new.project(&columns).expect("projects");
1038 old.project(&columns).expect("projects");
1039 let wanted: Vec<LogicalType> =
1040 columns.iter().map(|_| types[rng.below(types.len())].clone()).collect();
1041 new.retype(&wanted).expect("retypes");
1042 old.retype(&wanted).expect("retypes");
1043 }
1044 let expected = drained(&mut old, true);
1045 let found = drained(&mut new, false);
1046 assert_eq!(found.1, expected.1, "case {case}, block {block}");
1047 assert_eq!(found.0, expected.0, "case {case}, block {block}");
1048 }
1049 }
1050 }
1051
1052 #[test]
1057 #[ignore = "a measurement, run by hand"]
1058 fn reads_lineitem_faster_a_chunk_at_a_time() {
1059 use rudb_io::RealFilesystem;
1060 use std::time::Instant;
1061
1062 const ROWS: usize = 200_000;
1063 let mut rng = Rng(0x1234_5678_9abc_def1);
1064 let mut text = String::from(
1065 "l_orderkey,l_partkey,l_suppkey,l_linenumber,l_quantity,l_extendedprice,l_discount,\
1066 l_tax,l_returnflag,l_linestatus,l_shipdate,l_commitdate,l_receiptdate,\
1067 l_shipinstruct,l_shipmode,l_comment\n",
1068 );
1069 let words = ["carefully", "final", "deposits", "furiously", "regular", "ideas", "sleep"];
1070 for row in 0..ROWS {
1071 let n = rng.next();
1072 let date = |shift: u64| {
1073 format!(
1074 "{}-{:02}-{:02}",
1075 1992 + (n >> shift) % 7,
1076 1 + (n >> shift) % 12,
1077 1 + (n >> shift) % 28
1078 )
1079 };
1080 let comment: Vec<&str> =
1081 (0..3 + n % 4).map(|k| words[((n >> (k * 3)) % 7) as usize]).collect();
1082 text.push_str(&format!(
1083 "{},{},{},{},{}.00,{}.{:02},0.0{},0.0{},{},{},{},{},{},{},{},{}\n",
1084 row / 4 + 1,
1085 n % 200_000,
1086 n % 10_000,
1087 row % 4 + 1,
1088 1 + n % 50,
1089 900 + n % 100_000,
1090 n % 100,
1091 n % 10,
1092 (n >> 7) % 9,
1093 ["A", "N", "R"][(n % 3) as usize],
1094 ["O", "F"][(n % 2) as usize],
1095 date(3),
1096 date(11),
1097 date(19),
1098 ["DELIVER IN PERSON", "NONE", "TAKE BACK RETURN"][(n % 3) as usize],
1099 ["TRUCK", "MAIL", "AIR", "SHIP"][(n % 4) as usize],
1100 comment.join(" "),
1101 ));
1102 }
1103 let path = std::env::temp_dir().join(format!("rudb-lineitem-{}.csv", std::process::id()));
1104 std::fs::write(&path, &text).expect("writes");
1105 let megabytes = text.len() as f64 / 1e6;
1106 let filesystem = RealFilesystem::new();
1107 let mut best = [f64::MAX; 2];
1108 for _ in 0..5 {
1109 for (slot, old) in [(0, true), (1, false)] {
1110 let file = filesystem.open(&path, OpenMode::Read).expect("opens");
1111 let mut reader = Reader::open(file, "lineitem.csv").expect("sniffs");
1112 let started = Instant::now();
1113 let mut rows = 0;
1114 loop {
1115 let next =
1116 if old { reader.next_chunk_by_record() } else { reader.next_chunk() };
1117 let Some(chunk) = next.expect("reads") else { break };
1118 rows += chunk.len();
1119 }
1120 assert_eq!(rows, ROWS);
1121 best[slot] = best[slot].min(started.elapsed().as_secs_f64());
1122 }
1123 }
1124 std::fs::remove_file(&path).expect("removes");
1125 let (old, new) = (megabytes / best[0], megabytes / best[1]);
1126 println!(
1127 "{ROWS} rows, {megabytes:.1} MB: a record at a time {old:.1} MB/s, a chunk at a time"
1128 );
1129 println!("{new:.1} MB/s, {:.2}x", best[0] / best[1]);
1130 }
1131
1132 fn all_or_error(reader: &mut Reader) -> Result<Vec<Vec<Value>>> {
1133 let mut rows = Vec::new();
1134 while let Some(chunk) = reader.next_chunk()? {
1135 for row in 0..chunk.len() {
1136 rows.push((0..chunk.width()).map(|at| chunk.value_at(row, at)).collect());
1137 }
1138 }
1139 Ok(rows)
1140 }
1141}