1use std::borrow::Cow;
14use std::collections::{HashMap, HashSet};
15use std::sync::Arc;
16use std::sync::atomic::{AtomicUsize, Ordering};
17
18use polars::chunked_array::cast::CastOptions;
19use polars::prelude::{
20 DataType, Field, LazyFrame, PlRefPath, PlSmallStr, PolarsResult, Schema, TimeUnit, UnionArgs,
21 concat,
22};
23
24pub struct Pass<'a>(&'a FooterProgress);
26
27impl Pass<'_> {
28 pub fn advance(&self) {
30 self.0.advance();
31 }
32}
33
34impl Drop for Pass<'_> {
35 fn drop(&mut self) {
36 self.0.done();
37 }
38}
39
40pub struct Listing<'a>(&'a FooterProgress);
43
44impl Listing<'_> {
45 pub fn advance(&self) {
47 self.0.listed.fetch_add(1, Ordering::Relaxed);
48 }
49
50 pub fn counter(&self) -> std::sync::Arc<AtomicUsize> {
52 self.0.listed.clone()
53 }
54
55 pub fn add(&self, n: usize) {
57 self.0.listed.fetch_add(n, Ordering::Relaxed);
58 }
59
60 pub fn is_cancelled(&self) -> bool {
62 self.0.is_cancelled()
63 }
64
65 pub fn cancel_flag(&self) -> std::sync::Arc<std::sync::atomic::AtomicBool> {
67 self.0.cancel_flag()
68 }
69}
70
71impl Drop for Listing<'_> {
72 fn drop(&mut self) {
73 self.0.listing.store(false, Ordering::Release);
74 }
75}
76
77#[doc(hidden)]
80#[derive(Debug, Clone, Copy, PartialEq, Eq)]
81pub struct PassCount {
82 pub begun: usize,
84 pub read: usize,
86 pub total: usize,
88}
89
90#[derive(Debug, Default)]
94pub struct FooterProgress {
95 read: AtomicUsize,
96 total: AtomicUsize,
97 passes: AtomicUsize,
100 last_total: AtomicUsize,
102 cancelled: std::sync::Arc<std::sync::atomic::AtomicBool>,
105 listed: std::sync::Arc<AtomicUsize>,
108 listing: std::sync::atomic::AtomicBool,
109 at_once: AtomicUsize,
111 estimate: std::sync::Mutex<Option<RowEstimate>>,
113}
114
115#[derive(Debug, Clone, Copy, PartialEq, Eq)]
118pub struct RowEstimate {
119 pub rows: u64,
120 pub sampled: usize,
122 pub files: usize,
124}
125
126impl RowEstimate {
127 pub fn of<'a>(
129 files: usize,
130 footers: impl IntoIterator<Item = &'a Option<FileFooter>>,
131 ) -> Option<Self> {
132 let (sampled, rows) = footers
133 .into_iter()
134 .flatten()
135 .fold((0usize, 0u128), |(n, rows), f| {
136 (n + 1, rows + f.rows() as u128)
137 });
138 (sampled > 0).then(|| RowEstimate {
139 rows: u64::try_from(rows * files as u128 / sampled as u128).unwrap_or(u64::MAX),
140 sampled,
141 files,
142 })
143 }
144}
145
146pub const ESTIMATE_SAMPLE: usize = 2_000;
149
150pub const COUNT_AT_ONCE: usize = 256;
152
153pub fn random_sample(files: usize, n: usize, seed: u64) -> Vec<usize> {
156 if files <= n {
157 return (0..files).collect();
158 }
159 let mut state = seed;
161 let mut next = move || {
162 state = state.wrapping_add(0x9e37_79b9_7f4a_7c15);
163 let mut z = state;
164 z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
165 z = (z ^ (z >> 27)).wrapping_mul(0x94d0_49bb_1331_11eb);
166 z ^ (z >> 31)
167 };
168 let mut chosen = std::collections::BTreeSet::new();
169 while chosen.len() < n {
170 chosen.insert((next() % files as u64) as usize);
171 }
172 chosen.into_iter().collect()
173}
174
175impl FooterProgress {
176 pub fn listing(&self) -> Listing<'_> {
178 self.listed.store(0, Ordering::Relaxed);
179 self.listing.store(true, Ordering::Release);
180 Listing(self)
181 }
182
183 pub fn listed(&self) -> Option<usize> {
185 self.listing
186 .load(Ordering::Acquire)
187 .then(|| self.listed.load(Ordering::Relaxed))
188 }
189
190 pub fn begin(&self, total: usize) {
192 self.read.store(0, Ordering::Relaxed);
193 self.last_total.store(total, Ordering::Relaxed);
196 self.total.store(total, Ordering::Release);
197 self.passes.fetch_add(1, Ordering::Relaxed);
198 }
199
200 pub fn advance(&self) {
202 self.read.fetch_add(1, Ordering::Relaxed);
203 }
204
205 pub fn done(&self) {
207 self.total.store(0, Ordering::Relaxed);
208 }
209
210 pub fn pass(&self, total: usize) -> Pass<'_> {
213 self.begin(total);
214 Pass(self)
215 }
216
217 #[doc(hidden)]
221 pub fn last_pass(&self) -> PassCount {
222 PassCount {
223 begun: self.passes.load(Ordering::Relaxed),
224 read: self.read.load(Ordering::Relaxed),
225 total: self.last_total.load(Ordering::Relaxed),
226 }
227 }
228
229 pub fn reading(&self) -> Option<(usize, usize)> {
231 let total = self.total.load(Ordering::Acquire);
232 (total > 0).then(|| (self.read.load(Ordering::Relaxed).min(total), total))
233 }
234
235 pub fn cancel(&self) {
238 self.cancelled.store(true, Ordering::Relaxed);
239 }
240
241 pub fn is_cancelled(&self) -> bool {
242 self.cancelled.load(Ordering::Relaxed)
243 }
244
245 pub fn cancel_flag(&self) -> std::sync::Arc<std::sync::atomic::AtomicBool> {
247 self.cancelled.clone()
248 }
249
250 pub fn counting() -> Self {
252 let progress = Self::default();
253 progress.at_once.store(COUNT_AT_ONCE, Ordering::Relaxed);
254 progress
255 }
256
257 pub fn reads_at_once(&self) -> usize {
259 match self.at_once.load(Ordering::Relaxed) {
260 0 => FOOTERS_AT_ONCE,
261 n => n,
262 }
263 }
264
265 pub fn set_estimate(&self, estimate: Option<RowEstimate>) {
267 *self.estimate.lock().unwrap_or_else(|e| e.into_inner()) = estimate;
268 }
269
270 pub fn estimate(&self) -> Option<RowEstimate> {
271 *self.estimate.lock().unwrap_or_else(|e| e.into_inner())
272 }
273}
274
275#[derive(Debug, Clone)]
278pub struct FileFooter {
279 pub schema: Arc<Schema>,
280 pub row_group_rows: Vec<usize>,
282 pub row_group_bytes: Vec<usize>,
285 pub file_bytes: usize,
287 pub column_bytes: Vec<(String, usize)>,
290}
291
292impl FileFooter {
293 pub fn from_metadata(
295 schema: Schema,
296 metadata: &polars_parquet::parquet::metadata::FileMetadata,
297 file_bytes: usize,
298 widths: bool,
299 ) -> Self {
300 let column_bytes = if widths {
301 parquet_column_bytes(&schema, metadata)
302 } else {
303 Vec::new()
304 };
305 FileFooter {
306 schema: Arc::new(schema),
307 row_group_rows: metadata.row_groups.iter().map(|rg| rg.num_rows()).collect(),
308 row_group_bytes: metadata
309 .row_groups
310 .iter()
311 .map(|rg| rg.compressed_size())
312 .collect(),
313 file_bytes,
314 column_bytes,
315 }
316 }
317
318 pub fn from_tail(tail: &[u8], file_bytes: usize, widths: bool) -> color_eyre::Result<Self> {
321 use polars::prelude::{ParquetReader, SchemaExt, SerReader};
322 let mut cursor = std::io::Cursor::new(tail);
323 let mut reader = ParquetReader::new(&mut cursor);
324 let arrow_schema = reader
325 .schema()
326 .map_err(|e| color_eyre::eyre::eyre!("Parquet schema read failed: {e}"))?;
327 let metadata = reader
328 .get_metadata()
329 .map_err(|e| color_eyre::eyre::eyre!("Parquet footer read failed: {e}"))?;
330 Ok(Self::from_metadata(
331 Schema::from_arrow_schema(arrow_schema.as_ref()),
332 metadata,
333 file_bytes,
334 widths,
335 ))
336 }
337
338 pub fn rows(&self) -> usize {
340 self.row_group_rows.iter().sum()
341 }
342}
343
344pub fn parquet_column_bytes(
347 schema: &Schema,
348 metadata: &polars_parquet::parquet::metadata::FileMetadata,
349) -> Vec<(String, usize)> {
350 schema
351 .iter_names()
352 .map(|name| {
353 let bytes: i64 = metadata
354 .row_groups
355 .iter()
356 .flat_map(|rg| rg.columns_under_root_iter(name).into_iter().flatten())
357 .map(|chunk| chunk.uncompressed_size())
358 .sum();
359 (name.to_string(), bytes.max(0) as usize)
360 })
361 .collect()
362}
363
364pub fn column_bytes_per_row(footers: &[Option<FileFooter>]) -> Vec<(String, usize)> {
367 let rows: usize = footers.iter().flatten().map(FileFooter::rows).sum();
368 if rows == 0 {
369 return Vec::new();
370 }
371 let mut totals: Vec<(String, usize)> = Vec::new();
372 let mut at: HashMap<String, usize> = HashMap::new();
373 for (name, bytes) in footers.iter().flatten().flat_map(|f| &f.column_bytes) {
374 match at.get(name) {
375 Some(&i) => totals[i].1 += bytes,
376 None => {
377 at.insert(name.clone(), totals.len());
378 totals.push((name.clone(), *bytes));
379 }
380 }
381 }
382 totals
383 .into_iter()
384 .map(|(name, bytes)| (name, bytes / rows))
385 .collect()
386}
387
388#[derive(Debug, Clone, PartialEq, Eq)]
391pub enum SchemaOrigin {
392 AllFooters(usize),
394 FooterSample { read: usize, total: usize },
396}
397
398impl SchemaOrigin {
399 pub fn total_files(&self) -> usize {
402 match self {
403 SchemaOrigin::AllFooters(files) => *files,
404 SchemaOrigin::FooterSample { total, .. } => *total,
405 }
406 }
407}
408
409impl std::fmt::Display for SchemaOrigin {
410 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
411 match self {
412 SchemaOrigin::AllFooters(1) => write!(f, "one footer"),
413 SchemaOrigin::AllFooters(n) => {
414 write!(f, "all {} footers", crate::numfmt::group_chrome(*n))
415 }
416 SchemaOrigin::FooterSample { read, total } => write!(
417 f,
418 "{} of {} footers (sample)",
419 crate::numfmt::group_chrome(*read),
420 crate::numfmt::group_chrome(*total)
421 ),
422 }
423 }
424}
425
426#[derive(Debug, Clone, PartialEq)]
428pub struct ColumnDrift {
429 pub name: PlSmallStr,
430 pub dtype: DataType,
432 pub present_in: usize,
434 pub conflicting_files: usize,
436 pub conflicting_types: Vec<DataType>,
438 pub widened: bool,
440}
441
442impl ColumnDrift {
443 pub fn is_uniform(&self, files: usize) -> bool {
445 self.present_in == files && self.conflicting_files == 0 && !self.widened
446 }
447}
448
449#[derive(Debug, Clone, PartialEq, Eq)]
452pub enum ColumnRange {
453 Only(String),
455 NoneBefore(String),
457}
458
459#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
465pub struct SkippedFiles {
466 pub bookkeeping: usize,
468 pub not_parquet: usize,
470 pub empty: usize,
473}
474
475impl SkippedFiles {
476 pub fn count(&mut self, bookkeeping: bool) {
479 if bookkeeping {
480 self.bookkeeping += 1;
481 } else {
482 self.not_parquet += 1;
483 }
484 }
485}
486
487#[derive(Debug, Clone)]
492pub struct ReadAs {
493 pub delimiter: Option<u8>,
494 pub has_header: Option<bool>,
495 pub skip_rows: Option<usize>,
496 pub skip_lines: Option<usize>,
497 pub infer_schema_length: Option<usize>,
498 pub ignore_errors: bool,
499 pub try_parse_dates: bool,
500 pub comment_char: Option<String>,
501 pub header_rows: Vec<usize>,
502 pub header_join: String,
503}
504
505impl ReadAs {
506 fn open_options(&self, format: crate::FileFormat) -> crate::OpenOptions {
509 crate::OpenOptions {
510 delimiter: self.delimiter.or(format.separator()),
511 has_header: self.has_header,
512 skip_rows: self.skip_rows,
513 skip_lines: self.skip_lines,
514 infer_schema_length: self.infer_schema_length,
515 ignore_errors: self.ignore_errors,
516 parse_dates: self.try_parse_dates,
517 parse_strings: None,
518 comment_char: self.comment_char.clone(),
519 header_rows: self.header_rows.clone(),
520 header_join: self.header_join.clone(),
521 ..crate::OpenOptions::default()
522 }
523 }
524}
525
526impl Default for ReadAs {
527 fn default() -> Self {
531 Self {
532 delimiter: None,
533 has_header: None,
534 skip_rows: None,
535 skip_lines: None,
536 infer_schema_length: None,
537 ignore_errors: false,
538 try_parse_dates: true,
539 comment_char: None,
540 header_rows: Vec::new(),
541 header_join: crate::formats::csv_dialect::DEFAULT_HEADER_JOIN.to_string(),
542 }
543 }
544}
545
546pub fn column_schema_of(
552 path: &std::path::Path,
553 format: crate::FileFormat,
554 as_read: &ReadAs,
555) -> Option<Vec<(String, DataType)>> {
556 use polars::prelude::{LazyFileListReader, LazyJsonLineReader};
557 let lf = match format.descriptor().lines {
558 Some(crate::cli::Lines::Delimited(_)) => {
559 let options = as_read.open_options(format);
561 let header = crate::formats::csv_dialect::head(path, &options, None, false)
562 .ok()?
563 .names;
564 let reader = crate::formats::readers::csv::configure_csv_reader(
565 crate::formats::readers::csv::csv_reader_of(path).ok()?,
566 &options,
567 None,
568 );
569 crate::formats::csv_dialect::name_columns(reader.finish().ok()?, header.as_deref())
570 .ok()?
571 }
572 Some(crate::cli::Lines::Json) => {
573 LazyJsonLineReader::new(crate::cloud::source::polars_literal_path(path).ok()?)
574 .finish()
575 .ok()?
576 }
577 Some(crate::cli::Lines::Text) => {
579 return Some(
580 crate::formats::lines::schema(false)
581 .iter()
582 .map(|(name, dtype)| (name.to_string(), dtype.clone()))
583 .collect(),
584 );
585 }
586 None => return None,
587 };
588 let schema = lf.clone().collect_schema().ok()?;
589 let fields: Vec<(String, DataType)> = schema
592 .iter()
593 .map(|(name, dtype)| (name.trim().to_string(), dtype.clone()))
594 .collect();
595 Some(fields)
596}
597
598fn is_an_empty_file(schema: &[(String, DataType)]) -> bool {
603 match schema {
604 [] => true,
605 [(only, _)] => only.trim().is_empty(),
606 _ => false,
607 }
608}
609
610pub(crate) fn names_are_names(names: &[String]) -> bool {
625 !names.is_empty() && !names.iter().all(|n| n.trim().parse::<f64>().is_ok())
626}
627
628#[derive(Debug, Clone, Default)]
635pub struct Sampled {
636 pub columns: Vec<String>,
639 pub nests: Option<bool>,
642 pub columns_differ: bool,
645 pub types_differ: bool,
648 pub read: usize,
650 pub headerless: bool,
654}
655
656impl Sampled {
657 pub fn disagreement(&self) -> Disagreement {
659 if self.headerless {
661 return Disagreement {
662 headerless: true,
663 ..Default::default()
664 };
665 }
666 Disagreement {
667 columns: self.columns_differ,
668 types: self.types_differ,
669 headerless: false,
670 }
671 }
672}
673
674#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
678pub struct Disagreement {
679 pub columns: bool,
680 pub types: bool,
681 pub headerless: bool,
683}
684
685impl Disagreement {}
686
687pub fn sample_files(
691 files: &[std::path::PathBuf],
692 format: crate::FileFormat,
693 as_read: &ReadAs,
694) -> Sampled {
695 const WANTED: usize = 3;
698 const TRIES: usize = 12;
699 let last = files.len().saturating_sub(1);
700 const NEAR: usize = 4;
704 let anchors = [0usize, last / 2, last];
705
706 let mut read: Vec<Vec<(String, DataType)>> = Vec::new();
707 let mut tried = 0usize;
708 let mut seen: Vec<usize> = Vec::new();
709 'anchors: for anchor in anchors {
710 for step in 0..NEAR {
711 if read.len() >= WANTED || tried >= TRIES {
712 break 'anchors;
713 }
714 let i = anchor + step;
715 if i > last || seen.contains(&i) {
716 continue;
717 }
718 seen.push(i);
719 let Some(file) = files.get(i) else { continue };
720 tried += 1;
721 if let Some(schema) = column_schema_of(file, format, as_read)
723 && !is_an_empty_file(&schema)
724 {
725 read.push(schema);
726 continue 'anchors;
727 }
728 }
729 }
730
731 let mut out = Sampled {
732 read: read.len(),
733 ..Default::default()
734 };
735 for file in &read {
736 for (name, _) in file {
737 if !out.columns.iter().any(|c| c == name) {
738 out.columns.push(name.clone());
739 }
740 }
741 }
742 if read.len() < 2 {
743 return out;
744 }
745 let names: Vec<Vec<String>> = read
746 .iter()
747 .map(|f| f.iter().map(|(n, _)| n.clone()).collect())
748 .collect();
749 if names.iter().any(|f| !names_are_names(f)) {
753 out.columns.clear();
754 out.headerless = true;
755 out.nests = Some(false);
756 return out;
757 }
758 let nests = is_nested(&names);
759 out.nests = Some(nests);
760 let mut types: HashMap<&str, &DataType> = HashMap::new();
762 let mut typed_apart = false;
763 for (name, dtype) in read.iter().flatten() {
764 match types.get(name.as_str()) {
765 Some(seen) if *seen != dtype => typed_apart = true,
766 Some(_) => {}
767 None => {
768 types.insert(name.as_str(), dtype);
769 }
770 }
771 }
772 let widest = names.iter().map(|f| f.len()).max().unwrap_or(0);
775 out.columns_differ = names.iter().any(|f| f.len() != widest) || !nests;
776 out.types_differ = typed_apart;
777 out
778}
779
780pub fn is_nested(files: &[Vec<String>]) -> bool {
787 let Some(widest) = files.iter().max_by_key(|f| f.len()) else {
788 return true;
789 };
790 let widest: std::collections::BTreeSet<&str> = widest.iter().map(String::as_str).collect();
791 files
792 .iter()
793 .all(|file| file.iter().all(|name| widest.contains(name.as_str())))
794}
795
796pub fn top_level_columns(leaves: &[String]) -> Vec<String> {
800 let mut seen = std::collections::HashSet::new();
801 leaves
802 .iter()
803 .map(|leaf| leaf.split_once('.').map_or(leaf.as_str(), |(root, _)| root))
804 .filter(|root| seen.insert(root.to_string()))
805 .map(str::to_string)
806 .collect()
807}
808
809#[derive(Debug, Clone)]
811pub struct DatasetSchema {
812 pub schema: Arc<Schema>,
813 pub columns: Vec<ColumnDrift>,
814 pub omitted: Vec<Vec<(PlSmallStr, DataType)>>,
818 pub unreadable: Vec<usize>,
820 pub files: usize,
822 pub groups: Vec<DriftGroup>,
825 pub file_group: Vec<u32>,
827 pub origin: SchemaOrigin,
828 pub read_as_text: Vec<PlSmallStr>,
831 pub empty_files: usize,
833 pub median_row_group_bytes: Option<usize>,
836 pub median_file_bytes: Option<usize>,
838 pub column_ranges: HashMap<PlSmallStr, ColumnRange>,
842 pub partition_layouts: Vec<(Vec<String>, usize)>,
846 pub partition_layouts_dropped: (usize, usize),
848 pub skipped: SkippedFiles,
850 pub listed_files: usize,
852}
853
854#[derive(Debug, Clone, Default, PartialEq, Eq, Hash)]
857pub struct DriftGroup {
858 pub absent: Vec<PlSmallStr>,
860 pub unread: Vec<PlSmallStr>,
863}
864
865impl DriftGroup {
866 pub fn is_empty(&self) -> bool {
867 self.absent.is_empty() && self.unread.is_empty()
868 }
869}
870
871impl DatasetSchema {
872 pub fn drifting(&self) -> impl Iterator<Item = &ColumnDrift> {
874 let readable = self.files - self.unreadable.len();
875 self.columns.iter().filter(move |c| !c.is_uniform(readable))
876 }
877
878 pub fn with_partition_layouts(mut self, root: &str, paths: &[String]) -> DatasetSchema {
882 const KEPT: usize = 64;
885 let mut counts: HashMap<Vec<String>, usize> = HashMap::new();
887 for path in paths {
888 let Some(below) = path.strip_prefix(root) else {
891 continue;
892 };
893 let keys = partition_keys_of(below);
894 if keys.is_empty() {
895 continue;
896 }
897 *counts.entry(keys).or_insert(0) += 1;
898 }
899 let mut counts: Vec<(Vec<String>, usize)> = counts.into_iter().collect();
900 counts.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
902 let dropped = &counts[counts.len().min(KEPT)..];
904 self.partition_layouts_dropped =
905 (dropped.len(), dropped.iter().map(|(_, files)| files).sum());
906 counts.truncate(KEPT);
907 self.partition_layouts = counts;
908 self.listed_files = paths.len();
909 self.column_ranges = self.ranges_of_columns(root, paths);
910 self
911 }
912
913 fn ranges_of_columns(&self, root: &str, paths: &[String]) -> HashMap<PlSmallStr, ColumnRange> {
916 if matches!(self.origin, SchemaOrigin::FooterSample { .. })
919 || !self.unreadable.is_empty()
920 || self.file_group.len() != paths.len()
921 {
922 return HashMap::new();
923 }
924 struct Seen {
927 first_present: Option<usize>,
928 last_absent: Option<usize>,
929 with: Vec<String>,
932 without: Vec<String>,
933 unplaced: bool,
936 }
937 let partition_of = |index: usize| -> Option<String> {
938 let below = paths.get(index)?.strip_prefix(root)?;
939 let values = partition_values_of(below);
940 (!values.is_empty()).then(|| values.join("/"))
941 };
942 let reads_in_order = (0..paths.len())
945 .filter_map(&partition_of)
946 .collect::<Vec<_>>()
947 .windows(2)
948 .all(|pair| natural_cmp(&pair[0], &pair[1]) != std::cmp::Ordering::Greater);
949 let drifting: Vec<&ColumnDrift> = self
952 .columns
953 .iter()
954 .filter(|column| column.present_in > 0 && column.present_in < self.files)
955 .collect();
956 if drifting.is_empty() {
957 return HashMap::new();
958 }
959 let where_in_drifting: HashMap<&PlSmallStr, usize> = drifting
962 .iter()
963 .enumerate()
964 .map(|(at, column)| (&column.name, at))
965 .collect();
966 let missing_by_group: Vec<Vec<bool>> = self
967 .groups
968 .iter()
969 .map(|group| {
970 let mut missing = vec![false; drifting.len()];
971 for name in &group.absent {
972 if let Some(at) = where_in_drifting.get(name) {
973 missing[*at] = true;
974 }
975 }
976 missing
977 })
978 .collect();
979 let none_missing: Vec<bool> = vec![false; drifting.len()];
980 let mut seen: Vec<Seen> = (0..drifting.len())
982 .map(|_| Seen {
983 first_present: None,
984 last_absent: None,
985 with: Vec::new(),
986 without: Vec::new(),
987 unplaced: false,
988 })
989 .collect();
990
991 for (index, group) in self.file_group.iter().enumerate() {
992 let missing: &[bool] = missing_by_group
993 .get(*group as usize)
994 .map(Vec::as_slice)
995 .unwrap_or(&none_missing);
996 let here = partition_of(index);
998 for (at, absent) in missing.iter().enumerate() {
999 let entry = &mut seen[at];
1000 let seen_of = if *absent {
1001 entry.last_absent = Some(index);
1002 &mut entry.without
1003 } else {
1004 entry.first_present.get_or_insert(index);
1005 entry.unplaced |= here.is_none();
1006 &mut entry.with
1007 };
1008 if seen_of.len() < 2
1009 && let Some(partition) = here.as_ref()
1010 && !seen_of.contains(partition)
1011 {
1012 seen_of.push(partition.clone());
1013 }
1014 }
1015 }
1016 seen.into_iter()
1017 .zip(&drifting)
1018 .filter_map(|(entry, column)| {
1019 let name = column.name.clone();
1020 let first = entry.first_present?;
1021 if !entry.unplaced
1024 && entry.with.len() == 1
1025 && entry
1027 .without
1028 .iter()
1029 .any(|other| {
1030 !partition_holds(&entry.with[0], other)
1031 && !same_place(&entry.with[0], other)
1032 })
1033 {
1034 return Some((name, ColumnRange::Only(entry.with[0].clone())));
1035 }
1036 let last_absent = entry.last_absent?;
1039 if last_absent > first || !reads_in_order {
1042 return None;
1043 }
1044 let (ends, begins) = (partition_of(last_absent)?, partition_of(first)?);
1045 (!same_place(&ends, &begins)
1048 && !partition_holds(&begins, &ends)
1049 && !partition_holds(&ends, &begins))
1050 .then_some((name, ColumnRange::NoneBefore(begins)))
1051 })
1052 .collect()
1053 }
1054
1055 pub fn with_skipped(mut self, skipped: SkippedFiles) -> DatasetSchema {
1057 self.skipped = skipped;
1058 self
1059 }
1060
1061 pub fn reading_as_text(&self, as_text: &[PlSmallStr]) -> DatasetSchema {
1066 let mut out = self.clone();
1067 if as_text.is_empty() {
1068 return out;
1069 }
1070 out.schema = crate::formats::schema_union::text_schema(&self.schema, as_text);
1071 for column in &mut out.columns {
1072 if as_text.contains(&column.name) {
1073 column.dtype = DataType::String;
1074 column.conflicting_files = 0;
1075 column.conflicting_types.clear();
1076 }
1079 }
1080 for group in &mut out.groups {
1081 group.unread.retain(|name| !as_text.contains(name));
1082 }
1083 out.read_as_text = as_text.to_vec();
1084 out
1085 }
1086
1087 pub fn drifts(&self) -> bool {
1090 self.groups.iter().any(|g| !g.is_empty())
1091 }
1092}
1093
1094pub const FOOTERS_AT_ONCE: usize = 64;
1098
1099pub fn footers_to_cache(
1102 footers: &[Option<FileFooter>],
1103) -> (
1104 Vec<crate::cache::CachedFooter>,
1105 Vec<Vec<(String, DataType)>>,
1106) {
1107 let mut schemas = Vec::new();
1108 let cached = footers
1109 .iter()
1110 .map(|footer| match footer {
1111 None => crate::cache::CachedFooter::default(),
1112 Some(f) => crate::cache::CachedFooter {
1113 schema: Some(crate::cache::DatasetShape::intern_schema(
1114 &mut schemas,
1115 &f.schema,
1116 )),
1117 row_group_rows: f.row_group_rows.clone(),
1118 row_group_bytes: f.row_group_bytes.clone(),
1119 column_bytes: f.column_bytes.iter().map(|(_, bytes)| *bytes).collect(),
1120 },
1121 })
1122 .collect();
1123 (cached, schemas)
1124}
1125
1126pub fn footers_from_cache(
1130 cached: &[crate::cache::CachedFooter],
1131 schemas: &[Vec<(String, DataType)>],
1132 file_bytes: &[u64],
1133) -> Option<Vec<Option<FileFooter>>> {
1134 if cached.len() != file_bytes.len() {
1135 return None;
1136 }
1137 let shared: Vec<Arc<Schema>> = (0..schemas.len())
1139 .map(|at| crate::cache::DatasetShape::schema_at(schemas, at).map(Arc::new))
1140 .collect::<Option<_>>()?;
1141 cached
1142 .iter()
1143 .zip(file_bytes)
1144 .map(|(f, &bytes)| {
1145 let Some(at) = f.schema else {
1146 return Some(None);
1147 };
1148 let schema = shared.get(at)?;
1149 let column_bytes = schema
1150 .iter_names()
1151 .zip(&f.column_bytes)
1152 .map(|(name, bytes)| (name.to_string(), *bytes))
1153 .collect();
1154 Some(Some(FileFooter {
1155 schema: schema.clone(),
1156 row_group_rows: f.row_group_rows.clone(),
1157 row_group_bytes: f.row_group_bytes.clone(),
1158 file_bytes: bytes as usize,
1159 column_bytes,
1160 }))
1161 })
1162 .collect()
1163}
1164
1165type FooterHook = Arc<dyn Fn(&std::path::Path) + Send + Sync>;
1167
1168static FOOTER_HOOKS: std::sync::Mutex<Vec<(u64, std::path::PathBuf, FooterHook)>> =
1169 std::sync::Mutex::new(Vec::new());
1170static FOOTER_HOOKS_SET: AtomicUsize = AtomicUsize::new(0);
1172
1173#[doc(hidden)]
1176pub fn on_local_footer_read(
1177 dir: &std::path::Path,
1178 hook: impl Fn(&std::path::Path) + Send + Sync + 'static,
1179) -> FooterHookGuard {
1180 static NEXT: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1181 let id = NEXT.fetch_add(1, Ordering::Relaxed);
1182 FOOTER_HOOKS
1183 .lock()
1184 .unwrap_or_else(|e| e.into_inner())
1185 .push((id, dir.to_path_buf(), Arc::new(hook)));
1186 FOOTER_HOOKS_SET.fetch_add(1, Ordering::Release);
1187 FooterHookGuard(id)
1188}
1189
1190#[doc(hidden)]
1192pub struct FooterHookGuard(u64);
1193
1194impl Drop for FooterHookGuard {
1195 fn drop(&mut self) {
1196 FOOTER_HOOKS
1197 .lock()
1198 .unwrap_or_else(|e| e.into_inner())
1199 .retain(|(id, _, _)| *id != self.0);
1200 FOOTER_HOOKS_SET.fetch_sub(1, Ordering::Release);
1201 }
1202}
1203
1204pub(crate) fn before_local_footer_read(path: &std::path::Path) {
1206 if FOOTER_HOOKS_SET.load(Ordering::Acquire) == 0 {
1207 return;
1208 }
1209 let hooks: Vec<FooterHook> = FOOTER_HOOKS
1211 .lock()
1212 .unwrap_or_else(|e| e.into_inner())
1213 .iter()
1214 .filter(|(_, dir, _)| path.starts_with(dir))
1215 .map(|(_, _, hook)| hook.clone())
1216 .collect();
1217 for hook in hooks {
1218 hook(path);
1219 }
1220}
1221
1222pub const MAX_FOOTER_READS: usize = 20_000;
1225
1226pub fn footers_to_read(files: usize) -> Vec<usize> {
1230 if files <= MAX_FOOTER_READS {
1231 return (0..files).collect();
1232 }
1233 let last = files - 1;
1234 let mut sample: Vec<usize> = (0..MAX_FOOTER_READS)
1235 .map(|i| i * last / (MAX_FOOTER_READS - 1))
1236 .collect();
1237 sample.dedup();
1238 sample
1239}
1240
1241pub fn ends_of(files: usize) -> Vec<usize> {
1245 match files {
1246 0 => Vec::new(),
1247 1 => vec![0],
1248 n => vec![0, n - 1],
1249 }
1250}
1251
1252pub fn union_sampled(
1257 files: usize,
1258 read: &[usize],
1259 footers: &[Option<FileFooter>],
1260) -> DatasetSchema {
1261 let origin = if read.len() == files {
1262 SchemaOrigin::AllFooters(files)
1263 } else {
1264 SchemaOrigin::FooterSample {
1265 read: read.len(),
1266 total: files,
1267 }
1268 };
1269 let mut union = union_file_schemas(footers, origin);
1270 let mut omitted = vec![Vec::new(); files];
1271 let mut file_group = vec![0u32; files];
1272 for ((columns, group), &index) in union.omitted.iter().zip(union.file_group.iter()).zip(read) {
1273 omitted[index] = columns.clone();
1274 file_group[index] = *group;
1275 }
1276 union.omitted = omitted;
1277 union.file_group = file_group;
1278 union.unreadable = union
1279 .unreadable
1280 .iter()
1281 .filter_map(|i| read.get(*i).copied())
1282 .collect();
1283 union
1284}
1285
1286pub fn readable_paths<'a>(paths: &'a [String], unreadable: &[usize]) -> Cow<'a, [String]> {
1290 if unreadable.is_empty() {
1291 return Cow::Borrowed(paths);
1293 }
1294 debug_assert!(unreadable.windows(2).all(|pair| pair[0] < pair[1]));
1296 Cow::Owned(
1297 paths
1298 .iter()
1299 .enumerate()
1300 .filter(|(index, _)| unreadable.binary_search(index).is_err())
1301 .map(|(_, path)| path.clone())
1302 .collect(),
1303 )
1304}
1305
1306pub fn union_file_schemas(files: &[Option<FileFooter>], origin: SchemaOrigin) -> DatasetSchema {
1309 let unreadable = files
1310 .iter()
1311 .enumerate()
1312 .filter_map(|(i, f)| f.is_none().then_some(i))
1313 .collect();
1314
1315 let mut order: Vec<PlSmallStr> = Vec::new();
1317 let mut seen: HashMap<PlSmallStr, usize> = HashMap::new();
1318 let mut push = |name: &PlSmallStr, order: &mut Vec<PlSmallStr>| {
1319 if !seen.contains_key(name) {
1320 seen.insert(name.clone(), order.len());
1321 order.push(name.clone());
1322 }
1323 };
1324 if let Some(newest) = files.iter().rev().flatten().next() {
1325 for name in newest.schema.iter_names() {
1326 push(name, &mut order);
1327 }
1328 }
1329 for file in files.iter().flatten() {
1330 for name in file.schema.iter_names() {
1331 push(name, &mut order);
1332 }
1333 }
1334
1335 let mut sightings: Vec<Vec<(DataType, usize)>> = vec![Vec::new(); order.len()];
1337 for file in files.iter().flatten() {
1338 for (name, dtype) in file.schema.iter() {
1339 let Some(&index) = seen.get(name) else {
1340 continue;
1341 };
1342 sightings[index].push((dtype.clone(), file.rows()));
1343 }
1344 }
1345
1346 let mut schema = Schema::with_capacity(order.len());
1347 let mut columns = Vec::with_capacity(order.len());
1348 for (name, seen_types) in order.iter().zip(sightings.iter()) {
1349 let chosen = choose_dtype(seen_types);
1350 let conflicting_types = seen_types
1351 .iter()
1352 .map(|(d, _)| d)
1353 .filter(|d| !fits(d, &chosen))
1354 .fold(Vec::new(), |mut acc: Vec<DataType>, d| {
1355 if !acc.contains(d) {
1356 acc.push(d.clone());
1357 }
1358 acc
1359 });
1360 columns.push(ColumnDrift {
1361 name: name.clone(),
1362 present_in: seen_types.len(),
1363 conflicting_files: seen_types.iter().filter(|(d, _)| !fits(d, &chosen)).count(),
1364 widened: seen_types
1365 .iter()
1366 .any(|(d, _)| *d != chosen && fits(d, &chosen)),
1367 conflicting_types,
1368 dtype: chosen.clone(),
1369 });
1370 schema.with_column(name.clone(), chosen);
1371 }
1372
1373 let mut groups: Vec<DriftGroup> = vec![DriftGroup::default()];
1376 let mut group_of: HashMap<DriftGroup, u32> = HashMap::from([(DriftGroup::default(), 0)]);
1377 let mut file_group = Vec::with_capacity(files.len());
1378 let mut omitted = Vec::with_capacity(files.len());
1379 for file in files {
1380 let Some(file) = file else {
1381 file_group.push(0);
1382 omitted.push(Vec::new());
1383 continue;
1384 };
1385 let stored: Vec<(PlSmallStr, DataType)> = file
1386 .schema
1387 .iter()
1388 .filter(|(name, dtype)| schema.get(name).is_some_and(|target| !fits(dtype, target)))
1389 .map(|(name, dtype)| (name.clone(), dtype.clone()))
1390 .collect();
1391 let unread: Vec<PlSmallStr> = stored.iter().map(|(name, _)| name.clone()).collect();
1392 let absent: Vec<PlSmallStr> = schema
1393 .iter_names()
1394 .filter(|name| !file.schema.contains(name))
1395 .cloned()
1396 .collect();
1397 omitted.push(stored);
1398 let group = DriftGroup { absent, unread };
1399 let next = groups.len() as u32;
1400 let id = *group_of.entry(group.clone()).or_insert_with(|| {
1401 groups.push(group);
1402 next
1403 });
1404 file_group.push(id);
1405 }
1406
1407 DatasetSchema {
1408 schema: Arc::new(schema),
1409 columns,
1410 omitted,
1411 unreadable,
1412 files: files.len(),
1413 groups,
1414 file_group,
1415 origin,
1416 read_as_text: Vec::new(),
1417 empty_files: files.iter().flatten().filter(|f| f.rows() == 0).count(),
1418 median_file_bytes: median(files.iter().flatten().map(|f| f.file_bytes)),
1419 column_ranges: HashMap::new(),
1420 skipped: SkippedFiles::default(),
1421 partition_layouts: Vec::new(),
1422 partition_layouts_dropped: (0, 0),
1423 listed_files: 0,
1424 median_row_group_bytes: median(
1425 files
1426 .iter()
1427 .flatten()
1428 .flat_map(|f| f.row_group_bytes.iter().copied()),
1429 ),
1430 }
1431}
1432
1433fn partition_holds(outer: &str, inner: &str) -> bool {
1437 inner == outer
1438 || inner
1439 .strip_prefix(outer)
1440 .is_some_and(|rest| rest.starts_with('/'))
1441}
1442
1443fn same_place(a: &str, b: &str) -> bool {
1446 fn sorted(path: &str) -> Vec<&str> {
1447 let mut segments: Vec<&str> = path.split('/').collect();
1448 segments.sort_by_key(|segment| segment.split_once('=').map(|(key, _)| key));
1449 segments
1450 }
1451 let (a, b) = (sorted(a), sorted(b));
1452 a.len() == b.len()
1453 && a.iter()
1454 .zip(&b)
1455 .all(|(x, y)| natural_cmp(x, y) == std::cmp::Ordering::Equal)
1456}
1457
1458fn natural_cmp(a: &str, b: &str) -> std::cmp::Ordering {
1461 use std::cmp::Ordering;
1462 let (mut a, mut b) = (a.as_bytes(), b.as_bytes());
1463 loop {
1464 match (a.first(), b.first()) {
1465 (None, None) => return Ordering::Equal,
1466 (None, _) => return Ordering::Less,
1467 (_, None) => return Ordering::Greater,
1468 (Some(x), Some(y)) if x.is_ascii_digit() && y.is_ascii_digit() => {
1469 let digits = |s: &[u8]| s.iter().take_while(|c| c.is_ascii_digit()).count();
1470 let (na, nb) = (digits(a), digits(b));
1471 let (xs, ys) = (&a[..na], &b[..nb]);
1473 fn trim(s: &[u8]) -> &[u8] {
1474 let lead = s.iter().take_while(|c| **c == b'0').count();
1475 &s[lead.min(s.len().saturating_sub(1))..]
1476 }
1477 let (tx, ty) = (trim(xs), trim(ys));
1478 match tx.len().cmp(&ty.len()).then_with(|| tx.cmp(ty)) {
1479 Ordering::Equal => {}
1480 other => return other,
1481 }
1482 a = &a[na..];
1483 b = &b[nb..];
1484 }
1485 (Some(x), Some(y)) => match x.cmp(y) {
1486 Ordering::Equal => {
1487 a = &a[1..];
1488 b = &b[1..];
1489 }
1490 other => return other,
1491 },
1492 }
1493 }
1494}
1495
1496fn partition_values_of(path: &str) -> Vec<String> {
1500 #[cfg(windows)]
1501 let separators: &[char] = &['/', '\\'];
1502 #[cfg(not(windows))]
1503 let separators: &[char] = &['/'];
1504 let mut segments: Vec<&str> = path.split(separators).collect();
1505 segments.pop();
1507 segments
1508 .into_iter()
1509 .filter(|segment| {
1510 segment
1511 .split_once('=')
1512 .is_some_and(|(key, _)| !key.is_empty())
1513 })
1514 .map(|segment| segment.to_string())
1515 .collect()
1516}
1517
1518fn partition_keys_of(path: &str) -> Vec<String> {
1521 let mut keys: Vec<String> = Vec::new();
1522 #[cfg(windows)]
1524 let separators: &[char] = &['/', '\\'];
1525 #[cfg(not(windows))]
1526 let separators: &[char] = &['/'];
1527 let mut segments: Vec<&str> = path.split(separators).collect();
1528 segments.pop();
1529 for segment in segments {
1530 if let Some((key, _)) = segment.split_once('=')
1531 && !key.is_empty()
1532 {
1533 keys.push(key.to_string());
1534 }
1535 }
1536 keys.sort();
1538 keys.dedup();
1539 keys
1540}
1541
1542fn median(sizes: impl Iterator<Item = usize>) -> Option<usize> {
1545 let mut sizes: Vec<usize> = sizes.collect();
1546 if sizes.is_empty() {
1547 return None;
1548 }
1549 sizes.sort_unstable();
1550 Some(sizes[(sizes.len() - 1) / 2])
1551}
1552
1553fn choose_dtype(seen: &[(DataType, usize)]) -> DataType {
1557 let mut distinct: Vec<DataType> = Vec::new();
1558 for (dtype, _) in seen {
1559 if !distinct.contains(dtype) {
1560 distinct.push(dtype.clone());
1561 }
1562 }
1563 match distinct.as_slice() {
1564 [] => return DataType::Null,
1565 [only] => return only.clone(),
1566 _ => {}
1567 }
1568 let mut candidates = distinct.clone();
1571 for dtype in &distinct {
1572 let folded = distinct
1573 .iter()
1574 .filter(|other| widen(dtype, other).is_some())
1575 .try_fold(dtype.clone(), |acc, other| widen(&acc, other));
1576 if let Some(folded) = folded
1577 && !candidates.contains(&folded)
1578 {
1579 candidates.push(folded);
1580 }
1581 }
1582 let mut best: Option<(DataType, usize, usize)> = None;
1583 for candidate in candidates {
1584 let rows: usize = seen
1585 .iter()
1586 .filter(|(d, _)| fits(d, &candidate))
1587 .map(|(_, rows)| rows)
1588 .sum();
1589 let files = seen.iter().filter(|(d, _)| fits(d, &candidate)).count();
1590 let better = best
1591 .as_ref()
1592 .is_none_or(|(_, best_rows, best_files)| (rows, files) > (*best_rows, *best_files));
1593 if better {
1594 best = Some((candidate, rows, files));
1595 }
1596 }
1597 best.map(|(d, _, _)| d).unwrap_or(DataType::Null)
1598}
1599
1600pub fn fits(from: &DataType, to: &DataType) -> bool {
1604 widen(from, to).as_ref() == Some(to)
1605}
1606
1607pub fn widen(a: &DataType, b: &DataType) -> Option<DataType> {
1609 use DataType::*;
1610 if a == b {
1611 return Some(a.clone());
1612 }
1613 match (a, b) {
1614 (Null, other) | (other, Null) => Some(other.clone()),
1615 _ if a.is_integer() && b.is_integer() => widen_integers(a, b),
1616 _ if (a.is_integer() || a.is_float()) && (b.is_integer() || b.is_float()) => Some(Float64),
1617 (Datetime(a_unit, a_zone), Datetime(b_unit, b_zone)) if a_zone == b_zone => {
1618 Some(Datetime(finer_unit(*a_unit, *b_unit), a_zone.clone()))
1619 }
1620 (List(a_inner), List(b_inner)) => widen(a_inner, b_inner).map(|t| List(Box::new(t))),
1621 (Struct(a_fields), Struct(b_fields)) => widen_structs(a_fields, b_fields),
1622 _ => None,
1623 }
1624}
1625
1626fn widen_integers(a: &DataType, b: &DataType) -> Option<DataType> {
1629 use DataType::*;
1630 let signed = |d: &DataType| matches!(d, Int8 | Int16 | Int32 | Int64 | Int128);
1631 let bits = |d: &DataType| match d {
1632 Int8 | UInt8 => 8u32,
1633 Int16 | UInt16 => 16,
1634 Int32 | UInt32 => 32,
1635 Int64 | UInt64 => 64,
1636 _ => 128,
1637 };
1638 if signed(a) == signed(b) {
1639 let wider = if bits(a) >= bits(b) { a } else { b };
1640 return Some(wider.clone());
1641 }
1642 let (unsigned, sgn) = if signed(a) { (b, a) } else { (a, b) };
1644 let needed = match bits(unsigned) {
1645 8 => Int16,
1646 16 => Int32,
1647 32 => Int64,
1648 64 => return None,
1649 _ => return None,
1650 };
1651 Some(if bits(sgn) >= bits(&needed) {
1652 sgn.clone()
1653 } else {
1654 needed
1655 })
1656}
1657
1658fn widen_structs(a: &[Field], b: &[Field]) -> Option<DataType> {
1661 let mut fields: Vec<Field> = Vec::with_capacity(a.len() + b.len());
1662 for field in a {
1663 let widened = match b.iter().find(|other| other.name() == field.name()) {
1664 Some(other) => widen(field.dtype(), other.dtype())?,
1665 None => field.dtype().clone(),
1666 };
1667 fields.push(Field::new(field.name().clone(), widened));
1668 }
1669 for field in b {
1670 if !a.iter().any(|other| other.name() == field.name()) {
1671 fields.push(field.clone());
1672 }
1673 }
1674 Some(DataType::Struct(fields))
1675}
1676
1677fn finer_unit(a: TimeUnit, b: TimeUnit) -> TimeUnit {
1678 let rank = |u: TimeUnit| match u {
1679 TimeUnit::Milliseconds => 0,
1680 TimeUnit::Microseconds => 1,
1681 TimeUnit::Nanoseconds => 2,
1682 };
1683 if rank(a) >= rank(b) { a } else { b }
1684}
1685
1686pub fn with_partition_columns(
1688 file_schema: &Schema,
1689 partition_columns: &[String],
1690 values: &[(String, String)],
1691) -> Schema {
1692 let part_set: HashSet<&str> = partition_columns.iter().map(String::as_str).collect();
1693 let mut merged = Schema::with_capacity(partition_columns.len() + file_schema.len());
1694 for name in partition_columns {
1695 merged.with_column(
1696 name.clone().into(),
1697 crate::formats::readers::hive::partition_dtype(name, file_schema, values),
1698 );
1699 }
1700 for (name, dtype) in file_schema.iter() {
1701 if !part_set.contains(name.as_str()) {
1702 merged.with_column(name.clone(), dtype.clone());
1703 }
1704 }
1705 merged
1706}
1707
1708pub fn partition_columns_of_key(key: &str) -> Vec<String> {
1710 let mut columns = Vec::new();
1711 let mut seen = HashSet::new();
1712 for segment in key.split('/') {
1713 if let Some((name, _)) = segment.split_once('=')
1714 && !name.is_empty()
1715 && seen.insert(name.to_string())
1716 {
1717 columns.push(name.to_string());
1718 }
1719 }
1720 columns
1721}
1722
1723pub fn partitions_of_listing(first: &str, newest: &str) -> (Vec<String>, Vec<(String, String)>) {
1727 let values = [first, newest]
1728 .iter()
1729 .flat_map(|key| key.split('/'))
1730 .filter_map(|segment| segment.split_once('='))
1731 .map(|(k, v)| (k.to_string(), v.to_string()))
1732 .collect();
1733 (partition_columns_of_key(newest), values)
1734}
1735
1736pub struct FooterCount<F> {
1740 files: usize,
1741 counted: Vec<usize>,
1743 footers: std::sync::Mutex<Vec<Option<F>>>,
1745}
1746
1747pub struct Counted<F> {
1749 pub row_groups: Vec<Vec<usize>>,
1751 pub whole: Option<Vec<Option<F>>>,
1753}
1754
1755impl<F: Clone> FooterCount<F> {
1756 pub fn new(
1759 files: usize,
1760 counted: Vec<usize>,
1761 known: impl IntoIterator<Item = (usize, Option<F>)>,
1762 ) -> Self {
1763 let mut footers = vec![None; files];
1764 for (index, footer) in known {
1765 if let Some(slot) = footers.get_mut(index) {
1766 *slot = footer;
1767 }
1768 }
1769 Self {
1770 files,
1771 counted,
1772 footers: std::sync::Mutex::new(footers),
1773 }
1774 }
1775
1776 pub fn count(
1780 &self,
1781 read: impl FnOnce(&[usize]) -> Option<Vec<Option<F>>>,
1782 row_groups: impl Fn(&F) -> Vec<usize>,
1783 ) -> Option<Counted<F>> {
1784 let mut footers = self.footers.lock().unwrap_or_else(|e| e.into_inner());
1785 let missing = self.missing(&mut footers);
1786 let read = if missing.is_empty() {
1787 Vec::new()
1788 } else {
1789 read(&missing)?
1790 };
1791 Some(self.settle(&mut footers, missing, read, row_groups))
1792 }
1793
1794 fn missing(&self, footers: &mut Vec<Option<F>>) -> Vec<usize> {
1795 if footers.is_empty() {
1796 *footers = vec![None; self.files];
1798 }
1799 self.counted
1800 .iter()
1801 .copied()
1802 .filter(|&index| footers[index].is_none())
1803 .collect()
1804 }
1805
1806 fn settle(
1807 &self,
1808 footers: &mut Vec<Option<F>>,
1809 missing: Vec<usize>,
1810 read: Vec<Option<F>>,
1811 row_groups: impl Fn(&F) -> Vec<usize>,
1812 ) -> Counted<F> {
1813 if footers.is_empty() {
1814 *footers = vec![None; self.files];
1815 }
1816 for (index, footer) in missing.into_iter().zip(read) {
1817 footers[index] = footer;
1818 }
1819 let groups: Vec<Vec<usize>> = self
1820 .counted
1821 .iter()
1822 .map(|&index| footers[index].as_ref().map(&row_groups).unwrap_or_default())
1823 .collect();
1824 let whole = if footers.iter().all(Option::is_some) {
1825 Some(std::mem::take(footers))
1826 } else {
1827 if self.counted.iter().all(|&index| footers[index].is_some()) {
1828 footers.clear();
1830 }
1831 None
1832 };
1833 Counted {
1834 row_groups: groups,
1835 whole,
1836 }
1837 }
1838}
1839
1840pub const DRIFT_COLUMN: &str = "__datui_row";
1844
1845static NOTHING_MISSING: DriftGroup = DriftGroup {
1847 absent: Vec::new(),
1848 unread: Vec::new(),
1849};
1850
1851#[derive(Debug, Clone, Default)]
1854pub struct ScanDrift {
1855 group_of: HashMap<String, u32>,
1856 row_of: HashMap<String, usize>,
1859 stored_of: HashMap<String, Vec<(PlSmallStr, DataType)>>,
1861 pub groups: Vec<DriftGroup>,
1862}
1863
1864impl ScanDrift {
1865 pub fn new(paths: &[String], dataset: &DatasetSchema, file_rows: &[usize]) -> Option<Self> {
1868 if !dataset.drifts() || file_rows.len() != paths.len() {
1869 return None;
1870 }
1871 let group_of = paths
1872 .iter()
1873 .zip(dataset.file_group.iter())
1874 .filter(|(_, group)| **group != 0)
1875 .map(|(path, group)| (path.clone(), *group))
1876 .collect();
1877 let mut row = 0usize;
1878 let mut row_of = HashMap::with_capacity(paths.len());
1879 for (path, rows) in paths.iter().zip(file_rows) {
1880 row_of.insert(path.clone(), row);
1881 row += rows;
1882 }
1883 let stored_of = paths
1884 .iter()
1885 .zip(dataset.omitted.iter())
1886 .filter(|(_, stored)| !stored.is_empty())
1887 .map(|(path, stored)| (path.clone(), stored.clone()))
1888 .collect();
1889 Some(ScanDrift {
1890 group_of,
1891 row_of,
1892 stored_of,
1893 groups: dataset.groups.clone(),
1894 })
1895 }
1896
1897 pub fn group(&self, path: &str) -> u32 {
1899 self.group_of.get(path).copied().unwrap_or(0)
1900 }
1901
1902 fn first_row(&self, path: &str) -> usize {
1904 self.row_of.get(path).copied().unwrap_or(0)
1905 }
1906
1907 fn stored_type(&self, path: &str, column: &PlSmallStr) -> Option<&DataType> {
1909 self.stored_of
1910 .get(path)?
1911 .iter()
1912 .find(|(name, _)| name == column)
1913 .map(|(_, dtype)| dtype)
1914 }
1915
1916 fn unread(&self, path: &str) -> &[PlSmallStr] {
1918 self.groups
1919 .get(self.group(path) as usize)
1920 .unwrap_or(&NOTHING_MISSING)
1921 .unread
1922 .as_slice()
1923 }
1924}
1925
1926pub fn lenient_scan(
1943 paths: &[String],
1944 schema: Arc<Schema>,
1945 cloud_options: Option<polars::io::cloud::CloudOptions>,
1946 drift: Option<&ScanDrift>,
1947 as_text: &[PlSmallStr],
1948) -> PolarsResult<LazyFrame> {
1949 let Some(drift) = drift else {
1950 return scan_run(paths, &schema, cloud_options, &[], None, &[]);
1951 };
1952 let as_text: Vec<PlSmallStr> = as_text
1955 .iter()
1956 .filter(|name| {
1957 schema.get(name).is_some_and(can_read_as_text)
1958 && paths
1959 .iter()
1960 .all(|path| drift.stored_type(path, name).is_none_or(can_read_as_text))
1961 })
1962 .cloned()
1963 .collect();
1964 let as_text = as_text.as_slice();
1965 let unread_of = |path: &str| -> Vec<PlSmallStr> {
1967 drift
1968 .unread(path)
1969 .iter()
1970 .filter(|name| !as_text.contains(name))
1971 .cloned()
1972 .collect()
1973 };
1974 let key_of = |path: &str| -> (Vec<PlSmallStr>, Vec<Option<DataType>>) {
1976 (
1977 unread_of(path),
1978 as_text
1979 .iter()
1980 .map(|name| drift.stored_type(path, name).cloned())
1981 .collect(),
1982 )
1983 };
1984 let mut runs: Vec<LazyFrame> = Vec::new();
1985 let mut start = 0;
1986 while start < paths.len() {
1987 let key = key_of(&paths[start]);
1988 let end = paths[start..]
1989 .iter()
1990 .position(|path| key_of(path) != key)
1991 .map_or(paths.len(), |offset| start + offset);
1992 let (omit, stored) = key;
1993 let read_as: Vec<(PlSmallStr, DataType)> = as_text
1995 .iter()
1996 .zip(stored)
1997 .map(|(name, stored)| {
1998 let dtype = stored.or_else(|| schema.get(name).cloned());
1999 (name.clone(), dtype.unwrap_or(DataType::String))
2000 })
2001 .collect();
2002 runs.push(scan_run(
2003 &paths[start..end],
2004 &schema,
2005 cloud_options.clone(),
2006 &omit,
2007 Some(drift.first_row(&paths[start])),
2008 &read_as,
2009 )?);
2010 start = end;
2011 }
2012 match runs.len() {
2013 1 => Ok(runs.remove(0)),
2014 _ => concat(
2015 runs,
2016 UnionArgs {
2017 rechunk: false,
2018 parallel: true,
2019 ..Default::default()
2020 },
2021 ),
2022 }
2023}
2024
2025pub fn can_read_as_text(dtype: &DataType) -> bool {
2031 match dtype {
2032 DataType::Binary | DataType::BinaryOffset => false,
2034 DataType::Duration(_) => false,
2035 DataType::List(_) | DataType::Array(_, _) => false,
2037 DataType::Struct(_) => true,
2039 DataType::Unknown(_) => false,
2040 _ => true,
2041 }
2042}
2043
2044impl ColumnDrift {
2045 pub fn can_read_as_text(&self) -> bool {
2048 can_read_as_text(&self.dtype) && self.conflicting_types.iter().all(can_read_as_text)
2049 }
2050}
2051
2052pub fn text_schema(schema: &Arc<Schema>, as_text: &[PlSmallStr]) -> Arc<Schema> {
2055 if as_text.is_empty() {
2056 return schema.clone();
2057 }
2058 let mut out = Schema::with_capacity(schema.len());
2059 for (name, dtype) in schema.iter() {
2060 let dtype = if as_text.contains(name) {
2061 DataType::String
2062 } else {
2063 dtype.clone()
2064 };
2065 out.with_column(name.clone(), dtype);
2066 }
2067 Arc::new(out)
2068}
2069
2070fn scan_run(
2074 urls: &[String],
2075 schema: &Arc<Schema>,
2076 cloud_options: Option<polars::io::cloud::CloudOptions>,
2077 omit: &[PlSmallStr],
2078 first_row: Option<usize>,
2079 read_as: &[(PlSmallStr, DataType)],
2080) -> PolarsResult<LazyFrame> {
2081 use polars::lazy::dsl::{
2082 CastColumnsPolicy, DslBuilder, ExtraColumnsPolicy, MissingColumnsPolicy, ScanSources,
2083 UnifiedScanArgs,
2084 };
2085 use polars::prelude::{Expr, NULL, col, lit};
2086 let sources = ScanSources::Paths(
2087 urls.iter()
2088 .map(|url| PlRefPath::new(url.as_str()))
2089 .collect(),
2090 );
2091 let target = if omit.is_empty() && read_as.is_empty() {
2092 schema.clone()
2093 } else {
2094 let mut reduced = Schema::with_capacity(schema.len());
2095 for (name, dtype) in schema.iter() {
2096 if omit.contains(name) {
2097 continue;
2098 }
2099 let dtype = read_as
2101 .iter()
2102 .find(|(column, _)| column == name)
2103 .map(|(_, dtype)| dtype)
2104 .unwrap_or(dtype);
2105 reduced.with_column(name.clone(), dtype.clone());
2106 }
2107 Arc::new(reduced)
2108 };
2109 let options = polars::io::parquet::read::ParquetOptions {
2110 schema: Some(target),
2111 ..Default::default()
2112 };
2113 let args = UnifiedScanArgs {
2114 cloud_options,
2115 hive_options: polars::io::HiveOptions::new_enabled(),
2116 glob: false,
2117 cast_columns_policy: CastColumnsPolicy {
2118 integer_upcast: true,
2119 integer_to_float_cast: true,
2120 float_upcast: true,
2121 datetime_nanoseconds_downcast: true,
2122 datetime_microseconds_downcast: true,
2123 datetime_milliseconds_upcast: true,
2124 datetime_microseconds_upcast: true,
2125 null_upcast: true,
2126 missing_struct_fields: MissingColumnsPolicy::Insert,
2127 extra_struct_fields: ExtraColumnsPolicy::Ignore,
2128 ..CastColumnsPolicy::ERROR_ON_MISMATCH
2129 },
2130 missing_columns_policy: MissingColumnsPolicy::Insert,
2131 extra_columns_policy: ExtraColumnsPolicy::Ignore,
2132 row_index: first_row.map(|first| polars::io::RowIndex {
2135 name: DRIFT_COLUMN.into(),
2136 offset: first as polars::prelude::IdxSize,
2137 }),
2138 ..Default::default()
2139 };
2140 let mut lf: LazyFrame = DslBuilder::scan_parquet(sources, options, args)?
2141 .build()
2142 .into();
2143 if !omit.is_empty() {
2144 let nulls: Vec<Expr> = omit
2145 .iter()
2146 .filter_map(|name| {
2147 let dtype = schema.get(name)?;
2148 Some(lit(NULL).cast(dtype.clone()).alias(name.clone()))
2149 })
2150 .collect();
2151 lf = lf.with_columns(nulls);
2152 }
2153 if !read_as.is_empty() {
2154 let texts: Vec<Expr> = read_as
2157 .iter()
2158 .map(|(name, _)| {
2159 crate::past_calendar::text_expr(col(name.clone()), CastOptions::NonStrict)
2160 .alias(name.clone())
2161 })
2162 .collect();
2163 lf = lf.with_columns(texts);
2164 }
2165 if first_row.is_some() {
2166 let mut ordered: Vec<Expr> = schema.iter_names().map(|name| col(name.clone())).collect();
2169 ordered.push(col(DRIFT_COLUMN));
2170 lf = lf.select(ordered);
2171 }
2172 Ok(lf)
2173}
2174
2175#[cfg(test)]
2176mod tests;