1use std::borrow::Cow;
18use std::collections::{HashMap, HashSet};
19use std::sync::Arc;
20use std::sync::atomic::{AtomicUsize, Ordering};
21
22use polars::chunked_array::cast::CastOptions;
23use polars::prelude::{
24 DataType, Field, LazyFrame, PlRefPath, PlSmallStr, PolarsResult, Schema, TimeUnit, UnionArgs,
25 concat,
26};
27
28pub struct Pass<'a>(&'a FooterProgress);
30
31impl Pass<'_> {
32 pub fn advance(&self) {
34 self.0.advance();
35 }
36}
37
38impl Drop for Pass<'_> {
39 fn drop(&mut self) {
40 self.0.done();
41 }
42}
43
44pub struct Listing<'a>(&'a FooterProgress);
47
48impl Listing<'_> {
49 pub fn advance(&self) {
51 self.0.listed.fetch_add(1, Ordering::Relaxed);
52 }
53
54 pub fn counter(&self) -> std::sync::Arc<AtomicUsize> {
56 self.0.listed.clone()
57 }
58
59 pub fn add(&self, n: usize) {
61 self.0.listed.fetch_add(n, Ordering::Relaxed);
62 }
63
64 pub fn is_cancelled(&self) -> bool {
66 self.0.is_cancelled()
67 }
68
69 pub fn cancel_flag(&self) -> std::sync::Arc<std::sync::atomic::AtomicBool> {
71 self.0.cancel_flag()
72 }
73}
74
75impl Drop for Listing<'_> {
76 fn drop(&mut self) {
77 self.0.listing.store(false, Ordering::Release);
78 }
79}
80
81#[doc(hidden)]
86#[derive(Debug, Clone, Copy, PartialEq, Eq)]
87pub struct PassCount {
88 pub begun: usize,
90 pub read: usize,
92 pub total: usize,
94}
95
96#[derive(Debug, Default)]
106pub struct FooterProgress {
107 read: AtomicUsize,
108 total: AtomicUsize,
109 passes: AtomicUsize,
113 last_total: AtomicUsize,
115 cancelled: std::sync::Arc<std::sync::atomic::AtomicBool>,
118 listed: std::sync::Arc<AtomicUsize>,
121 listing: std::sync::atomic::AtomicBool,
122 at_once: AtomicUsize,
124 estimate: std::sync::Mutex<Option<RowEstimate>>,
126}
127
128#[derive(Debug, Clone, Copy, PartialEq, Eq)]
131pub struct RowEstimate {
132 pub rows: u64,
133 pub sampled: usize,
135 pub files: usize,
137}
138
139impl RowEstimate {
140 pub fn of<'a>(
143 files: usize,
144 footers: impl IntoIterator<Item = &'a Option<FileFooter>>,
145 ) -> Option<Self> {
146 let (sampled, rows) = footers
147 .into_iter()
148 .flatten()
149 .fold((0usize, 0u128), |(n, rows), f| {
150 (n + 1, rows + f.rows() as u128)
151 });
152 (sampled > 0).then(|| RowEstimate {
153 rows: u64::try_from(rows * files as u128 / sampled as u128).unwrap_or(u64::MAX),
154 sampled,
155 files,
156 })
157 }
158}
159
160pub const ESTIMATE_SAMPLE: usize = 2_000;
163
164pub const COUNT_AT_ONCE: usize = 256;
167
168pub fn random_sample(files: usize, n: usize, seed: u64) -> Vec<usize> {
171 if files <= n {
172 return (0..files).collect();
173 }
174 let mut state = seed;
176 let mut next = move || {
177 state = state.wrapping_add(0x9e37_79b9_7f4a_7c15);
178 let mut z = state;
179 z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
180 z = (z ^ (z >> 27)).wrapping_mul(0x94d0_49bb_1331_11eb);
181 z ^ (z >> 31)
182 };
183 let mut chosen = std::collections::BTreeSet::new();
184 while chosen.len() < n {
185 chosen.insert((next() % files as u64) as usize);
186 }
187 chosen.into_iter().collect()
188}
189
190impl FooterProgress {
191 pub fn listing(&self) -> Listing<'_> {
193 self.listed.store(0, Ordering::Relaxed);
194 self.listing.store(true, Ordering::Release);
195 Listing(self)
196 }
197
198 pub fn listed(&self) -> Option<usize> {
200 self.listing
201 .load(Ordering::Acquire)
202 .then(|| self.listed.load(Ordering::Relaxed))
203 }
204
205 pub fn begin(&self, total: usize) {
207 self.read.store(0, Ordering::Relaxed);
208 self.last_total.store(total, Ordering::Relaxed);
212 self.total.store(total, Ordering::Release);
213 self.passes.fetch_add(1, Ordering::Relaxed);
214 }
215
216 pub fn advance(&self) {
218 self.read.fetch_add(1, Ordering::Relaxed);
219 }
220
221 pub fn done(&self) {
223 self.total.store(0, Ordering::Relaxed);
224 }
225
226 pub fn pass(&self, total: usize) -> Pass<'_> {
232 self.begin(total);
233 Pass(self)
234 }
235
236 #[doc(hidden)]
245 pub fn last_pass(&self) -> PassCount {
246 PassCount {
247 begun: self.passes.load(Ordering::Relaxed),
248 read: self.read.load(Ordering::Relaxed),
249 total: self.last_total.load(Ordering::Relaxed),
250 }
251 }
252
253 pub fn reading(&self) -> Option<(usize, usize)> {
255 let total = self.total.load(Ordering::Acquire);
256 (total > 0).then(|| (self.read.load(Ordering::Relaxed).min(total), total))
257 }
258
259 pub fn cancel(&self) {
262 self.cancelled.store(true, Ordering::Relaxed);
263 }
264
265 pub fn is_cancelled(&self) -> bool {
266 self.cancelled.load(Ordering::Relaxed)
267 }
268
269 pub fn cancel_flag(&self) -> std::sync::Arc<std::sync::atomic::AtomicBool> {
271 self.cancelled.clone()
272 }
273
274 pub fn counting() -> Self {
276 let progress = Self::default();
277 progress.at_once.store(COUNT_AT_ONCE, Ordering::Relaxed);
278 progress
279 }
280
281 pub fn reads_at_once(&self) -> usize {
283 match self.at_once.load(Ordering::Relaxed) {
284 0 => FOOTERS_AT_ONCE,
285 n => n,
286 }
287 }
288
289 pub fn set_estimate(&self, estimate: Option<RowEstimate>) {
291 *self.estimate.lock().unwrap_or_else(|e| e.into_inner()) = estimate;
292 }
293
294 pub fn estimate(&self) -> Option<RowEstimate> {
295 *self.estimate.lock().unwrap_or_else(|e| e.into_inner())
296 }
297}
298
299#[derive(Debug, Clone)]
302pub struct FileFooter {
303 pub schema: Arc<Schema>,
304 pub row_group_rows: Vec<usize>,
306 pub row_group_bytes: Vec<usize>,
313 pub file_bytes: usize,
315 pub column_bytes: Vec<(String, usize)>,
319}
320
321impl FileFooter {
322 pub fn from_metadata(
325 schema: Schema,
326 metadata: &polars_parquet::parquet::metadata::FileMetadata,
327 file_bytes: usize,
328 widths: bool,
329 ) -> Self {
330 let column_bytes = if widths {
331 parquet_column_bytes(&schema, metadata)
332 } else {
333 Vec::new()
334 };
335 FileFooter {
336 schema: Arc::new(schema),
337 row_group_rows: metadata.row_groups.iter().map(|rg| rg.num_rows()).collect(),
338 row_group_bytes: metadata
339 .row_groups
340 .iter()
341 .map(|rg| rg.compressed_size())
342 .collect(),
343 file_bytes,
344 column_bytes,
345 }
346 }
347
348 pub fn from_tail(tail: &[u8], file_bytes: usize, widths: bool) -> color_eyre::Result<Self> {
351 use polars::prelude::{ParquetReader, SchemaExt, SerReader};
352 let mut cursor = std::io::Cursor::new(tail);
353 let mut reader = ParquetReader::new(&mut cursor);
354 let arrow_schema = reader
355 .schema()
356 .map_err(|e| color_eyre::eyre::eyre!("Parquet schema read failed: {e}"))?;
357 let metadata = reader
358 .get_metadata()
359 .map_err(|e| color_eyre::eyre::eyre!("Parquet footer read failed: {e}"))?;
360 Ok(Self::from_metadata(
361 Schema::from_arrow_schema(arrow_schema.as_ref()),
362 metadata,
363 file_bytes,
364 widths,
365 ))
366 }
367
368 pub fn rows(&self) -> usize {
370 self.row_group_rows.iter().sum()
371 }
372}
373
374pub fn parquet_column_bytes(
378 schema: &Schema,
379 metadata: &polars_parquet::parquet::metadata::FileMetadata,
380) -> Vec<(String, usize)> {
381 schema
382 .iter_names()
383 .map(|name| {
384 let bytes: i64 = metadata
385 .row_groups
386 .iter()
387 .flat_map(|rg| rg.columns_under_root_iter(name).into_iter().flatten())
388 .map(|chunk| chunk.uncompressed_size())
389 .sum();
390 (name.to_string(), bytes.max(0) as usize)
391 })
392 .collect()
393}
394
395pub fn column_bytes_per_row(footers: &[Option<FileFooter>]) -> Vec<(String, usize)> {
398 let rows: usize = footers.iter().flatten().map(FileFooter::rows).sum();
399 if rows == 0 {
400 return Vec::new();
401 }
402 let mut totals: Vec<(String, usize)> = Vec::new();
403 let mut at: HashMap<String, usize> = HashMap::new();
404 for (name, bytes) in footers.iter().flatten().flat_map(|f| &f.column_bytes) {
405 match at.get(name) {
406 Some(&i) => totals[i].1 += bytes,
407 None => {
408 at.insert(name.clone(), totals.len());
409 totals.push((name.clone(), *bytes));
410 }
411 }
412 }
413 totals
414 .into_iter()
415 .map(|(name, bytes)| (name, bytes / rows))
416 .collect()
417}
418
419#[derive(Debug, Clone, PartialEq, Eq)]
422pub enum SchemaOrigin {
423 AllFooters(usize),
425 FooterSample { read: usize, total: usize },
427}
428
429impl SchemaOrigin {
430 pub fn total_files(&self) -> usize {
436 match self {
437 SchemaOrigin::AllFooters(files) => *files,
438 SchemaOrigin::FooterSample { total, .. } => *total,
439 }
440 }
441}
442
443impl std::fmt::Display for SchemaOrigin {
444 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
445 match self {
446 SchemaOrigin::AllFooters(1) => write!(f, "one footer"),
447 SchemaOrigin::AllFooters(n) => {
448 write!(f, "all {} footers", crate::numfmt::group_chrome(*n))
449 }
450 SchemaOrigin::FooterSample { read, total } => write!(
451 f,
452 "{} of {} footers (sample)",
453 crate::numfmt::group_chrome(*read),
454 crate::numfmt::group_chrome(*total)
455 ),
456 }
457 }
458}
459
460#[derive(Debug, Clone, PartialEq)]
462pub struct ColumnDrift {
463 pub name: PlSmallStr,
464 pub dtype: DataType,
466 pub present_in: usize,
468 pub conflicting_files: usize,
470 pub conflicting_types: Vec<DataType>,
472 pub widened: bool,
474}
475
476impl ColumnDrift {
477 pub fn is_uniform(&self, files: usize) -> bool {
479 self.present_in == files && self.conflicting_files == 0 && !self.widened
480 }
481}
482
483#[derive(Debug, Clone, PartialEq, Eq)]
490pub enum ColumnRange {
491 Only(String),
493 NoneBefore(String),
495}
496
497#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
515pub struct SkippedFiles {
516 pub bookkeeping: usize,
518 pub not_parquet: usize,
520 pub empty: usize,
526}
527
528impl SkippedFiles {
529 pub fn count(&mut self, bookkeeping: bool) {
535 if bookkeeping {
536 self.bookkeeping += 1;
537 } else {
538 self.not_parquet += 1;
539 }
540 }
541}
542
543#[derive(Debug, Clone)]
556pub struct ReadAs {
557 pub delimiter: Option<u8>,
558 pub has_header: Option<bool>,
559 pub skip_rows: Option<usize>,
560 pub skip_lines: Option<usize>,
561 pub infer_schema_length: Option<usize>,
562 pub ignore_errors: bool,
563 pub try_parse_dates: bool,
564 pub comment_char: Option<String>,
565 pub header_rows: Vec<usize>,
566 pub header_join: String,
567}
568
569impl ReadAs {
570 fn open_options(&self, format: crate::FileFormat) -> crate::OpenOptions {
573 crate::OpenOptions {
574 delimiter: self.delimiter.or(format.separator()),
575 has_header: self.has_header,
576 skip_rows: self.skip_rows,
577 skip_lines: self.skip_lines,
578 infer_schema_length: self.infer_schema_length,
579 ignore_errors: self.ignore_errors,
580 parse_dates: self.try_parse_dates,
581 parse_strings: None,
582 comment_char: self.comment_char.clone(),
583 header_rows: self.header_rows.clone(),
584 header_join: self.header_join.clone(),
585 ..crate::OpenOptions::default()
586 }
587 }
588}
589
590impl Default for ReadAs {
591 fn default() -> Self {
600 Self {
601 delimiter: None,
602 has_header: None,
603 skip_rows: None,
604 skip_lines: None,
605 infer_schema_length: None,
606 ignore_errors: false,
607 try_parse_dates: true,
608 comment_char: None,
609 header_rows: Vec::new(),
610 header_join: crate::csv_dialect::DEFAULT_HEADER_JOIN.to_string(),
611 }
612 }
613}
614
615pub fn column_schema_of(
632 path: &std::path::Path,
633 format: crate::FileFormat,
634 as_read: &ReadAs,
635) -> Option<Vec<(String, DataType)>> {
636 use polars::prelude::{LazyFileListReader, LazyJsonLineReader};
637 let lf = match format.descriptor().lines {
638 Some(crate::cli::Lines::Delimited(_)) => {
639 let options = as_read.open_options(format);
643 let header = crate::widgets::datatable::DataTableState::csv_header_names_of(
644 &options, path, None,
645 )
646 .ok()?;
647 let reader = crate::widgets::datatable::DataTableState::configure_csv_reader(
648 crate::widgets::datatable::DataTableState::csv_reader_of(path).ok()?,
649 &options,
650 None,
651 );
652 crate::csv_dialect::name_columns(reader.finish().ok()?, header.as_deref()).ok()?
653 }
654 Some(crate::cli::Lines::Json) => {
655 LazyJsonLineReader::new(crate::source::polars_literal_path(path).ok()?)
656 .finish()
657 .ok()?
658 }
659 Some(crate::cli::Lines::Text) => {
661 return Some(
662 crate::lines::schema(false)
663 .iter()
664 .map(|(name, dtype)| (name.to_string(), dtype.clone()))
665 .collect(),
666 );
667 }
668 None => return None,
669 };
670 let schema = lf.clone().collect_schema().ok()?;
671 let fields: Vec<(String, DataType)> = schema
676 .iter()
677 .map(|(name, dtype)| (name.trim().to_string(), dtype.clone()))
678 .collect();
679 Some(fields)
680}
681
682fn is_an_empty_file(schema: &[(String, DataType)]) -> bool {
698 match schema {
699 [] => true,
700 [(only, _)] => only.trim().is_empty(),
701 _ => false,
702 }
703}
704
705pub(crate) fn names_are_names(names: &[String]) -> bool {
720 !names.is_empty() && !names.iter().all(|n| n.trim().parse::<f64>().is_ok())
721}
722
723#[derive(Debug, Clone, Default)]
730pub struct Sampled {
731 pub columns: Vec<String>,
734 pub nests: Option<bool>,
737 pub columns_differ: bool,
740 pub types_differ: bool,
745 pub read: usize,
748 pub headerless: bool,
756}
757
758impl Sampled {
759 pub fn disagreement(&self) -> Disagreement {
761 if self.headerless {
764 return Disagreement {
765 headerless: true,
766 ..Default::default()
767 };
768 }
769 Disagreement {
770 columns: self.columns_differ,
771 types: self.types_differ,
772 headerless: false,
773 }
774 }
775}
776
777#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
784pub struct Disagreement {
785 pub columns: bool,
786 pub types: bool,
787 pub headerless: bool,
789}
790
791impl Disagreement {
792 pub fn any(&self) -> bool {
793 self.columns || self.types || self.headerless
794 }
795}
796
797pub fn sample_files(
804 files: &[std::path::PathBuf],
805 format: crate::FileFormat,
806 as_read: &ReadAs,
807) -> Sampled {
808 const WANTED: usize = 3;
815 const TRIES: usize = 12;
816 let last = files.len().saturating_sub(1);
817 const NEAR: usize = 4;
828 let anchors = [0usize, last / 2, last];
829
830 let mut read: Vec<Vec<(String, DataType)>> = Vec::new();
831 let mut tried = 0usize;
832 let mut seen: Vec<usize> = Vec::new();
833 'anchors: for anchor in anchors {
834 for step in 0..NEAR {
835 if read.len() >= WANTED || tried >= TRIES {
836 break 'anchors;
837 }
838 let i = anchor + step;
839 if i > last || seen.contains(&i) {
840 continue;
841 }
842 seen.push(i);
843 let Some(file) = files.get(i) else { continue };
844 tried += 1;
845 if let Some(schema) = column_schema_of(file, format, as_read)
848 && !is_an_empty_file(&schema)
849 {
850 read.push(schema);
851 continue 'anchors;
852 }
853 }
854 }
855
856 let mut out = Sampled {
857 read: read.len(),
858 ..Default::default()
859 };
860 for file in &read {
861 for (name, _) in file {
862 if !out.columns.iter().any(|c| c == name) {
863 out.columns.push(name.clone());
864 }
865 }
866 }
867 if read.len() < 2 {
868 return out;
869 }
870 let names: Vec<Vec<String>> = read
871 .iter()
872 .map(|f| f.iter().map(|(n, _)| n.clone()).collect())
873 .collect();
874 if names.iter().any(|f| !names_are_names(f)) {
881 out.columns.clear();
882 out.headerless = true;
883 out.nests = Some(false);
884 return out;
885 }
886 let nests = is_nested(&names);
887 out.nests = Some(nests);
888 let mut types: HashMap<&str, &DataType> = HashMap::new();
890 let mut typed_apart = false;
891 for (name, dtype) in read.iter().flatten() {
892 match types.get(name.as_str()) {
893 Some(seen) if *seen != dtype => typed_apart = true,
894 Some(_) => {}
895 None => {
896 types.insert(name.as_str(), dtype);
897 }
898 }
899 }
900 let widest = names.iter().map(|f| f.len()).max().unwrap_or(0);
903 out.columns_differ = names.iter().any(|f| f.len() != widest) || !nests;
904 out.types_differ = typed_apart;
905 out
906}
907
908pub fn is_nested(files: &[Vec<String>]) -> bool {
932 let Some(widest) = files.iter().max_by_key(|f| f.len()) else {
933 return true;
934 };
935 let widest: std::collections::BTreeSet<&str> = widest.iter().map(String::as_str).collect();
936 files
937 .iter()
938 .all(|file| file.iter().all(|name| widest.contains(name.as_str())))
939}
940
941pub fn top_level_columns(leaves: &[String]) -> Vec<String> {
950 let mut seen = std::collections::HashSet::new();
951 leaves
952 .iter()
953 .map(|leaf| leaf.split_once('.').map_or(leaf.as_str(), |(root, _)| root))
954 .filter(|root| seen.insert(root.to_string()))
955 .map(str::to_string)
956 .collect()
957}
958
959#[derive(Debug, Clone)]
961pub struct DatasetSchema {
962 pub schema: Arc<Schema>,
963 pub columns: Vec<ColumnDrift>,
964 pub omitted: Vec<Vec<(PlSmallStr, DataType)>>,
971 pub unreadable: Vec<usize>,
973 pub files: usize,
975 pub groups: Vec<DriftGroup>,
978 pub file_group: Vec<u32>,
980 pub origin: SchemaOrigin,
981 pub read_as_text: Vec<PlSmallStr>,
984 pub empty_files: usize,
986 pub median_row_group_bytes: Option<usize>,
989 pub median_file_bytes: Option<usize>,
991 pub column_ranges: HashMap<PlSmallStr, ColumnRange>,
996 pub partition_layouts: Vec<(Vec<String>, usize)>,
1003 pub partition_layouts_dropped: (usize, usize),
1007 pub skipped: SkippedFiles,
1010 pub listed_files: usize,
1013}
1014
1015#[derive(Debug, Clone, Default, PartialEq, Eq, Hash)]
1018pub struct DriftGroup {
1019 pub absent: Vec<PlSmallStr>,
1021 pub unread: Vec<PlSmallStr>,
1024}
1025
1026impl DriftGroup {
1027 pub fn is_empty(&self) -> bool {
1028 self.absent.is_empty() && self.unread.is_empty()
1029 }
1030}
1031
1032impl DatasetSchema {
1033 pub fn drifting(&self) -> impl Iterator<Item = &ColumnDrift> {
1035 let readable = self.files - self.unreadable.len();
1036 self.columns.iter().filter(move |c| !c.is_uniform(readable))
1037 }
1038
1039 pub fn with_partition_layouts(mut self, root: &str, paths: &[String]) -> DatasetSchema {
1050 const KEPT: usize = 64;
1054 let mut counts: HashMap<Vec<String>, usize> = HashMap::new();
1058 for path in paths {
1059 let Some(below) = path.strip_prefix(root) else {
1064 continue;
1065 };
1066 let keys = partition_keys_of(below);
1067 if keys.is_empty() {
1068 continue;
1069 }
1070 *counts.entry(keys).or_insert(0) += 1;
1071 }
1072 let mut counts: Vec<(Vec<String>, usize)> = counts.into_iter().collect();
1073 counts.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
1077 let dropped = &counts[counts.len().min(KEPT)..];
1080 self.partition_layouts_dropped =
1081 (dropped.len(), dropped.iter().map(|(_, files)| files).sum());
1082 counts.truncate(KEPT);
1083 self.partition_layouts = counts;
1084 self.listed_files = paths.len();
1085 self.column_ranges = self.ranges_of_columns(root, paths);
1086 self
1087 }
1088
1089 fn ranges_of_columns(&self, root: &str, paths: &[String]) -> HashMap<PlSmallStr, ColumnRange> {
1096 if matches!(self.origin, SchemaOrigin::FooterSample { .. })
1101 || !self.unreadable.is_empty()
1102 || self.file_group.len() != paths.len()
1103 {
1104 return HashMap::new();
1105 }
1106 struct Seen {
1110 first_present: Option<usize>,
1111 last_absent: Option<usize>,
1112 with: Vec<String>,
1118 without: Vec<String>,
1119 unplaced: bool,
1123 }
1124 let partition_of = |index: usize| -> Option<String> {
1125 let below = paths.get(index)?.strip_prefix(root)?;
1126 let values = partition_values_of(below);
1127 (!values.is_empty()).then(|| values.join("/"))
1128 };
1129 let reads_in_order = (0..paths.len())
1134 .filter_map(&partition_of)
1135 .collect::<Vec<_>>()
1136 .windows(2)
1137 .all(|pair| natural_cmp(&pair[0], &pair[1]) != std::cmp::Ordering::Greater);
1138 let drifting: Vec<&ColumnDrift> = self
1143 .columns
1144 .iter()
1145 .filter(|column| column.present_in > 0 && column.present_in < self.files)
1146 .collect();
1147 if drifting.is_empty() {
1148 return HashMap::new();
1149 }
1150 let where_in_drifting: HashMap<&PlSmallStr, usize> = drifting
1154 .iter()
1155 .enumerate()
1156 .map(|(at, column)| (&column.name, at))
1157 .collect();
1158 let missing_by_group: Vec<Vec<bool>> = self
1159 .groups
1160 .iter()
1161 .map(|group| {
1162 let mut missing = vec![false; drifting.len()];
1163 for name in &group.absent {
1164 if let Some(at) = where_in_drifting.get(name) {
1165 missing[*at] = true;
1166 }
1167 }
1168 missing
1169 })
1170 .collect();
1171 let none_missing: Vec<bool> = vec![false; drifting.len()];
1172 let mut seen: Vec<Seen> = (0..drifting.len())
1176 .map(|_| Seen {
1177 first_present: None,
1178 last_absent: None,
1179 with: Vec::new(),
1180 without: Vec::new(),
1181 unplaced: false,
1182 })
1183 .collect();
1184
1185 for (index, group) in self.file_group.iter().enumerate() {
1186 let missing: &[bool] = missing_by_group
1187 .get(*group as usize)
1188 .map(Vec::as_slice)
1189 .unwrap_or(&none_missing);
1190 let here = partition_of(index);
1192 for (at, absent) in missing.iter().enumerate() {
1193 let entry = &mut seen[at];
1194 let seen_of = if *absent {
1195 entry.last_absent = Some(index);
1196 &mut entry.without
1197 } else {
1198 entry.first_present.get_or_insert(index);
1199 entry.unplaced |= here.is_none();
1200 &mut entry.with
1201 };
1202 if seen_of.len() < 2
1203 && let Some(partition) = here.as_ref()
1204 && !seen_of.contains(partition)
1205 {
1206 seen_of.push(partition.clone());
1207 }
1208 }
1209 }
1210 seen.into_iter()
1211 .zip(&drifting)
1212 .filter_map(|(entry, column)| {
1213 let name = column.name.clone();
1214 let first = entry.first_present?;
1215 if !entry.unplaced
1219 && entry.with.len() == 1
1220 && entry
1223 .without
1224 .iter()
1225 .any(|other| {
1226 !partition_holds(&entry.with[0], other)
1227 && !same_place(&entry.with[0], other)
1228 })
1229 {
1230 return Some((name, ColumnRange::Only(entry.with[0].clone())));
1231 }
1232 let last_absent = entry.last_absent?;
1236 if last_absent > first || !reads_in_order {
1243 return None;
1244 }
1245 let (ends, begins) = (partition_of(last_absent)?, partition_of(first)?);
1246 (!same_place(&ends, &begins)
1251 && !partition_holds(&begins, &ends)
1252 && !partition_holds(&ends, &begins))
1253 .then_some((name, ColumnRange::NoneBefore(begins)))
1254 })
1255 .collect()
1256 }
1257
1258 pub fn with_skipped(mut self, skipped: SkippedFiles) -> DatasetSchema {
1260 self.skipped = skipped;
1261 self
1262 }
1263
1264 pub fn reading_as_text(&self, as_text: &[PlSmallStr]) -> DatasetSchema {
1275 let mut out = self.clone();
1276 if as_text.is_empty() {
1277 return out;
1278 }
1279 out.schema = crate::schema_union::text_schema(&self.schema, as_text);
1280 for column in &mut out.columns {
1281 if as_text.contains(&column.name) {
1282 column.dtype = DataType::String;
1283 column.conflicting_files = 0;
1284 column.conflicting_types.clear();
1285 }
1290 }
1291 for group in &mut out.groups {
1292 group.unread.retain(|name| !as_text.contains(name));
1293 }
1294 out.read_as_text = as_text.to_vec();
1295 out
1296 }
1297
1298 pub fn drifts(&self) -> bool {
1301 self.groups.iter().any(|g| !g.is_empty())
1302 }
1303}
1304
1305pub const FOOTERS_AT_ONCE: usize = 64;
1309
1310pub fn footers_to_cache(
1315 footers: &[Option<FileFooter>],
1316) -> (
1317 Vec<crate::cache::CachedFooter>,
1318 Vec<Vec<(String, DataType)>>,
1319) {
1320 let mut schemas = Vec::new();
1321 let cached = footers
1322 .iter()
1323 .map(|footer| match footer {
1324 None => crate::cache::CachedFooter::default(),
1325 Some(f) => crate::cache::CachedFooter {
1326 schema: Some(crate::cache::DatasetShape::intern_schema(
1327 &mut schemas,
1328 &f.schema,
1329 )),
1330 row_group_rows: f.row_group_rows.clone(),
1331 row_group_bytes: f.row_group_bytes.clone(),
1332 column_bytes: f.column_bytes.iter().map(|(_, bytes)| *bytes).collect(),
1333 },
1334 })
1335 .collect();
1336 (cached, schemas)
1337}
1338
1339pub fn footers_from_cache(
1346 cached: &[crate::cache::CachedFooter],
1347 schemas: &[Vec<(String, DataType)>],
1348 file_bytes: &[u64],
1349) -> Option<Vec<Option<FileFooter>>> {
1350 if cached.len() != file_bytes.len() {
1351 return None;
1352 }
1353 let shared: Vec<Arc<Schema>> = (0..schemas.len())
1355 .map(|at| crate::cache::DatasetShape::schema_at(schemas, at).map(Arc::new))
1356 .collect::<Option<_>>()?;
1357 cached
1358 .iter()
1359 .zip(file_bytes)
1360 .map(|(f, &bytes)| {
1361 let Some(at) = f.schema else {
1362 return Some(None);
1363 };
1364 let schema = shared.get(at)?;
1365 let column_bytes = schema
1366 .iter_names()
1367 .zip(&f.column_bytes)
1368 .map(|(name, bytes)| (name.to_string(), *bytes))
1369 .collect();
1370 Some(Some(FileFooter {
1371 schema: schema.clone(),
1372 row_group_rows: f.row_group_rows.clone(),
1373 row_group_bytes: f.row_group_bytes.clone(),
1374 file_bytes: bytes as usize,
1375 column_bytes,
1376 }))
1377 })
1378 .collect()
1379}
1380
1381type FooterHook = Arc<dyn Fn(&std::path::Path) + Send + Sync>;
1383
1384static FOOTER_HOOKS: std::sync::Mutex<Vec<(u64, std::path::PathBuf, FooterHook)>> =
1385 std::sync::Mutex::new(Vec::new());
1386static FOOTER_HOOKS_SET: AtomicUsize = AtomicUsize::new(0);
1388
1389#[doc(hidden)]
1393pub fn on_local_footer_read(
1394 dir: &std::path::Path,
1395 hook: impl Fn(&std::path::Path) + Send + Sync + 'static,
1396) -> FooterHookGuard {
1397 static NEXT: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1398 let id = NEXT.fetch_add(1, Ordering::Relaxed);
1399 FOOTER_HOOKS
1400 .lock()
1401 .unwrap_or_else(|e| e.into_inner())
1402 .push((id, dir.to_path_buf(), Arc::new(hook)));
1403 FOOTER_HOOKS_SET.fetch_add(1, Ordering::Release);
1404 FooterHookGuard(id)
1405}
1406
1407#[doc(hidden)]
1409pub struct FooterHookGuard(u64);
1410
1411impl Drop for FooterHookGuard {
1412 fn drop(&mut self) {
1413 FOOTER_HOOKS
1414 .lock()
1415 .unwrap_or_else(|e| e.into_inner())
1416 .retain(|(id, _, _)| *id != self.0);
1417 FOOTER_HOOKS_SET.fetch_sub(1, Ordering::Release);
1418 }
1419}
1420
1421pub(crate) fn before_local_footer_read(path: &std::path::Path) {
1423 if FOOTER_HOOKS_SET.load(Ordering::Acquire) == 0 {
1424 return;
1425 }
1426 let hooks: Vec<FooterHook> = FOOTER_HOOKS
1428 .lock()
1429 .unwrap_or_else(|e| e.into_inner())
1430 .iter()
1431 .filter(|(_, dir, _)| path.starts_with(dir))
1432 .map(|(_, _, hook)| hook.clone())
1433 .collect();
1434 for hook in hooks {
1435 hook(path);
1436 }
1437}
1438
1439pub const MAX_FOOTER_READS: usize = 20_000;
1443
1444pub fn footers_to_read(files: usize) -> Vec<usize> {
1448 if files <= MAX_FOOTER_READS {
1449 return (0..files).collect();
1450 }
1451 let last = files - 1;
1452 let mut sample: Vec<usize> = (0..MAX_FOOTER_READS)
1453 .map(|i| i * last / (MAX_FOOTER_READS - 1))
1454 .collect();
1455 sample.dedup();
1456 sample
1457}
1458
1459pub fn ends_of(files: usize) -> Vec<usize> {
1467 match files {
1468 0 => Vec::new(),
1469 1 => vec![0],
1470 n => vec![0, n - 1],
1471 }
1472}
1473
1474pub fn union_sampled(
1486 files: usize,
1487 read: &[usize],
1488 footers: &[Option<FileFooter>],
1489) -> DatasetSchema {
1490 let origin = if read.len() == files {
1491 SchemaOrigin::AllFooters(files)
1492 } else {
1493 SchemaOrigin::FooterSample {
1494 read: read.len(),
1495 total: files,
1496 }
1497 };
1498 let mut union = union_file_schemas(footers, origin);
1499 let mut omitted = vec![Vec::new(); files];
1500 let mut file_group = vec![0u32; files];
1501 for ((columns, group), &index) in union.omitted.iter().zip(union.file_group.iter()).zip(read) {
1502 omitted[index] = columns.clone();
1503 file_group[index] = *group;
1504 }
1505 union.omitted = omitted;
1506 union.file_group = file_group;
1507 union.unreadable = union
1508 .unreadable
1509 .iter()
1510 .filter_map(|i| read.get(*i).copied())
1511 .collect();
1512 union
1513}
1514
1515pub fn readable_paths<'a>(paths: &'a [String], unreadable: &[usize]) -> Cow<'a, [String]> {
1523 if unreadable.is_empty() {
1524 return Cow::Borrowed(paths);
1526 }
1527 debug_assert!(unreadable.windows(2).all(|pair| pair[0] < pair[1]));
1531 Cow::Owned(
1532 paths
1533 .iter()
1534 .enumerate()
1535 .filter(|(index, _)| unreadable.binary_search(index).is_err())
1536 .map(|(_, path)| path.clone())
1537 .collect(),
1538 )
1539}
1540
1541pub fn union_file_schemas(files: &[Option<FileFooter>], origin: SchemaOrigin) -> DatasetSchema {
1544 let unreadable = files
1545 .iter()
1546 .enumerate()
1547 .filter_map(|(i, f)| f.is_none().then_some(i))
1548 .collect();
1549
1550 let mut order: Vec<PlSmallStr> = Vec::new();
1552 let mut seen: HashMap<PlSmallStr, usize> = HashMap::new();
1553 let mut push = |name: &PlSmallStr, order: &mut Vec<PlSmallStr>| {
1554 if !seen.contains_key(name) {
1555 seen.insert(name.clone(), order.len());
1556 order.push(name.clone());
1557 }
1558 };
1559 if let Some(newest) = files.iter().rev().flatten().next() {
1560 for name in newest.schema.iter_names() {
1561 push(name, &mut order);
1562 }
1563 }
1564 for file in files.iter().flatten() {
1565 for name in file.schema.iter_names() {
1566 push(name, &mut order);
1567 }
1568 }
1569
1570 let mut sightings: Vec<Vec<(DataType, usize)>> = vec![Vec::new(); order.len()];
1572 for file in files.iter().flatten() {
1573 for (name, dtype) in file.schema.iter() {
1574 let Some(&index) = seen.get(name) else {
1575 continue;
1576 };
1577 sightings[index].push((dtype.clone(), file.rows()));
1578 }
1579 }
1580
1581 let mut schema = Schema::with_capacity(order.len());
1582 let mut columns = Vec::with_capacity(order.len());
1583 for (name, seen_types) in order.iter().zip(sightings.iter()) {
1584 let chosen = choose_dtype(seen_types);
1585 let conflicting_types = seen_types
1586 .iter()
1587 .map(|(d, _)| d)
1588 .filter(|d| !fits(d, &chosen))
1589 .fold(Vec::new(), |mut acc: Vec<DataType>, d| {
1590 if !acc.contains(d) {
1591 acc.push(d.clone());
1592 }
1593 acc
1594 });
1595 columns.push(ColumnDrift {
1596 name: name.clone(),
1597 present_in: seen_types.len(),
1598 conflicting_files: seen_types.iter().filter(|(d, _)| !fits(d, &chosen)).count(),
1599 widened: seen_types
1600 .iter()
1601 .any(|(d, _)| *d != chosen && fits(d, &chosen)),
1602 conflicting_types,
1603 dtype: chosen.clone(),
1604 });
1605 schema.with_column(name.clone(), chosen);
1606 }
1607
1608 let mut groups: Vec<DriftGroup> = vec![DriftGroup::default()];
1612 let mut group_of: HashMap<DriftGroup, u32> = HashMap::from([(DriftGroup::default(), 0)]);
1613 let mut file_group = Vec::with_capacity(files.len());
1614 let mut omitted = Vec::with_capacity(files.len());
1615 for file in files {
1616 let Some(file) = file else {
1617 file_group.push(0);
1618 omitted.push(Vec::new());
1619 continue;
1620 };
1621 let stored: Vec<(PlSmallStr, DataType)> = file
1622 .schema
1623 .iter()
1624 .filter(|(name, dtype)| schema.get(name).is_some_and(|target| !fits(dtype, target)))
1625 .map(|(name, dtype)| (name.clone(), dtype.clone()))
1626 .collect();
1627 let unread: Vec<PlSmallStr> = stored.iter().map(|(name, _)| name.clone()).collect();
1628 let absent: Vec<PlSmallStr> = schema
1629 .iter_names()
1630 .filter(|name| !file.schema.contains(name))
1631 .cloned()
1632 .collect();
1633 omitted.push(stored);
1634 let group = DriftGroup { absent, unread };
1635 let next = groups.len() as u32;
1636 let id = *group_of.entry(group.clone()).or_insert_with(|| {
1637 groups.push(group);
1638 next
1639 });
1640 file_group.push(id);
1641 }
1642
1643 DatasetSchema {
1644 schema: Arc::new(schema),
1645 columns,
1646 omitted,
1647 unreadable,
1648 files: files.len(),
1649 groups,
1650 file_group,
1651 origin,
1652 read_as_text: Vec::new(),
1653 empty_files: files.iter().flatten().filter(|f| f.rows() == 0).count(),
1654 median_file_bytes: median(files.iter().flatten().map(|f| f.file_bytes)),
1655 column_ranges: HashMap::new(),
1656 skipped: SkippedFiles::default(),
1657 partition_layouts: Vec::new(),
1658 partition_layouts_dropped: (0, 0),
1659 listed_files: 0,
1660 median_row_group_bytes: median(
1661 files
1662 .iter()
1663 .flatten()
1664 .flat_map(|f| f.row_group_bytes.iter().copied()),
1665 ),
1666 }
1667}
1668
1669fn partition_holds(outer: &str, inner: &str) -> bool {
1676 inner == outer
1677 || inner
1678 .strip_prefix(outer)
1679 .is_some_and(|rest| rest.starts_with('/'))
1680}
1681
1682fn same_place(a: &str, b: &str) -> bool {
1689 fn sorted(path: &str) -> Vec<&str> {
1690 let mut segments: Vec<&str> = path.split('/').collect();
1691 segments.sort_by_key(|segment| segment.split_once('=').map(|(key, _)| key));
1692 segments
1693 }
1694 let (a, b) = (sorted(a), sorted(b));
1695 a.len() == b.len()
1696 && a.iter()
1697 .zip(&b)
1698 .all(|(x, y)| natural_cmp(x, y) == std::cmp::Ordering::Equal)
1699}
1700
1701fn natural_cmp(a: &str, b: &str) -> std::cmp::Ordering {
1707 use std::cmp::Ordering;
1708 let (mut a, mut b) = (a.as_bytes(), b.as_bytes());
1709 loop {
1710 match (a.first(), b.first()) {
1711 (None, None) => return Ordering::Equal,
1712 (None, _) => return Ordering::Less,
1713 (_, None) => return Ordering::Greater,
1714 (Some(x), Some(y)) if x.is_ascii_digit() && y.is_ascii_digit() => {
1715 let digits = |s: &[u8]| s.iter().take_while(|c| c.is_ascii_digit()).count();
1716 let (na, nb) = (digits(a), digits(b));
1717 let (xs, ys) = (&a[..na], &b[..nb]);
1721 fn trim(s: &[u8]) -> &[u8] {
1722 let lead = s.iter().take_while(|c| **c == b'0').count();
1723 &s[lead.min(s.len().saturating_sub(1))..]
1724 }
1725 let (tx, ty) = (trim(xs), trim(ys));
1726 match tx.len().cmp(&ty.len()).then_with(|| tx.cmp(ty)) {
1727 Ordering::Equal => {}
1728 other => return other,
1729 }
1730 a = &a[na..];
1731 b = &b[nb..];
1732 }
1733 (Some(x), Some(y)) => match x.cmp(y) {
1734 Ordering::Equal => {
1735 a = &a[1..];
1736 b = &b[1..];
1737 }
1738 other => return other,
1739 },
1740 }
1741 }
1742}
1743
1744fn partition_values_of(path: &str) -> Vec<String> {
1755 #[cfg(windows)]
1756 let separators: &[char] = &['/', '\\'];
1757 #[cfg(not(windows))]
1758 let separators: &[char] = &['/'];
1759 let mut segments: Vec<&str> = path.split(separators).collect();
1760 segments.pop();
1762 segments
1763 .into_iter()
1764 .filter(|segment| {
1765 segment
1766 .split_once('=')
1767 .is_some_and(|(key, _)| !key.is_empty())
1768 })
1769 .map(|segment| segment.to_string())
1770 .collect()
1771}
1772
1773fn partition_keys_of(path: &str) -> Vec<String> {
1778 let mut keys: Vec<String> = Vec::new();
1779 #[cfg(windows)]
1783 let separators: &[char] = &['/', '\\'];
1784 #[cfg(not(windows))]
1785 let separators: &[char] = &['/'];
1786 let mut segments: Vec<&str> = path.split(separators).collect();
1787 segments.pop();
1788 for segment in segments {
1789 if let Some((key, _)) = segment.split_once('=')
1790 && !key.is_empty()
1791 {
1792 keys.push(key.to_string());
1793 }
1794 }
1795 keys.sort();
1800 keys.dedup();
1801 keys
1802}
1803
1804fn median(sizes: impl Iterator<Item = usize>) -> Option<usize> {
1810 let mut sizes: Vec<usize> = sizes.collect();
1811 if sizes.is_empty() {
1812 return None;
1813 }
1814 sizes.sort_unstable();
1815 Some(sizes[(sizes.len() - 1) / 2])
1816}
1817
1818fn choose_dtype(seen: &[(DataType, usize)]) -> DataType {
1824 let mut distinct: Vec<DataType> = Vec::new();
1825 for (dtype, _) in seen {
1826 if !distinct.contains(dtype) {
1827 distinct.push(dtype.clone());
1828 }
1829 }
1830 match distinct.as_slice() {
1831 [] => return DataType::Null,
1832 [only] => return only.clone(),
1833 _ => {}
1834 }
1835 let mut candidates = distinct.clone();
1838 for dtype in &distinct {
1839 let folded = distinct
1840 .iter()
1841 .filter(|other| widen(dtype, other).is_some())
1842 .try_fold(dtype.clone(), |acc, other| widen(&acc, other));
1843 if let Some(folded) = folded
1844 && !candidates.contains(&folded)
1845 {
1846 candidates.push(folded);
1847 }
1848 }
1849 let mut best: Option<(DataType, usize, usize)> = None;
1850 for candidate in candidates {
1851 let rows: usize = seen
1852 .iter()
1853 .filter(|(d, _)| fits(d, &candidate))
1854 .map(|(_, rows)| rows)
1855 .sum();
1856 let files = seen.iter().filter(|(d, _)| fits(d, &candidate)).count();
1857 let better = best
1858 .as_ref()
1859 .is_none_or(|(_, best_rows, best_files)| (rows, files) > (*best_rows, *best_files));
1860 if better {
1861 best = Some((candidate, rows, files));
1862 }
1863 }
1864 best.map(|(d, _, _)| d).unwrap_or(DataType::Null)
1865}
1866
1867pub fn fits(from: &DataType, to: &DataType) -> bool {
1871 widen(from, to).as_ref() == Some(to)
1872}
1873
1874pub fn widen(a: &DataType, b: &DataType) -> Option<DataType> {
1877 use DataType::*;
1878 if a == b {
1879 return Some(a.clone());
1880 }
1881 match (a, b) {
1882 (Null, other) | (other, Null) => Some(other.clone()),
1883 _ if a.is_integer() && b.is_integer() => widen_integers(a, b),
1884 _ if (a.is_integer() || a.is_float()) && (b.is_integer() || b.is_float()) => Some(Float64),
1885 (Datetime(a_unit, a_zone), Datetime(b_unit, b_zone)) if a_zone == b_zone => {
1886 Some(Datetime(finer_unit(*a_unit, *b_unit), a_zone.clone()))
1887 }
1888 (List(a_inner), List(b_inner)) => widen(a_inner, b_inner).map(|t| List(Box::new(t))),
1889 (Struct(a_fields), Struct(b_fields)) => widen_structs(a_fields, b_fields),
1890 _ => None,
1891 }
1892}
1893
1894fn widen_integers(a: &DataType, b: &DataType) -> Option<DataType> {
1897 use DataType::*;
1898 let signed = |d: &DataType| matches!(d, Int8 | Int16 | Int32 | Int64 | Int128);
1899 let bits = |d: &DataType| match d {
1900 Int8 | UInt8 => 8u32,
1901 Int16 | UInt16 => 16,
1902 Int32 | UInt32 => 32,
1903 Int64 | UInt64 => 64,
1904 _ => 128,
1905 };
1906 if signed(a) == signed(b) {
1907 let wider = if bits(a) >= bits(b) { a } else { b };
1908 return Some(wider.clone());
1909 }
1910 let (unsigned, sgn) = if signed(a) { (b, a) } else { (a, b) };
1912 let needed = match bits(unsigned) {
1913 8 => Int16,
1914 16 => Int32,
1915 32 => Int64,
1916 64 => return None,
1917 _ => return None,
1918 };
1919 Some(if bits(sgn) >= bits(&needed) {
1920 sgn.clone()
1921 } else {
1922 needed
1923 })
1924}
1925
1926fn widen_structs(a: &[Field], b: &[Field]) -> Option<DataType> {
1929 let mut fields: Vec<Field> = Vec::with_capacity(a.len() + b.len());
1930 for field in a {
1931 let widened = match b.iter().find(|other| other.name() == field.name()) {
1932 Some(other) => widen(field.dtype(), other.dtype())?,
1933 None => field.dtype().clone(),
1934 };
1935 fields.push(Field::new(field.name().clone(), widened));
1936 }
1937 for field in b {
1938 if !a.iter().any(|other| other.name() == field.name()) {
1939 fields.push(field.clone());
1940 }
1941 }
1942 Some(DataType::Struct(fields))
1943}
1944
1945fn finer_unit(a: TimeUnit, b: TimeUnit) -> TimeUnit {
1946 let rank = |u: TimeUnit| match u {
1947 TimeUnit::Milliseconds => 0,
1948 TimeUnit::Microseconds => 1,
1949 TimeUnit::Nanoseconds => 2,
1950 };
1951 if rank(a) >= rank(b) { a } else { b }
1952}
1953
1954pub fn with_partition_columns(
1956 file_schema: &Schema,
1957 partition_columns: &[String],
1958 values: &[(String, String)],
1959) -> Schema {
1960 let part_set: HashSet<&str> = partition_columns.iter().map(String::as_str).collect();
1961 let mut merged = Schema::with_capacity(partition_columns.len() + file_schema.len());
1962 for name in partition_columns {
1963 merged.with_column(
1964 name.clone().into(),
1965 crate::widgets::datatable::partition_dtype(name, file_schema, values),
1966 );
1967 }
1968 for (name, dtype) in file_schema.iter() {
1969 if !part_set.contains(name.as_str()) {
1970 merged.with_column(name.clone(), dtype.clone());
1971 }
1972 }
1973 merged
1974}
1975
1976pub fn partition_columns_of_key(key: &str) -> Vec<String> {
1979 let mut columns = Vec::new();
1980 let mut seen = HashSet::new();
1981 for segment in key.split('/') {
1982 if let Some((name, _)) = segment.split_once('=')
1983 && !name.is_empty()
1984 && seen.insert(name.to_string())
1985 {
1986 columns.push(name.to_string());
1987 }
1988 }
1989 columns
1990}
1991
1992pub fn partitions_of_listing(first: &str, newest: &str) -> (Vec<String>, Vec<(String, String)>) {
1997 let values = [first, newest]
1998 .iter()
1999 .flat_map(|key| key.split('/'))
2000 .filter_map(|segment| segment.split_once('='))
2001 .map(|(k, v)| (k.to_string(), v.to_string()))
2002 .collect();
2003 (partition_columns_of_key(newest), values)
2004}
2005
2006pub struct FooterCount<F> {
2012 files: usize,
2013 counted: Vec<usize>,
2016 footers: std::sync::Mutex<Vec<Option<F>>>,
2018}
2019
2020pub struct Counted<F> {
2022 pub row_groups: Vec<Vec<usize>>,
2024 pub whole: Option<Vec<Option<F>>>,
2026}
2027
2028impl<F: Clone> FooterCount<F> {
2029 pub fn new(
2032 files: usize,
2033 counted: Vec<usize>,
2034 known: impl IntoIterator<Item = (usize, Option<F>)>,
2035 ) -> Self {
2036 let mut footers = vec![None; files];
2037 for (index, footer) in known {
2038 if let Some(slot) = footers.get_mut(index) {
2039 *slot = footer;
2040 }
2041 }
2042 Self {
2043 files,
2044 counted,
2045 footers: std::sync::Mutex::new(footers),
2046 }
2047 }
2048
2049 pub fn count(
2055 &self,
2056 read: impl FnOnce(&[usize]) -> Option<Vec<Option<F>>>,
2057 row_groups: impl Fn(&F) -> Vec<usize>,
2058 ) -> Option<Counted<F>> {
2059 let mut footers = self.footers.lock().unwrap_or_else(|e| e.into_inner());
2060 let missing = self.missing(&mut footers);
2061 let read = if missing.is_empty() {
2062 Vec::new()
2063 } else {
2064 read(&missing)?
2065 };
2066 Some(self.settle(&mut footers, missing, read, row_groups))
2067 }
2068
2069 fn missing(&self, footers: &mut Vec<Option<F>>) -> Vec<usize> {
2070 if footers.is_empty() {
2071 *footers = vec![None; self.files];
2073 }
2074 self.counted
2075 .iter()
2076 .copied()
2077 .filter(|&index| footers[index].is_none())
2078 .collect()
2079 }
2080
2081 fn settle(
2082 &self,
2083 footers: &mut Vec<Option<F>>,
2084 missing: Vec<usize>,
2085 read: Vec<Option<F>>,
2086 row_groups: impl Fn(&F) -> Vec<usize>,
2087 ) -> Counted<F> {
2088 if footers.is_empty() {
2089 *footers = vec![None; self.files];
2090 }
2091 for (index, footer) in missing.into_iter().zip(read) {
2092 footers[index] = footer;
2093 }
2094 let groups: Vec<Vec<usize>> = self
2095 .counted
2096 .iter()
2097 .map(|&index| footers[index].as_ref().map(&row_groups).unwrap_or_default())
2098 .collect();
2099 let whole = if footers.iter().all(Option::is_some) {
2100 Some(std::mem::take(footers))
2101 } else {
2102 if self.counted.iter().all(|&index| footers[index].is_some()) {
2103 footers.clear();
2106 }
2107 None
2108 };
2109 Counted {
2110 row_groups: groups,
2111 whole,
2112 }
2113 }
2114}
2115
2116pub const DRIFT_COLUMN: &str = "__datui_row";
2121
2122static NOTHING_MISSING: DriftGroup = DriftGroup {
2124 absent: Vec::new(),
2125 unread: Vec::new(),
2126};
2127
2128#[derive(Debug, Clone, Default)]
2131pub struct ScanDrift {
2132 group_of: HashMap<String, u32>,
2133 row_of: HashMap<String, usize>,
2136 stored_of: HashMap<String, Vec<(PlSmallStr, DataType)>>,
2139 pub groups: Vec<DriftGroup>,
2140}
2141
2142impl ScanDrift {
2143 pub fn new(paths: &[String], dataset: &DatasetSchema, file_rows: &[usize]) -> Option<Self> {
2148 if !dataset.drifts() || file_rows.len() != paths.len() {
2149 return None;
2150 }
2151 let group_of = paths
2152 .iter()
2153 .zip(dataset.file_group.iter())
2154 .filter(|(_, group)| **group != 0)
2155 .map(|(path, group)| (path.clone(), *group))
2156 .collect();
2157 let mut row = 0usize;
2158 let mut row_of = HashMap::with_capacity(paths.len());
2159 for (path, rows) in paths.iter().zip(file_rows) {
2160 row_of.insert(path.clone(), row);
2161 row += rows;
2162 }
2163 let stored_of = paths
2164 .iter()
2165 .zip(dataset.omitted.iter())
2166 .filter(|(_, stored)| !stored.is_empty())
2167 .map(|(path, stored)| (path.clone(), stored.clone()))
2168 .collect();
2169 Some(ScanDrift {
2170 group_of,
2171 row_of,
2172 stored_of,
2173 groups: dataset.groups.clone(),
2174 })
2175 }
2176
2177 pub fn group(&self, path: &str) -> u32 {
2179 self.group_of.get(path).copied().unwrap_or(0)
2180 }
2181
2182 fn first_row(&self, path: &str) -> usize {
2184 self.row_of.get(path).copied().unwrap_or(0)
2185 }
2186
2187 fn stored_type(&self, path: &str, column: &PlSmallStr) -> Option<&DataType> {
2190 self.stored_of
2191 .get(path)?
2192 .iter()
2193 .find(|(name, _)| name == column)
2194 .map(|(_, dtype)| dtype)
2195 }
2196
2197 fn unread(&self, path: &str) -> &[PlSmallStr] {
2199 self.groups
2200 .get(self.group(path) as usize)
2201 .unwrap_or(&NOTHING_MISSING)
2202 .unread
2203 .as_slice()
2204 }
2205}
2206
2207pub fn lenient_scan(
2242 paths: &[String],
2243 schema: Arc<Schema>,
2244 cloud_options: Option<polars::io::cloud::CloudOptions>,
2245 drift: Option<&ScanDrift>,
2246 as_text: &[PlSmallStr],
2247) -> PolarsResult<LazyFrame> {
2248 let Some(drift) = drift else {
2249 return scan_run(paths, &schema, cloud_options, &[], None, &[]);
2250 };
2251 let as_text: Vec<PlSmallStr> = as_text
2256 .iter()
2257 .filter(|name| {
2258 schema.get(name).is_some_and(can_read_as_text)
2259 && paths
2260 .iter()
2261 .all(|path| drift.stored_type(path, name).is_none_or(can_read_as_text))
2262 })
2263 .cloned()
2264 .collect();
2265 let as_text = as_text.as_slice();
2266 let unread_of = |path: &str| -> Vec<PlSmallStr> {
2268 drift
2269 .unread(path)
2270 .iter()
2271 .filter(|name| !as_text.contains(name))
2272 .cloned()
2273 .collect()
2274 };
2275 let key_of = |path: &str| -> (Vec<PlSmallStr>, Vec<Option<DataType>>) {
2278 (
2279 unread_of(path),
2280 as_text
2281 .iter()
2282 .map(|name| drift.stored_type(path, name).cloned())
2283 .collect(),
2284 )
2285 };
2286 let mut runs: Vec<LazyFrame> = Vec::new();
2287 let mut start = 0;
2288 while start < paths.len() {
2289 let key = key_of(&paths[start]);
2290 let end = paths[start..]
2291 .iter()
2292 .position(|path| key_of(path) != key)
2293 .map_or(paths.len(), |offset| start + offset);
2294 let (omit, stored) = key;
2295 let read_as: Vec<(PlSmallStr, DataType)> = as_text
2297 .iter()
2298 .zip(stored)
2299 .map(|(name, stored)| {
2300 let dtype = stored.or_else(|| schema.get(name).cloned());
2301 (name.clone(), dtype.unwrap_or(DataType::String))
2302 })
2303 .collect();
2304 runs.push(scan_run(
2305 &paths[start..end],
2306 &schema,
2307 cloud_options.clone(),
2308 &omit,
2309 Some(drift.first_row(&paths[start])),
2310 &read_as,
2311 )?);
2312 start = end;
2313 }
2314 match runs.len() {
2315 1 => Ok(runs.remove(0)),
2316 _ => concat(
2317 runs,
2318 UnionArgs {
2319 rechunk: false,
2320 parallel: true,
2321 ..Default::default()
2322 },
2323 ),
2324 }
2325}
2326
2327pub fn can_read_as_text(dtype: &DataType) -> bool {
2340 match dtype {
2341 DataType::Binary | DataType::BinaryOffset => false,
2343 DataType::Duration(_) => false,
2344 DataType::List(_) | DataType::Array(_, _) => false,
2346 DataType::Struct(_) => true,
2349 DataType::Unknown(_) => false,
2350 _ => true,
2351 }
2352}
2353
2354impl ColumnDrift {
2355 pub fn can_read_as_text(&self) -> bool {
2360 can_read_as_text(&self.dtype) && self.conflicting_types.iter().all(can_read_as_text)
2361 }
2362}
2363
2364pub fn text_schema(schema: &Arc<Schema>, as_text: &[PlSmallStr]) -> Arc<Schema> {
2371 if as_text.is_empty() {
2372 return schema.clone();
2373 }
2374 let mut out = Schema::with_capacity(schema.len());
2375 for (name, dtype) in schema.iter() {
2376 let dtype = if as_text.contains(name) {
2377 DataType::String
2378 } else {
2379 dtype.clone()
2380 };
2381 out.with_column(name.clone(), dtype);
2382 }
2383 Arc::new(out)
2384}
2385
2386fn scan_run(
2390 urls: &[String],
2391 schema: &Arc<Schema>,
2392 cloud_options: Option<polars::io::cloud::CloudOptions>,
2393 omit: &[PlSmallStr],
2394 first_row: Option<usize>,
2395 read_as: &[(PlSmallStr, DataType)],
2396) -> PolarsResult<LazyFrame> {
2397 use polars::lazy::dsl::{
2398 CastColumnsPolicy, DslBuilder, ExtraColumnsPolicy, MissingColumnsPolicy, ScanSources,
2399 UnifiedScanArgs,
2400 };
2401 use polars::prelude::{Expr, NULL, col, lit};
2402 let sources = ScanSources::Paths(
2403 urls.iter()
2404 .map(|url| PlRefPath::new(url.as_str()))
2405 .collect(),
2406 );
2407 let target = if omit.is_empty() && read_as.is_empty() {
2408 schema.clone()
2409 } else {
2410 let mut reduced = Schema::with_capacity(schema.len());
2411 for (name, dtype) in schema.iter() {
2412 if omit.contains(name) {
2413 continue;
2414 }
2415 let dtype = read_as
2418 .iter()
2419 .find(|(column, _)| column == name)
2420 .map(|(_, dtype)| dtype)
2421 .unwrap_or(dtype);
2422 reduced.with_column(name.clone(), dtype.clone());
2423 }
2424 Arc::new(reduced)
2425 };
2426 let options = polars::io::parquet::read::ParquetOptions {
2427 schema: Some(target),
2428 ..Default::default()
2429 };
2430 let args = UnifiedScanArgs {
2431 cloud_options,
2432 hive_options: polars::io::HiveOptions::new_enabled(),
2433 glob: false,
2434 cast_columns_policy: CastColumnsPolicy {
2435 integer_upcast: true,
2436 integer_to_float_cast: true,
2437 float_upcast: true,
2438 datetime_nanoseconds_downcast: true,
2439 datetime_microseconds_downcast: true,
2440 datetime_milliseconds_upcast: true,
2441 datetime_microseconds_upcast: true,
2442 null_upcast: true,
2443 missing_struct_fields: MissingColumnsPolicy::Insert,
2444 extra_struct_fields: ExtraColumnsPolicy::Ignore,
2445 ..CastColumnsPolicy::ERROR_ON_MISMATCH
2446 },
2447 missing_columns_policy: MissingColumnsPolicy::Insert,
2448 extra_columns_policy: ExtraColumnsPolicy::Ignore,
2449 row_index: first_row.map(|first| polars::io::RowIndex {
2453 name: DRIFT_COLUMN.into(),
2454 offset: first as polars::prelude::IdxSize,
2455 }),
2456 ..Default::default()
2457 };
2458 let mut lf: LazyFrame = DslBuilder::scan_parquet(sources, options, args)?
2459 .build()
2460 .into();
2461 if !omit.is_empty() {
2462 let nulls: Vec<Expr> = omit
2463 .iter()
2464 .filter_map(|name| {
2465 let dtype = schema.get(name)?;
2466 Some(lit(NULL).cast(dtype.clone()).alias(name.clone()))
2467 })
2468 .collect();
2469 lf = lf.with_columns(nulls);
2470 }
2471 if !read_as.is_empty() {
2472 let texts: Vec<Expr> = read_as
2477 .iter()
2478 .map(|(name, _)| {
2479 crate::past_calendar::text_expr(col(name.clone()), CastOptions::NonStrict)
2480 .alias(name.clone())
2481 })
2482 .collect();
2483 lf = lf.with_columns(texts);
2484 }
2485 if first_row.is_some() {
2486 let mut ordered: Vec<Expr> = schema.iter_names().map(|name| col(name.clone())).collect();
2489 ordered.push(col(DRIFT_COLUMN));
2490 lf = lf.select(ordered);
2491 }
2492 Ok(lf)
2493}
2494
2495#[cfg(test)]
2496mod tests {
2497 use super::*;
2498
2499 #[test]
2505 fn the_sample_reads_a_file_named_like_a_glob() {
2506 let dir = tempfile::tempdir().unwrap();
2507 let names = |name: &str, format| {
2508 column_schema_of(&dir.path().join(name), format, &ReadAs::default())
2509 .unwrap()
2510 .into_iter()
2511 .map(|(n, _)| n)
2512 .collect::<Vec<_>>()
2513 };
2514 std::fs::write(dir.path().join("d[1].jsonl"), "{\"own\": 1}\n").unwrap();
2515 std::fs::write(dir.path().join("d1.jsonl"), "{\"other\": 1}\n").unwrap();
2516 std::fs::write(dir.path().join("d[1].csv"), "own\n1\n").unwrap();
2517 std::fs::write(dir.path().join("d1.csv"), "other\n1\n").unwrap();
2518 assert_eq!(names("d[1].jsonl", crate::FileFormat::Jsonl), ["own"]);
2519 assert_eq!(names("d[1].csv", crate::FileFormat::Csv), ["own"]);
2520 }
2521
2522 #[test]
2523 fn the_sample_splits_on_the_separator_the_open_uses() {
2524 let dir = tempfile::tempdir().unwrap();
2525 let names = |name: &str, body: &str, format, delimiter| {
2526 let path = dir.path().join(name);
2527 std::fs::write(&path, body).unwrap();
2528 let as_read = ReadAs {
2529 delimiter,
2530 ..ReadAs::default()
2531 };
2532 column_schema_of(&path, format, &as_read)
2533 .unwrap()
2534 .into_iter()
2535 .map(|(n, _)| n)
2536 .collect::<Vec<_>>()
2537 };
2538 use crate::FileFormat::{Csv, Psv, Tsv};
2539 assert_eq!(names("a.tsv", "a\tb\n1\t2\n", Tsv, None), ["a", "b"]);
2540 assert_eq!(names("a.psv", "a|b\n1|2\n", Psv, None), ["a", "b"]);
2541 assert_eq!(names("a.csv", "a|b\n1|2\n", Csv, None), ["a|b"]);
2542 assert_eq!(names("b.csv", "a|b\n1|2\n", Csv, Some(b'|')), ["a", "b"]);
2543 }
2544
2545 fn file(columns: &[(&str, DataType)], rows: usize) -> Option<FileFooter> {
2546 let mut schema = Schema::with_capacity(columns.len());
2547 for (name, dtype) in columns {
2548 schema.with_column((*name).into(), dtype.clone());
2549 }
2550 Some(FileFooter {
2551 schema: Arc::new(schema),
2552 row_group_rows: vec![rows],
2553 file_bytes: 0,
2554 row_group_bytes: Vec::new(),
2555 column_bytes: Vec::new(),
2556 })
2557 }
2558
2559 #[test]
2562 fn column_widths_average_over_the_rows_read() {
2563 let with = |rows, bytes: &[(&str, usize)]| {
2564 let mut footer = file(&[], rows)?;
2565 footer.column_bytes = bytes.iter().map(|(n, b)| (n.to_string(), *b)).collect();
2566 Some(footer)
2567 };
2568 let footers = [
2569 with(3, &[("blob", 3_000), ("id", 24)]),
2570 None,
2571 with(1, &[("id", 8)]),
2572 ];
2573 assert_eq!(
2574 column_bytes_per_row(&footers),
2575 [("blob".to_string(), 750), ("id".to_string(), 8)]
2576 );
2577 assert!(column_bytes_per_row(&[with(0, &[("id", 0)])]).is_empty());
2578 }
2579
2580 fn cols(files: &[&[&str]]) -> Vec<Vec<String>> {
2582 files
2583 .iter()
2584 .map(|f| f.iter().map(|c| (*c).to_string()).collect())
2585 .collect()
2586 }
2587
2588 fn grew(files: usize, from: usize, to: usize) -> Vec<Vec<String>> {
2591 let names = |n: usize| (0..n).map(|i| format!("c{i}")).collect::<Vec<_>>();
2592 let mut out: Vec<Vec<String>> = (0..files - 1).map(|_| names(from)).collect();
2593 out.push(names(to));
2594 out
2595 }
2596
2597 #[test]
2600 fn one_table_whatever_its_files_did_over_time() {
2601 for (what, files) in [
2602 (
2603 "identical part files",
2604 cols(&[&["a", "b", "c"], &["a", "b", "c"], &["a", "b", "c"]]),
2605 ),
2606 (
2607 "a column only one file has",
2608 cols(&[&["id"], &["id", "oops"], &["id"]]),
2609 ),
2610 (
2611 "a column that starts",
2612 cols(&[&["id", "ts"], &["id", "ts"], &["id", "ts", "fee"]]),
2613 ),
2614 (
2615 "a column that stops",
2616 cols(&[&["id", "ts", "fee"], &["id", "ts"], &["id", "ts"]]),
2617 ),
2618 (
2619 "a file truncated to one column",
2620 cols(&[&["a", "b", "c", "d"], &["a"], &["a", "b", "c", "d"]]),
2621 ),
2622 ("five columns grown to fifty", grew(10, 5, 50)),
2623 ("one file", cols(&[&["a", "b"]])),
2624 ("no files", Vec::new()),
2625 ] {
2626 assert!(is_nested(&files), "{what} should read as one table");
2627 }
2628 }
2629
2630 #[test]
2642 fn a_column_each_way_is_not_nesting_and_costs_a_keystroke() {
2643 for (what, files) in [
2644 (
2645 "one column each way",
2646 cols(&[&["a", "b", "c", "d", "e"], &["a", "b", "c", "d", "f"]]),
2647 ),
2648 (
2649 "a column renamed",
2650 cols(&[&["id", "ts", "amount"], &["id", "ts", "amt"]]),
2651 ),
2652 ] {
2653 assert!(
2654 !is_nested(&files),
2655 "{what} brings a column the widest file cannot account for"
2656 );
2657 }
2658 }
2659
2660 #[test]
2662 fn separate_tables_are_not_one_table() {
2663 for (what, files) in [
2664 ("two tables", cols(&[&["a", "b", "c"], &["x", "y", "z"]])),
2665 (
2666 "tables sharing a key",
2667 cols(&[&["id", "a", "b"], &["id", "x", "y"], &["id", "p", "q"]]),
2668 ),
2669 (
2670 "a season of six tables",
2676 cols(&[
2677 &[
2678 "season",
2679 "circuit_id",
2680 "url",
2681 "circuit_name",
2682 "lat",
2683 "lng",
2684 "locality",
2685 "country",
2686 ],
2687 &["season", "constructor_id", "url", "name", "nationality"],
2688 &[
2689 "season",
2690 "round",
2691 "driver_id",
2692 "position",
2693 "points",
2694 "wins",
2695 "constructor_id",
2696 ],
2697 &[
2698 "season",
2699 "driver_id",
2700 "permanent_number",
2701 "code",
2702 "url",
2703 "given_name",
2704 "family_name",
2705 "date_of_birth",
2706 "nationality",
2707 ],
2708 &[
2709 "season",
2710 "round",
2711 "race_name",
2712 "circuit_id",
2713 "race_date",
2714 "driver_id",
2715 "constructor_id",
2716 "number",
2717 "grid",
2718 "position",
2719 "position_text",
2720 "points",
2721 "laps",
2722 "status",
2723 "time_millis",
2724 "time_text",
2725 "fastest_lap_rank",
2726 "fastest_lap_number",
2727 "fastest_lap_time",
2728 "fastest_lap_avg_speed",
2729 ],
2730 &[
2731 "season",
2732 "round",
2733 "race_name",
2734 "circuit_id",
2735 "circuit_name",
2736 "locality",
2737 "country",
2738 "lat",
2739 "lng",
2740 "date",
2741 "time",
2742 "qualifying_date",
2743 "qualifying_time",
2744 "sprint_date",
2745 "sprint_time",
2746 "sprint_shootout_date",
2747 "sprint_shootout_time",
2748 "url",
2749 ],
2750 ]),
2751 ),
2752 ] {
2753 assert!(!is_nested(&files), "{what} should be separate tables");
2754 }
2755 }
2756
2757 #[test]
2762 fn growth_nests_where_unrelated_tables_do_not() {
2763 assert!(is_nested(&grew(10, 5, 50)));
2764 let unrelated = cols(&[&["id", "a", "b"], &["id", "x", "y"], &["id", "p", "q"]]);
2765 assert!(!is_nested(&unrelated));
2766 }
2767
2768 #[test]
2771 fn two_tables_sharing_a_key_do_not_nest() {
2772 let files = cols(&[&["id", "name"], &["id", "customer_id", "amount"]]);
2773 assert!(!is_nested(&files));
2774 }
2775
2776 #[test]
2779 fn a_nested_column_is_one_column_however_it_was_written() {
2780 let old_writer = vec![
2781 "id".to_string(),
2782 "tags.array".to_string(),
2783 "refs.array".to_string(),
2784 ];
2785 let new_writer = vec![
2786 "id".to_string(),
2787 "tags.list.element".to_string(),
2788 "refs.list.element".to_string(),
2789 ];
2790 let files = vec![
2791 top_level_columns(&old_writer),
2792 top_level_columns(&new_writer),
2793 ];
2794 assert_eq!(files[0], vec!["id", "tags", "refs"]);
2795 assert!(is_nested(&files), "the same three columns, written twice");
2796 assert!(!is_nested(&[old_writer, new_writer]));
2799 }
2800
2801 #[test]
2804 fn an_empty_file_does_not_decide_the_directory() {
2805 assert!(is_nested(&cols(&[&["a", "b"], &[], &["a", "b"]])));
2806 assert!(is_nested(&cols(&[&[], &[]])), "nothing to disagree about");
2807 }
2808
2809 fn union(files: &[Option<FileFooter>]) -> DatasetSchema {
2810 union_file_schemas(files, SchemaOrigin::AllFooters(files.len()))
2811 }
2812
2813 fn omitted_names(union: &DatasetSchema, index: usize) -> Vec<String> {
2815 union.omitted[index]
2816 .iter()
2817 .map(|(name, _)| name.to_string())
2818 .collect()
2819 }
2820
2821 fn names(schema: &Schema) -> Vec<String> {
2822 schema.iter_names().map(|n| n.to_string()).collect()
2823 }
2824
2825 #[test]
2832 fn a_conflicting_column_read_as_text_shows_every_file_s_values() {
2833 use polars::prelude::{ParquetWriter, df};
2834
2835 let dir = tempfile::tempdir().unwrap();
2836 let mut paths = Vec::new();
2837 let mut write = |name: &str, mut frame: polars::prelude::DataFrame| {
2838 let path = dir.path().join(name);
2839 let file = std::fs::File::create(&path).unwrap();
2840 ParquetWriter::new(file).finish(&mut frame).unwrap();
2841 paths.push(path.to_string_lossy().to_string());
2842 };
2843 write(
2845 "a.parquet",
2846 df!("id" => &[0i64, 1, 2], "n" => &[10i64, 20, 30]).unwrap(),
2847 );
2848 write("b.parquet", df!("id" => &[3i64], "n" => &["x"]).unwrap());
2849 write("c.parquet", df!("id" => &[4i64], "n" => &[true]).unwrap());
2852
2853 let footers: Vec<Option<FileFooter>> = vec![
2854 file(&[("id", DataType::Int64), ("n", DataType::Int64)], 3),
2855 file(&[("id", DataType::Int64), ("n", DataType::String)], 1),
2856 file(&[("id", DataType::Int64), ("n", DataType::Boolean)], 1),
2857 ];
2858 let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(3));
2859 assert_eq!(
2860 dataset.schema.get("n"),
2861 Some(&DataType::Int64),
2862 "the integer file has the most rows"
2863 );
2864 let drift = ScanDrift::new(&paths, &dataset, &[3, 1, 1]).expect("the files disagree");
2865
2866 let plain = lenient_scan(&paths, dataset.schema.clone(), None, Some(&drift), &[])
2868 .unwrap()
2869 .collect()
2870 .unwrap();
2871 let n = plain.column("n").unwrap();
2872 assert_eq!(
2873 (0..n.len())
2874 .map(|i| n.get(i).unwrap().to_string())
2875 .collect::<Vec<_>>(),
2876 ["10", "20", "30", "null", "null"],
2877 "the text and boolean files hold a value, and it is not one this column \
2878 can carry"
2879 );
2880
2881 let as_text = [PlSmallStr::from("n")];
2883 let text = lenient_scan(&paths, dataset.schema.clone(), None, Some(&drift), &as_text)
2884 .unwrap()
2885 .collect()
2886 .unwrap();
2887 assert_eq!(
2888 text.column("n").unwrap().dtype(),
2889 &DataType::String,
2890 "the column is text now"
2891 );
2892 let n = text.column("n").unwrap().str().unwrap();
2893 assert_eq!(
2894 n.iter().collect::<Vec<_>>(),
2895 [Some("10"), Some("20"), Some("30"), Some("x"), Some("true")],
2896 "and holds what each file wrote, spelled as that file's own type prints"
2897 );
2898 let ids = text.column("id").unwrap().i64().unwrap();
2899 assert_eq!(
2900 ids.into_no_null_iter().collect::<Vec<_>>(),
2901 [0, 1, 2, 3, 4],
2902 "in dataset order, with the rows still lined up against their ids"
2903 );
2904 assert_eq!(
2905 text.column(DRIFT_COLUMN)
2906 .unwrap()
2907 .u32()
2908 .unwrap()
2909 .into_no_null_iter()
2910 .collect::<Vec<_>>(),
2911 [0, 1, 2, 3, 4],
2912 "and each row still knows its place in the dataset"
2913 );
2914 }
2915
2916 #[test]
2919 fn a_date_past_the_calendar_read_as_text_is_its_stored_number() {
2920 use polars::prelude::{NamedFrom, ParquetWriter, Series, TimeZone};
2921
2922 let dir = tempfile::tempdir().unwrap();
2923 let mut paths = Vec::new();
2924 let mut write = |name: &str, n: Series| {
2925 let path = dir.path().join(name);
2926 let file = std::fs::File::create(&path).unwrap();
2927 let mut frame = polars::prelude::DataFrame::new_infer_height(vec![n.into()]).unwrap();
2928 ParquetWriter::new(file).finish(&mut frame).unwrap();
2929 paths.push(path.to_string_lossy().to_string());
2930 };
2931 let paris = TimeZone::opt_try_new(Some("Europe/Paris")).unwrap();
2932 let stamps = |dtype: DataType| {
2933 Series::new("n".into(), [0, i64::MIN + 1])
2934 .cast(&dtype)
2935 .unwrap()
2936 };
2937 let types = [
2938 DataType::Date,
2939 DataType::Datetime(TimeUnit::Milliseconds, None),
2940 DataType::Datetime(TimeUnit::Microseconds, paris),
2941 ];
2942 write("a.parquet", Series::new("n".into(), ["x", "y", "z"]));
2944 write(
2945 "b.parquet",
2946 Series::new("n".into(), [0, i32::MAX])
2947 .cast(&types[0])
2948 .unwrap(),
2949 );
2950 write("c.parquet", stamps(types[1].clone()));
2951 write("d.parquet", stamps(types[2].clone()));
2952
2953 let footers: Vec<Option<FileFooter>> = [DataType::String]
2954 .into_iter()
2955 .chain(types)
2956 .enumerate()
2957 .map(|(i, dtype)| file(&[("n", dtype)], if i == 0 { 3 } else { 2 }))
2958 .collect();
2959 let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(4));
2960 assert_eq!(dataset.schema.get("n"), Some(&DataType::String));
2961 let drift = ScanDrift::new(&paths, &dataset, &[3, 2, 2, 2]).expect("the files disagree");
2962 let as_text = [PlSmallStr::from("n")];
2963 let text = lenient_scan(&paths, dataset.schema.clone(), None, Some(&drift), &as_text)
2964 .unwrap()
2965 .collect()
2966 .unwrap();
2967 assert_eq!(
2968 text.column("n")
2969 .unwrap()
2970 .str()
2971 .unwrap()
2972 .iter()
2973 .collect::<Vec<_>>(),
2974 [
2975 Some("x"),
2976 Some("y"),
2977 Some("z"),
2978 Some("1970-01-01"),
2979 Some("2147483647 days since 1970-01-01"),
2980 Some("1970-01-01 00:00:00.000"),
2981 Some("-9223372036854775807 ms since 1970-01-01 UTC"),
2982 Some("1970-01-01 01:00:00.000000+01:00"),
2983 Some("-9223372036854775807 us since 1970-01-01 UTC"),
2984 ]
2985 );
2986 }
2987
2988 #[test]
2991 fn reading_as_text_leaves_a_file_without_the_column_alone() {
2992 use polars::prelude::{ParquetWriter, df};
2993
2994 let dir = tempfile::tempdir().unwrap();
2995 let mut paths = Vec::new();
2996 let mut write = |name: &str, mut frame: polars::prelude::DataFrame| {
2997 let path = dir.path().join(name);
2998 let f = std::fs::File::create(&path).unwrap();
2999 ParquetWriter::new(f).finish(&mut frame).unwrap();
3000 paths.push(path.to_string_lossy().to_string());
3001 };
3002 write(
3003 "a.parquet",
3004 df!("id" => &[0i64, 1], "n" => &[10i64, 20]).unwrap(),
3005 );
3006 write("b.parquet", df!("id" => &[2i64]).unwrap());
3008 write("c.parquet", df!("id" => &[3i64], "n" => &["x"]).unwrap());
3009
3010 let footers: Vec<Option<FileFooter>> = vec![
3011 file(&[("id", DataType::Int64), ("n", DataType::Int64)], 2),
3012 file(&[("id", DataType::Int64)], 1),
3013 file(&[("id", DataType::Int64), ("n", DataType::String)], 1),
3014 ];
3015 let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(3));
3016 let drift = ScanDrift::new(&paths, &dataset, &[2, 1, 1]).expect("the files disagree");
3017 let as_text = [PlSmallStr::from("n")];
3018 let text = lenient_scan(&paths, dataset.schema.clone(), None, Some(&drift), &as_text)
3019 .unwrap()
3020 .collect()
3021 .unwrap();
3022 assert_eq!(
3023 text.column("n")
3024 .unwrap()
3025 .str()
3026 .unwrap()
3027 .iter()
3028 .collect::<Vec<_>>(),
3029 [Some("10"), Some("20"), None, Some("x")],
3030 "the file with no `n` has none to show"
3031 );
3032 }
3033
3034 #[test]
3041 fn types_the_cast_agrees_with_are_exactly_the_ones_offered() {
3042 use polars::prelude::*;
3043
3044 let mk = |dtype: DataType| -> Column {
3045 Series::new("x".into(), [1i64, 2])
3046 .cast(&dtype)
3047 .unwrap_or_else(|e| panic!("cannot build a {dtype:?} column: {e}"))
3048 .into()
3049 };
3050 let mut cases: Vec<(DataType, Column)> = vec![
3051 DataType::Int64,
3052 DataType::Float64,
3053 DataType::Boolean,
3054 DataType::Date,
3055 DataType::Time,
3056 DataType::Datetime(TimeUnit::Microseconds, None),
3057 DataType::Duration(TimeUnit::Milliseconds),
3058 DataType::Decimal(10, 2),
3059 DataType::List(Box::new(DataType::Int64)),
3060 ]
3061 .into_iter()
3062 .map(|dtype| (dtype.clone(), mk(dtype)))
3063 .collect();
3064 cases.push((DataType::String, Series::new("x".into(), ["a", "b"]).into()));
3065 cases.push((
3067 DataType::Binary,
3068 Series::new("x".into(), [&[0xffu8, 0xfe][..], &[0x41][..]]).into(),
3069 ));
3070 let plain =
3071 StructChunked::from_series("x".into(), 2, [Series::new("a".into(), [1i64, 2])].iter())
3072 .unwrap()
3073 .into_series();
3074 cases.push((plain.dtype().clone(), plain.into()));
3075 let inners: [Series; 3] = [
3078 Series::new("a".into(), [1i64, 2])
3079 .cast(&DataType::Duration(TimeUnit::Milliseconds))
3080 .unwrap(),
3081 Series::new("a".into(), [1i64, 2])
3082 .cast(&DataType::List(Box::new(DataType::Int64)))
3083 .unwrap(),
3084 Series::new("a".into(), [&[0xffu8, 0xfe][..], &[0x41][..]]),
3086 ];
3087 for inner in inners {
3088 let nested = StructChunked::from_series("x".into(), 2, [inner].iter())
3089 .unwrap()
3090 .into_series();
3091 cases.push((nested.dtype().clone(), nested.into()));
3092 }
3093
3094 for (dtype, column) in cases {
3095 let cast_works = DataFrame::new(2, vec![column])
3096 .unwrap()
3097 .lazy()
3098 .select([col("x").cast(DataType::String)])
3099 .collect()
3100 .is_ok();
3101 assert_eq!(
3102 can_read_as_text(&dtype),
3103 cast_works,
3104 "{dtype:?}: the predicate and the cast must agree"
3105 );
3106 }
3107 }
3108
3109 #[test]
3112 fn a_column_the_cast_refuses_is_read_as_it_was() {
3113 use polars::prelude::{ParquetWriter, df};
3114
3115 let dir = tempfile::tempdir().unwrap();
3116 let mut paths = Vec::new();
3117 let mut write = |name: &str, mut frame: polars::prelude::DataFrame| {
3118 let path = dir.path().join(name);
3119 let f = std::fs::File::create(&path).unwrap();
3120 ParquetWriter::new(f).finish(&mut frame).unwrap();
3121 paths.push(path.to_string_lossy().to_string());
3122 };
3123 write(
3125 "a.parquet",
3126 df!("id" => &[0i64, 1], "n" => &[&[0xffu8, 0xfe][..], &[0x41][..]]).unwrap(),
3127 );
3128 write("b.parquet", df!("id" => &[2i64], "n" => &["x"]).unwrap());
3129
3130 let footers: Vec<Option<FileFooter>> = vec![
3131 file(&[("id", DataType::Int64), ("n", DataType::Binary)], 2),
3132 file(&[("id", DataType::Int64), ("n", DataType::String)], 1),
3133 ];
3134 let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(2));
3135 let drifting = dataset
3136 .columns
3137 .iter()
3138 .find(|column| column.name == "n")
3139 .unwrap();
3140 assert!(
3141 !drifting.can_read_as_text(),
3142 "so the Notes tab never offers it"
3143 );
3144
3145 let drift = ScanDrift::new(&paths, &dataset, &[2, 1]).expect("the files disagree");
3146 let as_text = [PlSmallStr::from("n")];
3147 let frame = lenient_scan(&paths, dataset.schema.clone(), None, Some(&drift), &as_text)
3148 .unwrap()
3149 .collect()
3150 .expect("the read still succeeds, which is the point");
3151 assert_eq!(
3152 frame
3153 .column("id")
3154 .unwrap()
3155 .i64()
3156 .unwrap()
3157 .into_no_null_iter()
3158 .collect::<Vec<_>>(),
3159 [0, 1, 2],
3160 "every file is still read, the agreeing one included"
3161 );
3162 assert_ne!(
3163 frame.column("n").unwrap().dtype(),
3164 &DataType::String,
3165 "and the column is as it was, not half-cast"
3166 );
3167 }
3168
3169 #[test]
3177 fn a_type_only_one_file_holds_can_rule_the_column_out() {
3178 use polars::prelude::{IntoLazy, ParquetWriter, col, df};
3179
3180 let dir = tempfile::tempdir().unwrap();
3181 let mut paths = Vec::new();
3182 let mut write = |name: &str, mut frame: polars::prelude::DataFrame| {
3183 let path = dir.path().join(name);
3184 let f = std::fs::File::create(&path).unwrap();
3185 ParquetWriter::new(f).finish(&mut frame).unwrap();
3186 paths.push(path.to_string_lossy().to_string());
3187 };
3188 write(
3189 "a.parquet",
3190 df!("id" => &[0i64, 1, 2], "n" => &[10i64, 20, 30]).unwrap(),
3191 );
3192 write(
3194 "b.parquet",
3195 df!("id" => &[3i64], "n" => &[9i64])
3196 .unwrap()
3197 .lazy()
3198 .group_by([col("id")])
3199 .agg([col("n")])
3200 .collect()
3201 .unwrap(),
3202 );
3203
3204 let footers: Vec<Option<FileFooter>> = vec![
3205 file(&[("id", DataType::Int64), ("n", DataType::Int64)], 3),
3206 file(
3207 &[
3208 ("id", DataType::Int64),
3209 ("n", DataType::List(Box::new(DataType::Int64))),
3210 ],
3211 1,
3212 ),
3213 ];
3214 let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(2));
3215 assert_eq!(
3216 dataset.schema.get("n"),
3217 Some(&DataType::Int64),
3218 "read as the integer the three rows have"
3219 );
3220 let drifting = dataset
3221 .columns
3222 .iter()
3223 .find(|column| column.name == "n")
3224 .unwrap();
3225 assert!(
3226 can_read_as_text(&drifting.dtype),
3227 "an integer column casts to text on its own account"
3228 );
3229 assert!(
3230 !drifting.can_read_as_text(),
3231 "but one file holds a list, and that file's cast is the one that fails"
3232 );
3233
3234 let drift = ScanDrift::new(&paths, &dataset, &[3, 1]).expect("the files disagree");
3235 let as_text = [PlSmallStr::from("n")];
3236 let frame = lenient_scan(&paths, dataset.schema.clone(), None, Some(&drift), &as_text)
3237 .unwrap()
3238 .collect()
3239 .expect("asking anyway must not cost the read");
3240 assert_eq!(
3241 frame
3242 .column("id")
3243 .unwrap()
3244 .i64()
3245 .unwrap()
3246 .into_no_null_iter()
3247 .collect::<Vec<_>>(),
3248 [0, 1, 2, 3],
3249 "every file is read, the three that agreed included"
3250 );
3251 assert_eq!(
3252 frame.column("n").unwrap().dtype(),
3253 &DataType::Int64,
3254 "and the column is as it was"
3255 );
3256 }
3257
3258 #[test]
3260 fn text_schema_respells_without_reordering() {
3261 let mut schema = Schema::with_capacity(3);
3262 schema.with_column("a".into(), DataType::Int64);
3263 schema.with_column("n".into(), DataType::Int64);
3264 schema.with_column("z".into(), DataType::Float64);
3265 let schema = Arc::new(schema);
3266
3267 let text = text_schema(&schema, &[PlSmallStr::from("n")]);
3268 assert_eq!(
3269 names(&text),
3270 ["a", "n", "z"],
3271 "a column read differently does not move"
3272 );
3273 assert_eq!(text.get("n"), Some(&DataType::String));
3274 assert_eq!(
3275 text.get("a"),
3276 Some(&DataType::Int64),
3277 "nor do its neighbours change"
3278 );
3279 assert_eq!(text.get("z"), Some(&DataType::Float64));
3280
3281 assert!(
3282 Arc::ptr_eq(&schema, &text_schema(&schema, &[])),
3283 "asking for nothing is the schema itself"
3284 );
3285 assert_eq!(
3286 names(&text_schema(&schema, &[PlSmallStr::from("ghost")])),
3287 ["a", "n", "z"],
3288 "a name the schema does not have adds nothing"
3289 );
3290 }
3291
3292 #[test]
3301 fn a_sampled_dataset_counts_what_it_read_and_not_what_it_did_not() {
3302 let footers = vec![
3303 file(&[("id", DataType::Int64)], 0),
3304 file(&[("id", DataType::Int64), ("x", DataType::String)], 5),
3305 None,
3306 ];
3307 let read = [0usize, 250, 499];
3309 let union = union_sampled(500, &read, &footers);
3310
3311 assert_eq!(
3312 union.files, 3,
3313 "the population is the footers read, not the files there are"
3314 );
3315 assert_eq!(union.empty_files, 1, "one of the three held nothing");
3316 assert_eq!(
3317 union.origin,
3318 SchemaOrigin::FooterSample {
3319 read: 3,
3320 total: 500
3321 }
3322 );
3323 assert_eq!(
3324 union.unreadable,
3325 [499],
3326 "and the footer that would not parse is named by its place among the files"
3327 );
3328 assert_eq!(union.file_group.len(), 500);
3331 assert_eq!(union.omitted.len(), 500);
3332 }
3333
3334 #[test]
3340 fn row_groups_are_noted_by_their_middle_size_and_only_when_it_is_large() {
3341 const MIB: usize = 1024 * 1024;
3342 let note = |groups: &[&[usize]]| -> Option<String> {
3343 let files: Vec<Option<FileFooter>> = groups
3344 .iter()
3345 .map(|sizes| {
3346 Some(FileFooter {
3347 schema: Arc::new(Schema::with_capacity(0)),
3348 row_group_rows: vec![1],
3349 file_bytes: 0,
3350 row_group_bytes: sizes.to_vec(),
3351 column_bytes: Vec::new(),
3352 })
3353 })
3354 .collect();
3355 let dataset = union_file_schemas(&files, SchemaOrigin::AllFooters(files.len()));
3356 crate::notes::from_dataset(&dataset)
3357 .into_iter()
3358 .find(|note| note.summary.starts_with("median row group"))
3359 .map(|note| note.summary)
3360 };
3361
3362 assert_eq!(note(&[&[MIB], &[2 * MIB]]), None, "ordinary row groups");
3363 assert_eq!(
3364 note(&[&[64 * MIB]]),
3365 None,
3366 "the threshold itself is not past it"
3367 );
3368 assert_eq!(
3369 note(&[&[65 * MIB]]).as_deref(),
3370 Some("median row group 65.0 MiB, each read whole"),
3371 );
3372 assert_eq!(
3373 note(&[&[MIB, MIB, 4096 * MIB]]),
3374 None,
3375 "one huge row group among small ones does not describe the dataset"
3376 );
3377 assert_eq!(
3378 note(&[&[100 * MIB, 100 * MIB], &[MIB]]).as_deref(),
3379 Some("median row group 100.0 MiB, each read whole"),
3380 "the middle of every row group of every file, not the middle of the files"
3381 );
3382 assert_eq!(
3383 note(&[&[MIB], &[100 * MIB, 100 * MIB]]).as_deref(),
3384 Some("median row group 100.0 MiB, each read whole"),
3385 "including when the large ones are not in the first file"
3386 );
3387 assert_eq!(
3390 note(&[&[100 * MIB], &[MIB], &[100 * MIB]]).as_deref(),
3391 Some("median row group 100.0 MiB, each read whole"),
3392 "and when they arrive out of order"
3393 );
3394 assert_eq!(
3395 note(&[&[MIB], &[100 * MIB], &[MIB]]),
3396 None,
3397 "which cuts both ways: one big group between two small ones is not the middle"
3398 );
3399 assert_eq!(note(&[&[]]), None, "a file with no row groups says nothing");
3400 assert_eq!(
3403 note(&[&[64 * MIB, 65 * MIB]]),
3404 None,
3405 "two row groups either side of the line: the lower one decides"
3406 );
3407 assert_eq!(
3408 note(&[&[65 * MIB, 66 * MIB]]).as_deref(),
3409 Some("median row group 65.0 MiB, each read whole"),
3410 "and when it decides the other way it is still the lower one"
3411 );
3412 }
3413
3414 #[test]
3423 fn a_real_footer_reports_the_compressed_size_of_each_row_group() {
3424 use polars::prelude::{ParquetWriter, df};
3425
3426 let dir = tempfile::tempdir().unwrap();
3427 let rows: Vec<String> = (0..20_000)
3428 .map(|i| format!("{i:0>6}{}", "abcdefghij".repeat(19)))
3429 .collect();
3430 let mut frame = df!("s" => rows).unwrap();
3431 let file = std::fs::File::create(dir.path().join("wide.parquet")).unwrap();
3432 ParquetWriter::new(file)
3433 .with_row_group_size(Some(20_000))
3434 .finish(&mut frame)
3435 .unwrap();
3436
3437 let footer = crate::dataset_files::local_footer(&dir.path().join("wide.parquet"))
3438 .expect("the footer reads");
3439 assert_eq!(footer.rows(), 20_000);
3440 assert_eq!(footer.row_group_bytes.len(), 1, "one row group");
3441
3442 let on_disk = std::fs::metadata(dir.path().join("wide.parquet"))
3446 .unwrap()
3447 .len();
3448 assert_eq!(
3449 footer.file_bytes as u64, on_disk,
3450 "the file's size, as the filesystem reports it"
3451 );
3452
3453 let size = footer.row_group_bytes[0];
3454 assert!(size > 0, "a size is reported");
3455 assert!(
3456 size < 1_000_000,
3457 "and it is the compressed size: 20,000 distinct strings of 200 characters \
3458 are about 4 MiB decoded and a small fraction of that on disk, so {size} \
3459 bytes is the decoded figure"
3460 );
3461 }
3462
3463 #[test]
3465 fn many_files_are_noted_only_when_they_are_also_small() {
3466 const KIB: usize = 1024;
3467 const MIB: usize = 1024 * KIB;
3468 let note = |files: usize, read: usize, sizes: &[usize]| -> Option<String> {
3471 let footers: Vec<Option<FileFooter>> = sizes
3472 .iter()
3473 .cycle()
3474 .take(if sizes.is_empty() { 0 } else { read })
3475 .map(|bytes| {
3476 Some(FileFooter {
3477 schema: Arc::new(Schema::with_capacity(0)),
3478 row_group_rows: vec![1],
3479 file_bytes: *bytes,
3480 row_group_bytes: Vec::new(),
3481 column_bytes: Vec::new(),
3482 })
3483 })
3484 .collect();
3485 let origin = if read == files {
3486 SchemaOrigin::AllFooters(files)
3487 } else {
3488 SchemaOrigin::FooterSample { read, total: files }
3489 };
3490 crate::notes::from_dataset(&union_file_schemas(&footers, origin))
3491 .into_iter()
3492 .find(|note| note.summary.contains("files, median"))
3493 .map(|note| note.summary)
3494 };
3495
3496 assert_eq!(
3497 note(10_000, 10_000, &[40 * KIB]),
3498 None,
3499 "a year of hourly partitions, and more, is an ordinary shape"
3500 );
3501 assert_eq!(
3502 note(10_001, 10_001, &[40 * KIB]).as_deref(),
3503 Some("10,001 files, median 40.0 KiB; every footer read before any row"),
3504 "one more is not"
3505 );
3506 assert_eq!(
3507 note(50_000, 50_000, &[MIB]),
3508 None,
3509 "a megabyte is not small by this measure"
3510 );
3511 assert!(
3512 note(50_000, 50_000, &[MIB - 1]).is_some(),
3513 "a byte under it is"
3514 );
3515 assert_eq!(
3516 note(50_000, 50_000, &[40 * KIB, 40 * KIB, 900 * MIB]).as_deref(),
3517 Some("50,000 files, median 40.0 KiB; every footer read before any row"),
3518 "a large minority does not move the middle"
3519 );
3520 assert_eq!(
3524 note(500_000, 2, &[40 * KIB, 40 * KIB]).as_deref(),
3525 Some("500,000 files, median 40.0 KiB; 2 footers read before any row"),
3526 "the count is the listing's; the footers read are their own number"
3527 );
3528 assert_eq!(note(50_000, 0, &[]), None, "no footer read, nothing to say");
3529
3530 let sampled = union_file_schemas(
3533 &[
3534 Some(FileFooter {
3535 schema: Arc::new(Schema::with_capacity(0)),
3536 row_group_rows: vec![1],
3537 file_bytes: 40 * KIB,
3538 row_group_bytes: Vec::new(),
3539 column_bytes: Vec::new(),
3540 }),
3541 Some(FileFooter {
3542 schema: Arc::new(Schema::with_capacity(0)),
3543 row_group_rows: vec![1],
3544 file_bytes: 40 * KIB,
3545 row_group_bytes: Vec::new(),
3546 column_bytes: Vec::new(),
3547 }),
3548 ],
3549 SchemaOrigin::FooterSample {
3550 read: 2,
3551 total: 500_000,
3552 },
3553 );
3554 let sampled_note = crate::notes::from_dataset(&sampled)
3555 .into_iter()
3556 .find(|note| note.summary.contains("files, median"))
3557 .expect("the note is made");
3558 assert_eq!(sampled_note.scope, "in 2 of 500,000 footers (sample)");
3559 assert_eq!(
3560 note(50_000, 2, &[0, 0]),
3561 None,
3562 "and a size of nothing means the size is not known, not that it is small"
3563 );
3564 }
3565
3566 #[test]
3574 fn partition_keys_are_the_key_equals_segments_above_the_file() {
3575 let keys = |path: &str| partition_keys_of(path);
3576 assert_eq!(keys("data/date=2024-01-01/a.parquet"), ["date"]);
3577 assert_eq!(keys("data/y=2024/m=05/a.parquet"), ["m", "y"]);
3578 assert_eq!(
3579 keys("data/m=05/y=2024/a.parquet"),
3580 keys("data/y=2024/m=05/a.parquet"),
3581 "the same two partitions, written in two orders"
3582 );
3583 assert_eq!(
3584 keys("data/x=1/x=2/a.parquet"),
3585 ["x"],
3586 "and a key repeated down the tree is one key"
3587 );
3588 assert_eq!(keys("data/a.parquet"), Vec::<String>::new());
3589 assert_eq!(
3590 keys("data/x=1/2024=05.parquet"),
3591 ["x"],
3592 "the file's own name is not a partition, whatever it looks like"
3593 );
3594 assert_eq!(
3595 keys("data/=2024/a.parquet"),
3596 Vec::<String>::new(),
3597 "nor is a segment with nothing before the equals"
3598 );
3599 #[cfg(windows)]
3604 assert_eq!(
3605 keys(r"data\date=2024-01-01\a.parquet"),
3606 ["date"],
3607 "a path written the other way round is the same path"
3608 );
3609 #[cfg(not(windows))]
3610 assert_eq!(
3611 keys(r"data/we\ird=1/f.parquet"),
3612 ["we\\ird"],
3613 "a backslash here is part of the name, not a separator"
3614 );
3615 #[cfg(not(windows))]
3616 assert_eq!(
3617 keys(r"data/x=1\y=2/f.parquet"),
3618 ["x"],
3619 "so one directory is one partition, however it is spelled"
3620 );
3621 }
3622
3623 #[test]
3629 fn directories_that_partition_differently_are_counted_each_way() {
3630 let note = |root: &str, paths: &[&str]| -> Option<crate::notes::Note> {
3631 let footers = vec![
3632 Some(FileFooter {
3633 schema: Arc::new(Schema::with_capacity(0)),
3634 row_group_rows: vec![1],
3635 file_bytes: 1,
3636 row_group_bytes: Vec::new(),
3637 column_bytes: Vec::new(),
3638 });
3639 paths.len()
3640 ];
3641 let owned: Vec<String> = paths.iter().map(|p| p.to_string()).collect();
3642 let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(paths.len()))
3643 .with_partition_layouts(root, &owned);
3644 crate::notes::from_dataset(&dataset)
3645 .into_iter()
3646 .find(|note| note.summary.contains("mixed partition keys"))
3647 };
3648
3649 assert_eq!(
3650 note("d", &["d/date=1/a.parquet", "d/date=2/b.parquet"]),
3651 None,
3652 "directories that agree have nothing to say"
3653 );
3654 assert_eq!(
3655 note("d", &["d/a.parquet", "d/b.parquet"]),
3656 None,
3657 "nor has a dataset with no partitions at all"
3658 );
3659 assert_eq!(
3660 note("d", &["d/y=1/m=1/a.parquet", "d/m=2/y=2/b.parquet"]),
3661 None,
3662 "nor two orders of the same two keys: hive matches columns by name, so \
3663 that dataset reads perfectly well and has nothing in dispute"
3664 );
3665 assert_eq!(
3666 note(
3667 "d/run=7",
3668 &["d/run=7/loose.parquet", "d/run=7/date=1/a.parquet"]
3669 ),
3670 None,
3671 "a key=value directory above the dataset as it was opened is not one of the \
3672 things its directories disagree about — and these two files are where that \
3673 matters, since counting `run` would make the one without a key of its \
3674 own a second layout"
3675 );
3676 assert_eq!(
3677 note(
3678 "s3://b//data/",
3679 &["s3://b/data/date=1/a.parquet", "s3://b/data/dt=2/b.parquet"]
3680 ),
3681 None,
3682 "and a path the root is not a prefix of — a typed URL with a doubled \
3683 slash rebuilds without it — is one this cannot place, so it is left out \
3684 rather than read from the top"
3685 );
3686
3687 let renamed = note(
3688 "d",
3689 &[
3690 "d/date=1/a.parquet",
3691 "d/date=2/b.parquet",
3692 "d/date=3/c.parquet",
3693 "d/dt=4/e.parquet",
3694 ],
3695 )
3696 .expect("the directories disagree");
3697 assert_eq!(
3698 renamed.summary,
3699 "mixed partition keys: 3 files by date, \
3700 1 file by dt"
3701 );
3702 assert_eq!(
3703 renamed.scope, "in the names of 4 files",
3704 "read off every name, not off the footers datui opened"
3705 );
3706
3707 let loose = note(
3710 "d",
3711 &[
3712 "d/y=1/m=1/a.parquet",
3713 "d/y=1/m=2/b.parquet",
3714 "d/date=3/c.parquet",
3715 "d/loose.parquet",
3716 ],
3717 )
3718 .expect("the directories disagree");
3719 assert_eq!(
3720 loose.summary,
3721 "mixed partition keys: 2 files by m/y, \
3722 1 file by date"
3723 );
3724 assert_eq!(
3725 loose.scope, "in the names of 4 files",
3726 "the unpartitioned file is one of the names read"
3727 );
3728 }
3729
3730 #[test]
3732 fn the_layouts_a_note_names_are_the_commonest_of_them() {
3733 let layouts = |paths: &[&str]| -> Vec<(Vec<String>, usize)> {
3734 let owned: Vec<String> = paths.iter().map(|p| p.to_string()).collect();
3735 union_file_schemas(&[], SchemaOrigin::AllFooters(0))
3736 .with_partition_layouts("d", &owned)
3737 .partition_layouts
3738 };
3739 let note = |paths: &[&str]| -> String {
3740 let footers = vec![
3741 Some(FileFooter {
3742 schema: Arc::new(Schema::with_capacity(0)),
3743 row_group_rows: vec![1],
3744 file_bytes: 1,
3745 row_group_bytes: Vec::new(),
3746 column_bytes: Vec::new(),
3747 });
3748 paths.len()
3749 ];
3750 let owned: Vec<String> = paths.iter().map(|p| p.to_string()).collect();
3751 let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(paths.len()))
3752 .with_partition_layouts("d", &owned);
3753 crate::notes::from_dataset(&dataset)
3754 .into_iter()
3755 .find(|note| note.summary.contains("mixed partition keys"))
3756 .expect("the directories disagree")
3757 .summary
3758 };
3759
3760 assert_eq!(
3764 layouts(&[
3765 "d/zzz=1/b.parquet",
3766 "d/aaa=1/a.parquet",
3767 "d/zzz=2/c.parquet",
3768 "d/zzz=3/e.parquet",
3769 ]),
3770 vec![(vec!["zzz".to_string()], 3), (vec!["aaa".to_string()], 1)],
3771 "commonest first, though the rare one sorts first and arrived first"
3772 );
3773 assert_eq!(
3774 layouts(&[
3775 "d/zz=1/a.parquet",
3776 "d/aa=1/b.parquet",
3777 "d/mm=1/c.parquet",
3778 "d/qq=1/e.parquet"
3779 ]),
3780 vec![
3781 (vec!["aa".to_string()], 1),
3782 (vec!["mm".to_string()], 1),
3783 (vec!["qq".to_string()], 1),
3784 (vec!["zz".to_string()], 1)
3785 ],
3786 "and equally common ones by their keys, so the same dataset reads the \
3787 same way every time it is opened"
3788 );
3789
3790 assert_eq!(
3791 note(&[
3792 "d/aa=1/a.parquet",
3793 "d/bb=1/b.parquet",
3794 "d/cc=1/c.parquet",
3795 "d/dd=1/e.parquet",
3796 ]),
3797 "mixed partition keys: 1 file by aa, \
3798 1 file by bb, 2 files by 2 other ways"
3799 );
3800 assert_eq!(
3801 note(&["d/aa=1/a.parquet", "d/bb=1/b.parquet", "d/cc=1/c.parquet"]),
3802 "mixed partition keys: 1 file by aa, \
3803 1 file by bb, 1 file by 1 other way",
3804 "and one of them is one way, not one ways"
3805 );
3806
3807 let many: Vec<String> = (0..100)
3811 .map(|i| format!("d/k{i:0>3}=1/f.parquet"))
3812 .collect();
3813 let many: Vec<&str> = many.iter().map(String::as_str).collect();
3814 assert_eq!(
3815 note(&many),
3816 "mixed partition keys: 1 file by k000, \
3817 1 file by k001, 98 files by 98 other ways"
3818 );
3819 let owned: Vec<String> = many.iter().map(|p| p.to_string()).collect();
3820 let dataset = union_file_schemas(&[], SchemaOrigin::AllFooters(0))
3821 .with_partition_layouts("d", &owned);
3822 assert!(
3823 dataset.partition_layouts.len() <= 64,
3824 "and it is not holding a hundred of them to say so: {}",
3825 dataset.partition_layouts.len()
3826 );
3827 }
3828
3829 #[test]
3834 fn the_footer_count_speaks_only_while_a_pass_is_running() {
3835 let progress = FooterProgress::default();
3836 assert_eq!(progress.reading(), None, "nothing has begun");
3837
3838 progress.begin(3);
3839 assert_eq!(progress.reading(), Some((0, 3)), "none read yet");
3840 progress.advance();
3841 progress.advance();
3842 assert_eq!(progress.reading(), Some((2, 3)));
3843
3844 progress.done();
3845 assert_eq!(progress.reading(), None, "and nothing once it has landed");
3846
3847 progress.begin(2);
3849 assert_eq!(progress.reading(), Some((0, 2)));
3850 }
3851
3852 #[test]
3862 fn the_footer_count_never_passes_its_total() {
3863 let progress = FooterProgress::default();
3864 progress.begin(2);
3865 for _ in 0..5 {
3866 progress.advance();
3867 }
3868 assert_eq!(progress.reading(), Some((2, 2)));
3869 }
3870
3871 #[test]
3877 fn a_pass_that_panics_still_says_it_has_finished() {
3878 let progress = FooterProgress::default();
3879 let caught = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
3880 let pass = progress.pass(3);
3881 pass.advance();
3882 panic!("a footer reader gave up");
3883 }));
3884 assert!(caught.is_err(), "the panic happened");
3885 assert_eq!(
3886 progress.reading(),
3887 None,
3888 "and the count went with it rather than sitting there"
3889 );
3890 assert_eq!(
3891 progress.last_pass().read,
3892 1,
3893 "with what it managed still readable"
3894 );
3895 }
3896
3897 #[test]
3904 fn absent_columns_alone_never_split_the_scan() {
3905 let files = 64;
3906 let per_file: Vec<Option<FileFooter>> = (0..files)
3907 .map(|i| {
3908 let mut s = Schema::with_capacity(2);
3909 s.with_column("id".into(), DataType::Int64);
3910 if i % 2 == 1 {
3911 s.with_column("extra".into(), DataType::String);
3912 }
3913 Some(FileFooter {
3914 schema: Arc::new(s),
3915 row_group_rows: vec![1],
3916 file_bytes: 0,
3917 row_group_bytes: Vec::new(),
3918 column_bytes: Vec::new(),
3919 })
3920 })
3921 .collect();
3922 let paths: Vec<String> = (0..files).map(|i| format!("part-{i:05}.parquet")).collect();
3923 let read: Vec<usize> = (0..files).collect();
3924 let dataset = union_sampled(files, &read, &per_file);
3925 let rows = vec![1usize; files];
3926 let drift = ScanDrift::new(&paths, &dataset, &rows).expect("this dataset drifts");
3927 assert!(dataset.drifts());
3928 assert_eq!(
3929 runs_of(&paths, &drift),
3930 1,
3931 "absent columns need no split, however they alternate"
3932 );
3933
3934 let mut with_conflict = per_file.clone();
3936 let mut odd = Schema::with_capacity(2);
3937 odd.with_column("id".into(), DataType::String);
3938 with_conflict[7] = Some(FileFooter {
3939 schema: Arc::new(odd),
3940 row_group_rows: vec![1],
3941 file_bytes: 0,
3942 row_group_bytes: Vec::new(),
3943 column_bytes: Vec::new(),
3944 });
3945 let dataset = union_sampled(files, &read, &with_conflict);
3946 let drift = ScanDrift::new(&paths, &dataset, &rows).unwrap();
3947 assert_eq!(runs_of(&paths, &drift), 3, "before it, it, and after it");
3948 }
3949
3950 fn runs_of(paths: &[String], drift: &ScanDrift) -> usize {
3952 let mut runs = 1;
3953 for pair in paths.windows(2) {
3954 if drift.unread(&pair[0]) != drift.unread(&pair[1]) {
3955 runs += 1;
3956 }
3957 }
3958 runs
3959 }
3960
3961 #[test]
3963 fn a_reader_puts_part_2_before_part_10() {
3964 use std::cmp::Ordering;
3965 let cmp = |a: &str, b: &str| natural_cmp(a, b);
3966 assert_eq!(
3967 cmp("part=2", "part=10"),
3968 Ordering::Less,
3969 "which bytes do not"
3970 );
3971 assert_eq!(cmp("part=10", "part=2"), Ordering::Greater);
3972 assert_eq!(cmp("date=2024-01-02", "date=2024-01-03"), Ordering::Less);
3973 assert_eq!(cmp("date=2024-01-02", "date=2024-01-02"), Ordering::Equal);
3974 assert_eq!(
3975 cmp("m=03", "m=3"),
3976 Ordering::Equal,
3977 "the same number written two ways is neither before nor after itself — a \
3978 dataset that spells one month both ways is past helping, and this at \
3979 least does not invent an order for it"
3980 );
3981 assert_eq!(cmp("a=1/b=2", "a=1/b=10"), Ordering::Less);
3982 assert_eq!(
3983 cmp("x=a", "x=b"),
3984 Ordering::Less,
3985 "and letters are still letters"
3986 );
3987 }
3988
3989 #[test]
3991 fn a_partition_holds_the_ones_below_it() {
3992 assert!(partition_holds("y=2024", "y=2024/m=03"));
3993 assert!(partition_holds("y=2024", "y=2024"));
3994 assert!(!partition_holds("y=2024", "y=2025"));
3995 assert!(
3996 !partition_holds("y=202", "y=2024"),
3997 "a prefix of the spelling is not a directory above it"
3998 );
3999 assert!(!partition_holds("y=2024/m=03", "y=2024"));
4000 }
4001
4002 #[test]
4009 fn a_file_name_is_not_a_partition_and_neither_is_a_bare_segment() {
4010 assert_eq!(partition_values_of("/x=1/f.parquet"), ["x=1"]);
4011 assert_eq!(
4012 partition_values_of("/x=1/2024=05.parquet"),
4013 ["x=1"],
4014 "the file's own name is never a partition, whatever it is called"
4015 );
4016 assert_eq!(
4017 partition_values_of("/raw/x=1/f.parquet"),
4018 ["x=1"],
4019 "and a segment with no key before the `=` is not one either"
4020 );
4021 assert_eq!(
4022 partition_values_of("/=1/f.parquet"),
4023 Vec::<String>::new(),
4024 "an empty key is no key"
4025 );
4026 assert_eq!(
4027 partition_values_of("/y=2024/m=03/f.parquet"),
4028 ["y=2024", "m=03"],
4029 "in the order written, because a partition is a place"
4030 );
4031 assert_eq!(
4032 partition_values_of("/date=2024=05/f.parquet"),
4033 ["date=2024=05"],
4034 "and a value may hold an `=` of its own"
4035 );
4036 }
4037
4038 #[test]
4039 fn a_column_only_a_middle_file_has_is_kept() {
4040 let files = [
4041 file(&[("id", DataType::Int64)], 10),
4042 file(&[("id", DataType::Int64), ("oops", DataType::String)], 10),
4043 file(&[("id", DataType::Int64)], 10),
4044 ];
4045 let union = union(&files);
4046 assert_eq!(names(&union.schema), ["id", "oops"]);
4047 let oops = union.columns.iter().find(|c| c.name == "oops").unwrap();
4048 assert_eq!(oops.present_in, 1);
4049 }
4050
4051 #[test]
4052 fn the_newest_files_order_leads_and_older_columns_follow() {
4053 let files = [
4054 file(&[("a", DataType::Int64), ("gone", DataType::Int64)], 1),
4055 file(&[("b", DataType::Int64), ("a", DataType::Int64)], 1),
4056 ];
4057 assert_eq!(names(&union(&files).schema), ["b", "a", "gone"]);
4058 }
4059
4060 #[test]
4061 fn integer_widths_widen_losslessly() {
4062 let files = [
4063 file(&[("n", DataType::Int32)], 100),
4064 file(&[("n", DataType::Int64)], 1),
4065 ];
4066 let union = union(&files);
4067 assert_eq!(union.schema.get("n"), Some(&DataType::Int64));
4068 assert!(union.columns[0].widened);
4069 assert_eq!(union.columns[0].conflicting_files, 0);
4070 assert!(union.omitted.iter().all(|o| o.is_empty()));
4071 }
4072
4073 #[test]
4074 fn an_integer_and_a_float_meet_at_float64() {
4075 let files = [
4076 file(&[("n", DataType::Int32)], 1),
4077 file(&[("n", DataType::Float32)], 1),
4078 ];
4079 assert_eq!(union(&files).schema.get("n"), Some(&DataType::Float64));
4080 }
4081
4082 #[test]
4083 fn datetime_units_widen_to_the_finer_one() {
4084 let ms = DataType::Datetime(TimeUnit::Milliseconds, None);
4085 let ns = DataType::Datetime(TimeUnit::Nanoseconds, None);
4086 let files = [file(&[("t", ms)], 1), file(&[("t", ns.clone())], 1)];
4087 assert_eq!(union(&files).schema.get("t"), Some(&ns));
4088 }
4089
4090 #[test]
4091 fn a_struct_has_every_field_either_file_has() {
4092 let old = DataType::Struct(vec![Field::new("a".into(), DataType::Int32)]);
4093 let new = DataType::Struct(vec![
4094 Field::new("a".into(), DataType::Int64),
4095 Field::new("b".into(), DataType::String),
4096 ]);
4097 let files = [file(&[("s", old)], 1), file(&[("s", new.clone())], 1)];
4098 assert_eq!(union(&files).schema.get("s"), Some(&new));
4099 }
4100
4101 #[test]
4102 fn a_type_conflict_goes_to_the_majority_of_rows() {
4103 let files = [
4104 file(&[("price", DataType::String)], 10),
4105 file(&[("price", DataType::Int64)], 90),
4106 ];
4107 let union = union(&files);
4108 assert_eq!(union.schema.get("price"), Some(&DataType::Int64));
4109 assert_eq!(union.columns[0].conflicting_files, 1);
4110 assert_eq!(union.columns[0].conflicting_types, [DataType::String]);
4111 assert_eq!(omitted_names(&union, 0), ["price"]);
4112 assert!(union.omitted[1].is_empty());
4113 }
4114
4115 #[test]
4116 fn the_majority_can_be_the_text_files() {
4117 let files = [
4118 file(&[("price", DataType::String)], 90),
4119 file(&[("price", DataType::Int64)], 10),
4120 ];
4121 let union = union(&files);
4122 assert_eq!(union.schema.get("price"), Some(&DataType::String));
4123 assert_eq!(omitted_names(&union, 1), ["price"]);
4124 }
4125
4126 #[test]
4127 fn a_type_that_covers_more_files_wins_over_one_that_covers_none_extra() {
4128 let files = [
4130 file(&[("n", DataType::Int32)], 30),
4131 file(&[("n", DataType::Float32)], 30),
4132 file(&[("n", DataType::String)], 50),
4133 ];
4134 let union = union(&files);
4135 assert_eq!(union.schema.get("n"), Some(&DataType::Float64));
4136 assert_eq!(omitted_names(&union, 2), ["n"]);
4137 }
4138
4139 #[test]
4140 fn names_differing_only_by_case_stay_two_columns() {
4141 let files = [file(
4142 &[("Price", DataType::Int64), ("price", DataType::Int64)],
4143 1,
4144 )];
4145 assert_eq!(names(&union(&files).schema), ["Price", "price"]);
4146 }
4147
4148 #[test]
4149 fn an_unreadable_footer_is_recorded_and_left_out() {
4150 let files = [
4151 file(&[("id", DataType::Int64)], 1),
4152 None,
4153 file(&[("id", DataType::Int64), ("late", DataType::Int64)], 1),
4154 ];
4155 let union = union(&files);
4156 assert_eq!(union.unreadable, [1]);
4157 assert_eq!(names(&union.schema), ["id", "late"]);
4158 assert!(union.omitted[1].is_empty());
4159 }
4160
4161 #[test]
4162 fn a_column_of_nulls_takes_the_other_files_type() {
4163 let files = [
4164 file(&[("x", DataType::Null)], 1),
4165 file(&[("x", DataType::Int64)], 1),
4166 ];
4167 let union = union(&files);
4168 assert_eq!(union.schema.get("x"), Some(&DataType::Int64));
4169 assert_eq!(union.columns[0].conflicting_files, 0);
4170 }
4171
4172 #[test]
4173 fn unsigned_and_signed_meet_in_a_wider_signed_type() {
4174 assert_eq!(
4175 widen(&DataType::UInt32, &DataType::Int32),
4176 Some(DataType::Int64)
4177 );
4178 assert_eq!(widen(&DataType::UInt64, &DataType::Int64), None);
4179 }
4180
4181 #[test]
4182 fn origins_read_as_sentences() {
4183 assert_eq!(
4184 SchemaOrigin::AllFooters(6541).to_string(),
4185 "all 6,541 footers"
4186 );
4187 assert_eq!(
4188 SchemaOrigin::FooterSample {
4189 read: 5000,
4190 total: 200_000
4191 }
4192 .to_string(),
4193 "5,000 of 200,000 footers (sample)"
4194 );
4195 }
4196
4197 #[test]
4198 fn a_sample_spans_the_files_and_keeps_the_first_and_newest() {
4199 assert_eq!(footers_to_read(3), [0, 1, 2]);
4200 assert_eq!(footers_to_read(MAX_FOOTER_READS).len(), MAX_FOOTER_READS);
4201 let sample = footers_to_read(MAX_FOOTER_READS * 10);
4202 assert_eq!(sample.len(), MAX_FOOTER_READS);
4203 assert_eq!(sample.first(), Some(&0));
4204 assert_eq!(sample.last(), Some(&(MAX_FOOTER_READS * 10 - 1)));
4205 assert!(sample.windows(2).all(|w| w[0] < w[1]), "ascending");
4206 }
4207
4208 #[test]
4212 fn types_the_scan_cannot_cast_are_not_widened() {
4213 let ms = DataType::Duration(TimeUnit::Milliseconds);
4214 let us = DataType::Duration(TimeUnit::Microseconds);
4215 assert_eq!(widen(&ms, &us), None);
4216 assert_eq!(widen(&DataType::Binary, &DataType::String), None);
4217 assert_eq!(widen(&DataType::Date, &ms), None);
4218 }
4219
4220 #[test]
4221 fn drifting_counts_against_the_files_read_not_the_busiest_column() {
4222 let files = [
4224 file(&[("a", DataType::Int64)], 1),
4225 file(&[("b", DataType::Int64)], 1),
4226 ];
4227 let union = union(&files);
4228 let drifting: Vec<_> = union.drifting().map(|c| c.name.to_string()).collect();
4229 assert_eq!(drifting, ["b", "a"]);
4230 }
4231
4232 #[test]
4233 fn drifting_names_only_the_columns_worth_a_note() {
4234 let files = [
4235 file(&[("id", DataType::Int64)], 1),
4236 file(&[("id", DataType::Int64), ("oops", DataType::String)], 1),
4237 ];
4238 let union = union(&files);
4239 let drifting: Vec<_> = union.drifting().map(|c| c.name.to_string()).collect();
4240 assert_eq!(drifting, ["oops"]);
4241 }
4242}