use std::borrow::Cow;
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use polars::chunked_array::cast::CastOptions;
use polars::prelude::{
DataType, Field, LazyFrame, PlRefPath, PlSmallStr, PolarsResult, Schema, TimeUnit, UnionArgs,
concat,
};
pub struct Pass<'a>(&'a FooterProgress);
impl Pass<'_> {
pub fn advance(&self) {
self.0.advance();
}
}
impl Drop for Pass<'_> {
fn drop(&mut self) {
self.0.done();
}
}
pub struct Listing<'a>(&'a FooterProgress);
impl Listing<'_> {
pub fn advance(&self) {
self.0.listed.fetch_add(1, Ordering::Relaxed);
}
pub fn counter(&self) -> std::sync::Arc<AtomicUsize> {
self.0.listed.clone()
}
pub fn add(&self, n: usize) {
self.0.listed.fetch_add(n, Ordering::Relaxed);
}
pub fn is_cancelled(&self) -> bool {
self.0.is_cancelled()
}
pub fn cancel_flag(&self) -> std::sync::Arc<std::sync::atomic::AtomicBool> {
self.0.cancel_flag()
}
}
impl Drop for Listing<'_> {
fn drop(&mut self) {
self.0.listing.store(false, Ordering::Release);
}
}
#[doc(hidden)]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct PassCount {
pub begun: usize,
pub read: usize,
pub total: usize,
}
#[derive(Debug, Default)]
pub struct FooterProgress {
read: AtomicUsize,
total: AtomicUsize,
passes: AtomicUsize,
last_total: AtomicUsize,
cancelled: std::sync::Arc<std::sync::atomic::AtomicBool>,
listed: std::sync::Arc<AtomicUsize>,
listing: std::sync::atomic::AtomicBool,
at_once: AtomicUsize,
estimate: std::sync::Mutex<Option<RowEstimate>>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct RowEstimate {
pub rows: u64,
pub sampled: usize,
pub files: usize,
}
impl RowEstimate {
pub fn of<'a>(
files: usize,
footers: impl IntoIterator<Item = &'a Option<FileFooter>>,
) -> Option<Self> {
let (sampled, rows) = footers
.into_iter()
.flatten()
.fold((0usize, 0u128), |(n, rows), f| {
(n + 1, rows + f.rows() as u128)
});
(sampled > 0).then(|| RowEstimate {
rows: u64::try_from(rows * files as u128 / sampled as u128).unwrap_or(u64::MAX),
sampled,
files,
})
}
}
pub const ESTIMATE_SAMPLE: usize = 2_000;
pub const COUNT_AT_ONCE: usize = 256;
pub fn random_sample(files: usize, n: usize, seed: u64) -> Vec<usize> {
if files <= n {
return (0..files).collect();
}
let mut state = seed;
let mut next = move || {
state = state.wrapping_add(0x9e37_79b9_7f4a_7c15);
let mut z = state;
z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
z = (z ^ (z >> 27)).wrapping_mul(0x94d0_49bb_1331_11eb);
z ^ (z >> 31)
};
let mut chosen = std::collections::BTreeSet::new();
while chosen.len() < n {
chosen.insert((next() % files as u64) as usize);
}
chosen.into_iter().collect()
}
impl FooterProgress {
pub fn listing(&self) -> Listing<'_> {
self.listed.store(0, Ordering::Relaxed);
self.listing.store(true, Ordering::Release);
Listing(self)
}
pub fn listed(&self) -> Option<usize> {
self.listing
.load(Ordering::Acquire)
.then(|| self.listed.load(Ordering::Relaxed))
}
pub fn begin(&self, total: usize) {
self.read.store(0, Ordering::Relaxed);
self.last_total.store(total, Ordering::Relaxed);
self.total.store(total, Ordering::Release);
self.passes.fetch_add(1, Ordering::Relaxed);
}
pub fn advance(&self) {
self.read.fetch_add(1, Ordering::Relaxed);
}
pub fn done(&self) {
self.total.store(0, Ordering::Relaxed);
}
pub fn pass(&self, total: usize) -> Pass<'_> {
self.begin(total);
Pass(self)
}
#[doc(hidden)]
pub fn last_pass(&self) -> PassCount {
PassCount {
begun: self.passes.load(Ordering::Relaxed),
read: self.read.load(Ordering::Relaxed),
total: self.last_total.load(Ordering::Relaxed),
}
}
pub fn reading(&self) -> Option<(usize, usize)> {
let total = self.total.load(Ordering::Acquire);
(total > 0).then(|| (self.read.load(Ordering::Relaxed).min(total), total))
}
pub fn cancel(&self) {
self.cancelled.store(true, Ordering::Relaxed);
}
pub fn is_cancelled(&self) -> bool {
self.cancelled.load(Ordering::Relaxed)
}
pub fn cancel_flag(&self) -> std::sync::Arc<std::sync::atomic::AtomicBool> {
self.cancelled.clone()
}
pub fn counting() -> Self {
let progress = Self::default();
progress.at_once.store(COUNT_AT_ONCE, Ordering::Relaxed);
progress
}
pub fn reads_at_once(&self) -> usize {
match self.at_once.load(Ordering::Relaxed) {
0 => FOOTERS_AT_ONCE,
n => n,
}
}
pub fn set_estimate(&self, estimate: Option<RowEstimate>) {
*self.estimate.lock().unwrap_or_else(|e| e.into_inner()) = estimate;
}
pub fn estimate(&self) -> Option<RowEstimate> {
*self.estimate.lock().unwrap_or_else(|e| e.into_inner())
}
}
#[derive(Debug, Clone)]
pub struct FileFooter {
pub schema: Arc<Schema>,
pub row_group_rows: Vec<usize>,
pub row_group_bytes: Vec<usize>,
pub file_bytes: usize,
pub column_bytes: Vec<(String, usize)>,
}
impl FileFooter {
pub fn from_metadata(
schema: Schema,
metadata: &polars_parquet::parquet::metadata::FileMetadata,
file_bytes: usize,
widths: bool,
) -> Self {
let column_bytes = if widths {
parquet_column_bytes(&schema, metadata)
} else {
Vec::new()
};
FileFooter {
schema: Arc::new(schema),
row_group_rows: metadata.row_groups.iter().map(|rg| rg.num_rows()).collect(),
row_group_bytes: metadata
.row_groups
.iter()
.map(|rg| rg.compressed_size())
.collect(),
file_bytes,
column_bytes,
}
}
pub fn from_tail(tail: &[u8], file_bytes: usize, widths: bool) -> color_eyre::Result<Self> {
use polars::prelude::{ParquetReader, SchemaExt, SerReader};
let mut cursor = std::io::Cursor::new(tail);
let mut reader = ParquetReader::new(&mut cursor);
let arrow_schema = reader
.schema()
.map_err(|e| color_eyre::eyre::eyre!("Parquet schema read failed: {e}"))?;
let metadata = reader
.get_metadata()
.map_err(|e| color_eyre::eyre::eyre!("Parquet footer read failed: {e}"))?;
Ok(Self::from_metadata(
Schema::from_arrow_schema(arrow_schema.as_ref()),
metadata,
file_bytes,
widths,
))
}
pub fn rows(&self) -> usize {
self.row_group_rows.iter().sum()
}
}
pub fn parquet_column_bytes(
schema: &Schema,
metadata: &polars_parquet::parquet::metadata::FileMetadata,
) -> Vec<(String, usize)> {
schema
.iter_names()
.map(|name| {
let bytes: i64 = metadata
.row_groups
.iter()
.flat_map(|rg| rg.columns_under_root_iter(name).into_iter().flatten())
.map(|chunk| chunk.uncompressed_size())
.sum();
(name.to_string(), bytes.max(0) as usize)
})
.collect()
}
pub fn column_bytes_per_row(footers: &[Option<FileFooter>]) -> Vec<(String, usize)> {
let rows: usize = footers.iter().flatten().map(FileFooter::rows).sum();
if rows == 0 {
return Vec::new();
}
let mut totals: Vec<(String, usize)> = Vec::new();
let mut at: HashMap<String, usize> = HashMap::new();
for (name, bytes) in footers.iter().flatten().flat_map(|f| &f.column_bytes) {
match at.get(name) {
Some(&i) => totals[i].1 += bytes,
None => {
at.insert(name.clone(), totals.len());
totals.push((name.clone(), *bytes));
}
}
}
totals
.into_iter()
.map(|(name, bytes)| (name, bytes / rows))
.collect()
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SchemaOrigin {
AllFooters(usize),
FooterSample { read: usize, total: usize },
}
impl SchemaOrigin {
pub fn total_files(&self) -> usize {
match self {
SchemaOrigin::AllFooters(files) => *files,
SchemaOrigin::FooterSample { total, .. } => *total,
}
}
}
impl std::fmt::Display for SchemaOrigin {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
SchemaOrigin::AllFooters(1) => write!(f, "one footer"),
SchemaOrigin::AllFooters(n) => {
write!(f, "all {} footers", crate::numfmt::group_chrome(*n))
}
SchemaOrigin::FooterSample { read, total } => write!(
f,
"{} of {} footers (sample)",
crate::numfmt::group_chrome(*read),
crate::numfmt::group_chrome(*total)
),
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct ColumnDrift {
pub name: PlSmallStr,
pub dtype: DataType,
pub present_in: usize,
pub conflicting_files: usize,
pub conflicting_types: Vec<DataType>,
pub widened: bool,
}
impl ColumnDrift {
pub fn is_uniform(&self, files: usize) -> bool {
self.present_in == files && self.conflicting_files == 0 && !self.widened
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ColumnRange {
Only(String),
NoneBefore(String),
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct SkippedFiles {
pub bookkeeping: usize,
pub not_parquet: usize,
pub empty: usize,
}
impl SkippedFiles {
pub fn count(&mut self, bookkeeping: bool) {
if bookkeeping {
self.bookkeeping += 1;
} else {
self.not_parquet += 1;
}
}
}
#[derive(Debug, Clone)]
pub struct ReadAs {
pub delimiter: Option<u8>,
pub has_header: Option<bool>,
pub skip_rows: Option<usize>,
pub skip_lines: Option<usize>,
pub infer_schema_length: Option<usize>,
pub ignore_errors: bool,
pub try_parse_dates: bool,
pub comment_char: Option<String>,
pub header_rows: Vec<usize>,
pub header_join: String,
}
impl ReadAs {
fn open_options(&self, format: crate::FileFormat) -> crate::OpenOptions {
crate::OpenOptions {
delimiter: self.delimiter.or(format.separator()),
has_header: self.has_header,
skip_rows: self.skip_rows,
skip_lines: self.skip_lines,
infer_schema_length: self.infer_schema_length,
ignore_errors: self.ignore_errors,
parse_dates: self.try_parse_dates,
parse_strings: None,
comment_char: self.comment_char.clone(),
header_rows: self.header_rows.clone(),
header_join: self.header_join.clone(),
..crate::OpenOptions::default()
}
}
}
impl Default for ReadAs {
fn default() -> Self {
Self {
delimiter: None,
has_header: None,
skip_rows: None,
skip_lines: None,
infer_schema_length: None,
ignore_errors: false,
try_parse_dates: true,
comment_char: None,
header_rows: Vec::new(),
header_join: crate::csv_dialect::DEFAULT_HEADER_JOIN.to_string(),
}
}
}
pub fn column_schema_of(
path: &std::path::Path,
format: crate::FileFormat,
as_read: &ReadAs,
) -> Option<Vec<(String, DataType)>> {
use polars::prelude::{LazyFileListReader, LazyJsonLineReader};
let lf = match format.descriptor().lines {
Some(crate::cli::Lines::Delimited(_)) => {
let options = as_read.open_options(format);
let header = crate::widgets::datatable::DataTableState::csv_header_names_of(
&options, path, None,
)
.ok()?;
let reader = crate::widgets::datatable::DataTableState::configure_csv_reader(
crate::widgets::datatable::DataTableState::csv_reader_of(path).ok()?,
&options,
None,
);
crate::csv_dialect::name_columns(reader.finish().ok()?, header.as_deref()).ok()?
}
Some(crate::cli::Lines::Json) => {
LazyJsonLineReader::new(crate::source::polars_literal_path(path).ok()?)
.finish()
.ok()?
}
Some(crate::cli::Lines::Text) => {
return Some(
crate::lines::schema(false)
.iter()
.map(|(name, dtype)| (name.to_string(), dtype.clone()))
.collect(),
);
}
None => return None,
};
let schema = lf.clone().collect_schema().ok()?;
let fields: Vec<(String, DataType)> = schema
.iter()
.map(|(name, dtype)| (name.trim().to_string(), dtype.clone()))
.collect();
Some(fields)
}
fn is_an_empty_file(schema: &[(String, DataType)]) -> bool {
match schema {
[] => true,
[(only, _)] => only.trim().is_empty(),
_ => false,
}
}
pub(crate) fn names_are_names(names: &[String]) -> bool {
!names.is_empty() && !names.iter().all(|n| n.trim().parse::<f64>().is_ok())
}
#[derive(Debug, Clone, Default)]
pub struct Sampled {
pub columns: Vec<String>,
pub nests: Option<bool>,
pub columns_differ: bool,
pub types_differ: bool,
pub read: usize,
pub headerless: bool,
}
impl Sampled {
pub fn disagreement(&self) -> Disagreement {
if self.headerless {
return Disagreement {
headerless: true,
..Default::default()
};
}
Disagreement {
columns: self.columns_differ,
types: self.types_differ,
headerless: false,
}
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct Disagreement {
pub columns: bool,
pub types: bool,
pub headerless: bool,
}
impl Disagreement {
pub fn any(&self) -> bool {
self.columns || self.types || self.headerless
}
}
pub fn sample_files(
files: &[std::path::PathBuf],
format: crate::FileFormat,
as_read: &ReadAs,
) -> Sampled {
const WANTED: usize = 3;
const TRIES: usize = 12;
let last = files.len().saturating_sub(1);
const NEAR: usize = 4;
let anchors = [0usize, last / 2, last];
let mut read: Vec<Vec<(String, DataType)>> = Vec::new();
let mut tried = 0usize;
let mut seen: Vec<usize> = Vec::new();
'anchors: for anchor in anchors {
for step in 0..NEAR {
if read.len() >= WANTED || tried >= TRIES {
break 'anchors;
}
let i = anchor + step;
if i > last || seen.contains(&i) {
continue;
}
seen.push(i);
let Some(file) = files.get(i) else { continue };
tried += 1;
if let Some(schema) = column_schema_of(file, format, as_read)
&& !is_an_empty_file(&schema)
{
read.push(schema);
continue 'anchors;
}
}
}
let mut out = Sampled {
read: read.len(),
..Default::default()
};
for file in &read {
for (name, _) in file {
if !out.columns.iter().any(|c| c == name) {
out.columns.push(name.clone());
}
}
}
if read.len() < 2 {
return out;
}
let names: Vec<Vec<String>> = read
.iter()
.map(|f| f.iter().map(|(n, _)| n.clone()).collect())
.collect();
if names.iter().any(|f| !names_are_names(f)) {
out.columns.clear();
out.headerless = true;
out.nests = Some(false);
return out;
}
let nests = is_nested(&names);
out.nests = Some(nests);
let mut types: HashMap<&str, &DataType> = HashMap::new();
let mut typed_apart = false;
for (name, dtype) in read.iter().flatten() {
match types.get(name.as_str()) {
Some(seen) if *seen != dtype => typed_apart = true,
Some(_) => {}
None => {
types.insert(name.as_str(), dtype);
}
}
}
let widest = names.iter().map(|f| f.len()).max().unwrap_or(0);
out.columns_differ = names.iter().any(|f| f.len() != widest) || !nests;
out.types_differ = typed_apart;
out
}
pub fn is_nested(files: &[Vec<String>]) -> bool {
let Some(widest) = files.iter().max_by_key(|f| f.len()) else {
return true;
};
let widest: std::collections::BTreeSet<&str> = widest.iter().map(String::as_str).collect();
files
.iter()
.all(|file| file.iter().all(|name| widest.contains(name.as_str())))
}
pub fn top_level_columns(leaves: &[String]) -> Vec<String> {
let mut seen = std::collections::HashSet::new();
leaves
.iter()
.map(|leaf| leaf.split_once('.').map_or(leaf.as_str(), |(root, _)| root))
.filter(|root| seen.insert(root.to_string()))
.map(str::to_string)
.collect()
}
#[derive(Debug, Clone)]
pub struct DatasetSchema {
pub schema: Arc<Schema>,
pub columns: Vec<ColumnDrift>,
pub omitted: Vec<Vec<(PlSmallStr, DataType)>>,
pub unreadable: Vec<usize>,
pub files: usize,
pub groups: Vec<DriftGroup>,
pub file_group: Vec<u32>,
pub origin: SchemaOrigin,
pub read_as_text: Vec<PlSmallStr>,
pub empty_files: usize,
pub median_row_group_bytes: Option<usize>,
pub median_file_bytes: Option<usize>,
pub column_ranges: HashMap<PlSmallStr, ColumnRange>,
pub partition_layouts: Vec<(Vec<String>, usize)>,
pub partition_layouts_dropped: (usize, usize),
pub skipped: SkippedFiles,
pub listed_files: usize,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Hash)]
pub struct DriftGroup {
pub absent: Vec<PlSmallStr>,
pub unread: Vec<PlSmallStr>,
}
impl DriftGroup {
pub fn is_empty(&self) -> bool {
self.absent.is_empty() && self.unread.is_empty()
}
}
impl DatasetSchema {
pub fn drifting(&self) -> impl Iterator<Item = &ColumnDrift> {
let readable = self.files - self.unreadable.len();
self.columns.iter().filter(move |c| !c.is_uniform(readable))
}
pub fn with_partition_layouts(mut self, root: &str, paths: &[String]) -> DatasetSchema {
const KEPT: usize = 64;
let mut counts: HashMap<Vec<String>, usize> = HashMap::new();
for path in paths {
let Some(below) = path.strip_prefix(root) else {
continue;
};
let keys = partition_keys_of(below);
if keys.is_empty() {
continue;
}
*counts.entry(keys).or_insert(0) += 1;
}
let mut counts: Vec<(Vec<String>, usize)> = counts.into_iter().collect();
counts.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
let dropped = &counts[counts.len().min(KEPT)..];
self.partition_layouts_dropped =
(dropped.len(), dropped.iter().map(|(_, files)| files).sum());
counts.truncate(KEPT);
self.partition_layouts = counts;
self.listed_files = paths.len();
self.column_ranges = self.ranges_of_columns(root, paths);
self
}
fn ranges_of_columns(&self, root: &str, paths: &[String]) -> HashMap<PlSmallStr, ColumnRange> {
if matches!(self.origin, SchemaOrigin::FooterSample { .. })
|| !self.unreadable.is_empty()
|| self.file_group.len() != paths.len()
{
return HashMap::new();
}
struct Seen {
first_present: Option<usize>,
last_absent: Option<usize>,
with: Vec<String>,
without: Vec<String>,
unplaced: bool,
}
let partition_of = |index: usize| -> Option<String> {
let below = paths.get(index)?.strip_prefix(root)?;
let values = partition_values_of(below);
(!values.is_empty()).then(|| values.join("/"))
};
let reads_in_order = (0..paths.len())
.filter_map(&partition_of)
.collect::<Vec<_>>()
.windows(2)
.all(|pair| natural_cmp(&pair[0], &pair[1]) != std::cmp::Ordering::Greater);
let drifting: Vec<&ColumnDrift> = self
.columns
.iter()
.filter(|column| column.present_in > 0 && column.present_in < self.files)
.collect();
if drifting.is_empty() {
return HashMap::new();
}
let where_in_drifting: HashMap<&PlSmallStr, usize> = drifting
.iter()
.enumerate()
.map(|(at, column)| (&column.name, at))
.collect();
let missing_by_group: Vec<Vec<bool>> = self
.groups
.iter()
.map(|group| {
let mut missing = vec![false; drifting.len()];
for name in &group.absent {
if let Some(at) = where_in_drifting.get(name) {
missing[*at] = true;
}
}
missing
})
.collect();
let none_missing: Vec<bool> = vec![false; drifting.len()];
let mut seen: Vec<Seen> = (0..drifting.len())
.map(|_| Seen {
first_present: None,
last_absent: None,
with: Vec::new(),
without: Vec::new(),
unplaced: false,
})
.collect();
for (index, group) in self.file_group.iter().enumerate() {
let missing: &[bool] = missing_by_group
.get(*group as usize)
.map(Vec::as_slice)
.unwrap_or(&none_missing);
let here = partition_of(index);
for (at, absent) in missing.iter().enumerate() {
let entry = &mut seen[at];
let seen_of = if *absent {
entry.last_absent = Some(index);
&mut entry.without
} else {
entry.first_present.get_or_insert(index);
entry.unplaced |= here.is_none();
&mut entry.with
};
if seen_of.len() < 2
&& let Some(partition) = here.as_ref()
&& !seen_of.contains(partition)
{
seen_of.push(partition.clone());
}
}
}
seen.into_iter()
.zip(&drifting)
.filter_map(|(entry, column)| {
let name = column.name.clone();
let first = entry.first_present?;
if !entry.unplaced
&& entry.with.len() == 1
&& entry
.without
.iter()
.any(|other| {
!partition_holds(&entry.with[0], other)
&& !same_place(&entry.with[0], other)
})
{
return Some((name, ColumnRange::Only(entry.with[0].clone())));
}
let last_absent = entry.last_absent?;
if last_absent > first || !reads_in_order {
return None;
}
let (ends, begins) = (partition_of(last_absent)?, partition_of(first)?);
(!same_place(&ends, &begins)
&& !partition_holds(&begins, &ends)
&& !partition_holds(&ends, &begins))
.then_some((name, ColumnRange::NoneBefore(begins)))
})
.collect()
}
pub fn with_skipped(mut self, skipped: SkippedFiles) -> DatasetSchema {
self.skipped = skipped;
self
}
pub fn reading_as_text(&self, as_text: &[PlSmallStr]) -> DatasetSchema {
let mut out = self.clone();
if as_text.is_empty() {
return out;
}
out.schema = crate::schema_union::text_schema(&self.schema, as_text);
for column in &mut out.columns {
if as_text.contains(&column.name) {
column.dtype = DataType::String;
column.conflicting_files = 0;
column.conflicting_types.clear();
}
}
for group in &mut out.groups {
group.unread.retain(|name| !as_text.contains(name));
}
out.read_as_text = as_text.to_vec();
out
}
pub fn drifts(&self) -> bool {
self.groups.iter().any(|g| !g.is_empty())
}
}
pub const FOOTERS_AT_ONCE: usize = 64;
pub fn footers_to_cache(
footers: &[Option<FileFooter>],
) -> (
Vec<crate::cache::CachedFooter>,
Vec<Vec<(String, DataType)>>,
) {
let mut schemas = Vec::new();
let cached = footers
.iter()
.map(|footer| match footer {
None => crate::cache::CachedFooter::default(),
Some(f) => crate::cache::CachedFooter {
schema: Some(crate::cache::DatasetShape::intern_schema(
&mut schemas,
&f.schema,
)),
row_group_rows: f.row_group_rows.clone(),
row_group_bytes: f.row_group_bytes.clone(),
column_bytes: f.column_bytes.iter().map(|(_, bytes)| *bytes).collect(),
},
})
.collect();
(cached, schemas)
}
pub fn footers_from_cache(
cached: &[crate::cache::CachedFooter],
schemas: &[Vec<(String, DataType)>],
file_bytes: &[u64],
) -> Option<Vec<Option<FileFooter>>> {
if cached.len() != file_bytes.len() {
return None;
}
let shared: Vec<Arc<Schema>> = (0..schemas.len())
.map(|at| crate::cache::DatasetShape::schema_at(schemas, at).map(Arc::new))
.collect::<Option<_>>()?;
cached
.iter()
.zip(file_bytes)
.map(|(f, &bytes)| {
let Some(at) = f.schema else {
return Some(None);
};
let schema = shared.get(at)?;
let column_bytes = schema
.iter_names()
.zip(&f.column_bytes)
.map(|(name, bytes)| (name.to_string(), *bytes))
.collect();
Some(Some(FileFooter {
schema: schema.clone(),
row_group_rows: f.row_group_rows.clone(),
row_group_bytes: f.row_group_bytes.clone(),
file_bytes: bytes as usize,
column_bytes,
}))
})
.collect()
}
type FooterHook = Arc<dyn Fn(&std::path::Path) + Send + Sync>;
static FOOTER_HOOKS: std::sync::Mutex<Vec<(u64, std::path::PathBuf, FooterHook)>> =
std::sync::Mutex::new(Vec::new());
static FOOTER_HOOKS_SET: AtomicUsize = AtomicUsize::new(0);
#[doc(hidden)]
pub fn on_local_footer_read(
dir: &std::path::Path,
hook: impl Fn(&std::path::Path) + Send + Sync + 'static,
) -> FooterHookGuard {
static NEXT: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let id = NEXT.fetch_add(1, Ordering::Relaxed);
FOOTER_HOOKS
.lock()
.unwrap_or_else(|e| e.into_inner())
.push((id, dir.to_path_buf(), Arc::new(hook)));
FOOTER_HOOKS_SET.fetch_add(1, Ordering::Release);
FooterHookGuard(id)
}
#[doc(hidden)]
pub struct FooterHookGuard(u64);
impl Drop for FooterHookGuard {
fn drop(&mut self) {
FOOTER_HOOKS
.lock()
.unwrap_or_else(|e| e.into_inner())
.retain(|(id, _, _)| *id != self.0);
FOOTER_HOOKS_SET.fetch_sub(1, Ordering::Release);
}
}
pub(crate) fn before_local_footer_read(path: &std::path::Path) {
if FOOTER_HOOKS_SET.load(Ordering::Acquire) == 0 {
return;
}
let hooks: Vec<FooterHook> = FOOTER_HOOKS
.lock()
.unwrap_or_else(|e| e.into_inner())
.iter()
.filter(|(_, dir, _)| path.starts_with(dir))
.map(|(_, _, hook)| hook.clone())
.collect();
for hook in hooks {
hook(path);
}
}
pub const MAX_FOOTER_READS: usize = 20_000;
pub fn footers_to_read(files: usize) -> Vec<usize> {
if files <= MAX_FOOTER_READS {
return (0..files).collect();
}
let last = files - 1;
let mut sample: Vec<usize> = (0..MAX_FOOTER_READS)
.map(|i| i * last / (MAX_FOOTER_READS - 1))
.collect();
sample.dedup();
sample
}
pub fn ends_of(files: usize) -> Vec<usize> {
match files {
0 => Vec::new(),
1 => vec![0],
n => vec![0, n - 1],
}
}
pub fn union_sampled(
files: usize,
read: &[usize],
footers: &[Option<FileFooter>],
) -> DatasetSchema {
let origin = if read.len() == files {
SchemaOrigin::AllFooters(files)
} else {
SchemaOrigin::FooterSample {
read: read.len(),
total: files,
}
};
let mut union = union_file_schemas(footers, origin);
let mut omitted = vec![Vec::new(); files];
let mut file_group = vec![0u32; files];
for ((columns, group), &index) in union.omitted.iter().zip(union.file_group.iter()).zip(read) {
omitted[index] = columns.clone();
file_group[index] = *group;
}
union.omitted = omitted;
union.file_group = file_group;
union.unreadable = union
.unreadable
.iter()
.filter_map(|i| read.get(*i).copied())
.collect();
union
}
pub fn readable_paths<'a>(paths: &'a [String], unreadable: &[usize]) -> Cow<'a, [String]> {
if unreadable.is_empty() {
return Cow::Borrowed(paths);
}
debug_assert!(unreadable.windows(2).all(|pair| pair[0] < pair[1]));
Cow::Owned(
paths
.iter()
.enumerate()
.filter(|(index, _)| unreadable.binary_search(index).is_err())
.map(|(_, path)| path.clone())
.collect(),
)
}
pub fn union_file_schemas(files: &[Option<FileFooter>], origin: SchemaOrigin) -> DatasetSchema {
let unreadable = files
.iter()
.enumerate()
.filter_map(|(i, f)| f.is_none().then_some(i))
.collect();
let mut order: Vec<PlSmallStr> = Vec::new();
let mut seen: HashMap<PlSmallStr, usize> = HashMap::new();
let mut push = |name: &PlSmallStr, order: &mut Vec<PlSmallStr>| {
if !seen.contains_key(name) {
seen.insert(name.clone(), order.len());
order.push(name.clone());
}
};
if let Some(newest) = files.iter().rev().flatten().next() {
for name in newest.schema.iter_names() {
push(name, &mut order);
}
}
for file in files.iter().flatten() {
for name in file.schema.iter_names() {
push(name, &mut order);
}
}
let mut sightings: Vec<Vec<(DataType, usize)>> = vec![Vec::new(); order.len()];
for file in files.iter().flatten() {
for (name, dtype) in file.schema.iter() {
let Some(&index) = seen.get(name) else {
continue;
};
sightings[index].push((dtype.clone(), file.rows()));
}
}
let mut schema = Schema::with_capacity(order.len());
let mut columns = Vec::with_capacity(order.len());
for (name, seen_types) in order.iter().zip(sightings.iter()) {
let chosen = choose_dtype(seen_types);
let conflicting_types = seen_types
.iter()
.map(|(d, _)| d)
.filter(|d| !fits(d, &chosen))
.fold(Vec::new(), |mut acc: Vec<DataType>, d| {
if !acc.contains(d) {
acc.push(d.clone());
}
acc
});
columns.push(ColumnDrift {
name: name.clone(),
present_in: seen_types.len(),
conflicting_files: seen_types.iter().filter(|(d, _)| !fits(d, &chosen)).count(),
widened: seen_types
.iter()
.any(|(d, _)| *d != chosen && fits(d, &chosen)),
conflicting_types,
dtype: chosen.clone(),
});
schema.with_column(name.clone(), chosen);
}
let mut groups: Vec<DriftGroup> = vec![DriftGroup::default()];
let mut group_of: HashMap<DriftGroup, u32> = HashMap::from([(DriftGroup::default(), 0)]);
let mut file_group = Vec::with_capacity(files.len());
let mut omitted = Vec::with_capacity(files.len());
for file in files {
let Some(file) = file else {
file_group.push(0);
omitted.push(Vec::new());
continue;
};
let stored: Vec<(PlSmallStr, DataType)> = file
.schema
.iter()
.filter(|(name, dtype)| schema.get(name).is_some_and(|target| !fits(dtype, target)))
.map(|(name, dtype)| (name.clone(), dtype.clone()))
.collect();
let unread: Vec<PlSmallStr> = stored.iter().map(|(name, _)| name.clone()).collect();
let absent: Vec<PlSmallStr> = schema
.iter_names()
.filter(|name| !file.schema.contains(name))
.cloned()
.collect();
omitted.push(stored);
let group = DriftGroup { absent, unread };
let next = groups.len() as u32;
let id = *group_of.entry(group.clone()).or_insert_with(|| {
groups.push(group);
next
});
file_group.push(id);
}
DatasetSchema {
schema: Arc::new(schema),
columns,
omitted,
unreadable,
files: files.len(),
groups,
file_group,
origin,
read_as_text: Vec::new(),
empty_files: files.iter().flatten().filter(|f| f.rows() == 0).count(),
median_file_bytes: median(files.iter().flatten().map(|f| f.file_bytes)),
column_ranges: HashMap::new(),
skipped: SkippedFiles::default(),
partition_layouts: Vec::new(),
partition_layouts_dropped: (0, 0),
listed_files: 0,
median_row_group_bytes: median(
files
.iter()
.flatten()
.flat_map(|f| f.row_group_bytes.iter().copied()),
),
}
}
fn partition_holds(outer: &str, inner: &str) -> bool {
inner == outer
|| inner
.strip_prefix(outer)
.is_some_and(|rest| rest.starts_with('/'))
}
fn same_place(a: &str, b: &str) -> bool {
fn sorted(path: &str) -> Vec<&str> {
let mut segments: Vec<&str> = path.split('/').collect();
segments.sort_by_key(|segment| segment.split_once('=').map(|(key, _)| key));
segments
}
let (a, b) = (sorted(a), sorted(b));
a.len() == b.len()
&& a.iter()
.zip(&b)
.all(|(x, y)| natural_cmp(x, y) == std::cmp::Ordering::Equal)
}
fn natural_cmp(a: &str, b: &str) -> std::cmp::Ordering {
use std::cmp::Ordering;
let (mut a, mut b) = (a.as_bytes(), b.as_bytes());
loop {
match (a.first(), b.first()) {
(None, None) => return Ordering::Equal,
(None, _) => return Ordering::Less,
(_, None) => return Ordering::Greater,
(Some(x), Some(y)) if x.is_ascii_digit() && y.is_ascii_digit() => {
let digits = |s: &[u8]| s.iter().take_while(|c| c.is_ascii_digit()).count();
let (na, nb) = (digits(a), digits(b));
let (xs, ys) = (&a[..na], &b[..nb]);
fn trim(s: &[u8]) -> &[u8] {
let lead = s.iter().take_while(|c| **c == b'0').count();
&s[lead.min(s.len().saturating_sub(1))..]
}
let (tx, ty) = (trim(xs), trim(ys));
match tx.len().cmp(&ty.len()).then_with(|| tx.cmp(ty)) {
Ordering::Equal => {}
other => return other,
}
a = &a[na..];
b = &b[nb..];
}
(Some(x), Some(y)) => match x.cmp(y) {
Ordering::Equal => {
a = &a[1..];
b = &b[1..];
}
other => return other,
},
}
}
}
fn partition_values_of(path: &str) -> Vec<String> {
#[cfg(windows)]
let separators: &[char] = &['/', '\\'];
#[cfg(not(windows))]
let separators: &[char] = &['/'];
let mut segments: Vec<&str> = path.split(separators).collect();
segments.pop();
segments
.into_iter()
.filter(|segment| {
segment
.split_once('=')
.is_some_and(|(key, _)| !key.is_empty())
})
.map(|segment| segment.to_string())
.collect()
}
fn partition_keys_of(path: &str) -> Vec<String> {
let mut keys: Vec<String> = Vec::new();
#[cfg(windows)]
let separators: &[char] = &['/', '\\'];
#[cfg(not(windows))]
let separators: &[char] = &['/'];
let mut segments: Vec<&str> = path.split(separators).collect();
segments.pop();
for segment in segments {
if let Some((key, _)) = segment.split_once('=')
&& !key.is_empty()
{
keys.push(key.to_string());
}
}
keys.sort();
keys.dedup();
keys
}
fn median(sizes: impl Iterator<Item = usize>) -> Option<usize> {
let mut sizes: Vec<usize> = sizes.collect();
if sizes.is_empty() {
return None;
}
sizes.sort_unstable();
Some(sizes[(sizes.len() - 1) / 2])
}
fn choose_dtype(seen: &[(DataType, usize)]) -> DataType {
let mut distinct: Vec<DataType> = Vec::new();
for (dtype, _) in seen {
if !distinct.contains(dtype) {
distinct.push(dtype.clone());
}
}
match distinct.as_slice() {
[] => return DataType::Null,
[only] => return only.clone(),
_ => {}
}
let mut candidates = distinct.clone();
for dtype in &distinct {
let folded = distinct
.iter()
.filter(|other| widen(dtype, other).is_some())
.try_fold(dtype.clone(), |acc, other| widen(&acc, other));
if let Some(folded) = folded
&& !candidates.contains(&folded)
{
candidates.push(folded);
}
}
let mut best: Option<(DataType, usize, usize)> = None;
for candidate in candidates {
let rows: usize = seen
.iter()
.filter(|(d, _)| fits(d, &candidate))
.map(|(_, rows)| rows)
.sum();
let files = seen.iter().filter(|(d, _)| fits(d, &candidate)).count();
let better = best
.as_ref()
.is_none_or(|(_, best_rows, best_files)| (rows, files) > (*best_rows, *best_files));
if better {
best = Some((candidate, rows, files));
}
}
best.map(|(d, _, _)| d).unwrap_or(DataType::Null)
}
pub fn fits(from: &DataType, to: &DataType) -> bool {
widen(from, to).as_ref() == Some(to)
}
pub fn widen(a: &DataType, b: &DataType) -> Option<DataType> {
use DataType::*;
if a == b {
return Some(a.clone());
}
match (a, b) {
(Null, other) | (other, Null) => Some(other.clone()),
_ if a.is_integer() && b.is_integer() => widen_integers(a, b),
_ if (a.is_integer() || a.is_float()) && (b.is_integer() || b.is_float()) => Some(Float64),
(Datetime(a_unit, a_zone), Datetime(b_unit, b_zone)) if a_zone == b_zone => {
Some(Datetime(finer_unit(*a_unit, *b_unit), a_zone.clone()))
}
(List(a_inner), List(b_inner)) => widen(a_inner, b_inner).map(|t| List(Box::new(t))),
(Struct(a_fields), Struct(b_fields)) => widen_structs(a_fields, b_fields),
_ => None,
}
}
fn widen_integers(a: &DataType, b: &DataType) -> Option<DataType> {
use DataType::*;
let signed = |d: &DataType| matches!(d, Int8 | Int16 | Int32 | Int64 | Int128);
let bits = |d: &DataType| match d {
Int8 | UInt8 => 8u32,
Int16 | UInt16 => 16,
Int32 | UInt32 => 32,
Int64 | UInt64 => 64,
_ => 128,
};
if signed(a) == signed(b) {
let wider = if bits(a) >= bits(b) { a } else { b };
return Some(wider.clone());
}
let (unsigned, sgn) = if signed(a) { (b, a) } else { (a, b) };
let needed = match bits(unsigned) {
8 => Int16,
16 => Int32,
32 => Int64,
64 => return None,
_ => return None,
};
Some(if bits(sgn) >= bits(&needed) {
sgn.clone()
} else {
needed
})
}
fn widen_structs(a: &[Field], b: &[Field]) -> Option<DataType> {
let mut fields: Vec<Field> = Vec::with_capacity(a.len() + b.len());
for field in a {
let widened = match b.iter().find(|other| other.name() == field.name()) {
Some(other) => widen(field.dtype(), other.dtype())?,
None => field.dtype().clone(),
};
fields.push(Field::new(field.name().clone(), widened));
}
for field in b {
if !a.iter().any(|other| other.name() == field.name()) {
fields.push(field.clone());
}
}
Some(DataType::Struct(fields))
}
fn finer_unit(a: TimeUnit, b: TimeUnit) -> TimeUnit {
let rank = |u: TimeUnit| match u {
TimeUnit::Milliseconds => 0,
TimeUnit::Microseconds => 1,
TimeUnit::Nanoseconds => 2,
};
if rank(a) >= rank(b) { a } else { b }
}
pub fn with_partition_columns(
file_schema: &Schema,
partition_columns: &[String],
values: &[(String, String)],
) -> Schema {
let part_set: HashSet<&str> = partition_columns.iter().map(String::as_str).collect();
let mut merged = Schema::with_capacity(partition_columns.len() + file_schema.len());
for name in partition_columns {
merged.with_column(
name.clone().into(),
crate::widgets::datatable::partition_dtype(name, file_schema, values),
);
}
for (name, dtype) in file_schema.iter() {
if !part_set.contains(name.as_str()) {
merged.with_column(name.clone(), dtype.clone());
}
}
merged
}
pub fn partition_columns_of_key(key: &str) -> Vec<String> {
let mut columns = Vec::new();
let mut seen = HashSet::new();
for segment in key.split('/') {
if let Some((name, _)) = segment.split_once('=')
&& !name.is_empty()
&& seen.insert(name.to_string())
{
columns.push(name.to_string());
}
}
columns
}
pub fn partitions_of_listing(first: &str, newest: &str) -> (Vec<String>, Vec<(String, String)>) {
let values = [first, newest]
.iter()
.flat_map(|key| key.split('/'))
.filter_map(|segment| segment.split_once('='))
.map(|(k, v)| (k.to_string(), v.to_string()))
.collect();
(partition_columns_of_key(newest), values)
}
pub struct FooterCount<F> {
files: usize,
counted: Vec<usize>,
footers: std::sync::Mutex<Vec<Option<F>>>,
}
pub struct Counted<F> {
pub row_groups: Vec<Vec<usize>>,
pub whole: Option<Vec<Option<F>>>,
}
impl<F: Clone> FooterCount<F> {
pub fn new(
files: usize,
counted: Vec<usize>,
known: impl IntoIterator<Item = (usize, Option<F>)>,
) -> Self {
let mut footers = vec![None; files];
for (index, footer) in known {
if let Some(slot) = footers.get_mut(index) {
*slot = footer;
}
}
Self {
files,
counted,
footers: std::sync::Mutex::new(footers),
}
}
pub fn count(
&self,
read: impl FnOnce(&[usize]) -> Option<Vec<Option<F>>>,
row_groups: impl Fn(&F) -> Vec<usize>,
) -> Option<Counted<F>> {
let mut footers = self.footers.lock().unwrap_or_else(|e| e.into_inner());
let missing = self.missing(&mut footers);
let read = if missing.is_empty() {
Vec::new()
} else {
read(&missing)?
};
Some(self.settle(&mut footers, missing, read, row_groups))
}
fn missing(&self, footers: &mut Vec<Option<F>>) -> Vec<usize> {
if footers.is_empty() {
*footers = vec![None; self.files];
}
self.counted
.iter()
.copied()
.filter(|&index| footers[index].is_none())
.collect()
}
fn settle(
&self,
footers: &mut Vec<Option<F>>,
missing: Vec<usize>,
read: Vec<Option<F>>,
row_groups: impl Fn(&F) -> Vec<usize>,
) -> Counted<F> {
if footers.is_empty() {
*footers = vec![None; self.files];
}
for (index, footer) in missing.into_iter().zip(read) {
footers[index] = footer;
}
let groups: Vec<Vec<usize>> = self
.counted
.iter()
.map(|&index| footers[index].as_ref().map(&row_groups).unwrap_or_default())
.collect();
let whole = if footers.iter().all(Option::is_some) {
Some(std::mem::take(footers))
} else {
if self.counted.iter().all(|&index| footers[index].is_some()) {
footers.clear();
}
None
};
Counted {
row_groups: groups,
whole,
}
}
}
pub const DRIFT_COLUMN: &str = "__datui_row";
static NOTHING_MISSING: DriftGroup = DriftGroup {
absent: Vec::new(),
unread: Vec::new(),
};
#[derive(Debug, Clone, Default)]
pub struct ScanDrift {
group_of: HashMap<String, u32>,
row_of: HashMap<String, usize>,
stored_of: HashMap<String, Vec<(PlSmallStr, DataType)>>,
pub groups: Vec<DriftGroup>,
}
impl ScanDrift {
pub fn new(paths: &[String], dataset: &DatasetSchema, file_rows: &[usize]) -> Option<Self> {
if !dataset.drifts() || file_rows.len() != paths.len() {
return None;
}
let group_of = paths
.iter()
.zip(dataset.file_group.iter())
.filter(|(_, group)| **group != 0)
.map(|(path, group)| (path.clone(), *group))
.collect();
let mut row = 0usize;
let mut row_of = HashMap::with_capacity(paths.len());
for (path, rows) in paths.iter().zip(file_rows) {
row_of.insert(path.clone(), row);
row += rows;
}
let stored_of = paths
.iter()
.zip(dataset.omitted.iter())
.filter(|(_, stored)| !stored.is_empty())
.map(|(path, stored)| (path.clone(), stored.clone()))
.collect();
Some(ScanDrift {
group_of,
row_of,
stored_of,
groups: dataset.groups.clone(),
})
}
pub fn group(&self, path: &str) -> u32 {
self.group_of.get(path).copied().unwrap_or(0)
}
fn first_row(&self, path: &str) -> usize {
self.row_of.get(path).copied().unwrap_or(0)
}
fn stored_type(&self, path: &str, column: &PlSmallStr) -> Option<&DataType> {
self.stored_of
.get(path)?
.iter()
.find(|(name, _)| name == column)
.map(|(_, dtype)| dtype)
}
fn unread(&self, path: &str) -> &[PlSmallStr] {
self.groups
.get(self.group(path) as usize)
.unwrap_or(&NOTHING_MISSING)
.unread
.as_slice()
}
}
pub fn lenient_scan(
paths: &[String],
schema: Arc<Schema>,
cloud_options: Option<polars::io::cloud::CloudOptions>,
drift: Option<&ScanDrift>,
as_text: &[PlSmallStr],
) -> PolarsResult<LazyFrame> {
let Some(drift) = drift else {
return scan_run(paths, &schema, cloud_options, &[], None, &[]);
};
let as_text: Vec<PlSmallStr> = as_text
.iter()
.filter(|name| {
schema.get(name).is_some_and(can_read_as_text)
&& paths
.iter()
.all(|path| drift.stored_type(path, name).is_none_or(can_read_as_text))
})
.cloned()
.collect();
let as_text = as_text.as_slice();
let unread_of = |path: &str| -> Vec<PlSmallStr> {
drift
.unread(path)
.iter()
.filter(|name| !as_text.contains(name))
.cloned()
.collect()
};
let key_of = |path: &str| -> (Vec<PlSmallStr>, Vec<Option<DataType>>) {
(
unread_of(path),
as_text
.iter()
.map(|name| drift.stored_type(path, name).cloned())
.collect(),
)
};
let mut runs: Vec<LazyFrame> = Vec::new();
let mut start = 0;
while start < paths.len() {
let key = key_of(&paths[start]);
let end = paths[start..]
.iter()
.position(|path| key_of(path) != key)
.map_or(paths.len(), |offset| start + offset);
let (omit, stored) = key;
let read_as: Vec<(PlSmallStr, DataType)> = as_text
.iter()
.zip(stored)
.map(|(name, stored)| {
let dtype = stored.or_else(|| schema.get(name).cloned());
(name.clone(), dtype.unwrap_or(DataType::String))
})
.collect();
runs.push(scan_run(
&paths[start..end],
&schema,
cloud_options.clone(),
&omit,
Some(drift.first_row(&paths[start])),
&read_as,
)?);
start = end;
}
match runs.len() {
1 => Ok(runs.remove(0)),
_ => concat(
runs,
UnionArgs {
rechunk: false,
parallel: true,
..Default::default()
},
),
}
}
pub fn can_read_as_text(dtype: &DataType) -> bool {
match dtype {
DataType::Binary | DataType::BinaryOffset => false,
DataType::Duration(_) => false,
DataType::List(_) | DataType::Array(_, _) => false,
DataType::Struct(_) => true,
DataType::Unknown(_) => false,
_ => true,
}
}
impl ColumnDrift {
pub fn can_read_as_text(&self) -> bool {
can_read_as_text(&self.dtype) && self.conflicting_types.iter().all(can_read_as_text)
}
}
pub fn text_schema(schema: &Arc<Schema>, as_text: &[PlSmallStr]) -> Arc<Schema> {
if as_text.is_empty() {
return schema.clone();
}
let mut out = Schema::with_capacity(schema.len());
for (name, dtype) in schema.iter() {
let dtype = if as_text.contains(name) {
DataType::String
} else {
dtype.clone()
};
out.with_column(name.clone(), dtype);
}
Arc::new(out)
}
fn scan_run(
urls: &[String],
schema: &Arc<Schema>,
cloud_options: Option<polars::io::cloud::CloudOptions>,
omit: &[PlSmallStr],
first_row: Option<usize>,
read_as: &[(PlSmallStr, DataType)],
) -> PolarsResult<LazyFrame> {
use polars::lazy::dsl::{
CastColumnsPolicy, DslBuilder, ExtraColumnsPolicy, MissingColumnsPolicy, ScanSources,
UnifiedScanArgs,
};
use polars::prelude::{Expr, NULL, col, lit};
let sources = ScanSources::Paths(
urls.iter()
.map(|url| PlRefPath::new(url.as_str()))
.collect(),
);
let target = if omit.is_empty() && read_as.is_empty() {
schema.clone()
} else {
let mut reduced = Schema::with_capacity(schema.len());
for (name, dtype) in schema.iter() {
if omit.contains(name) {
continue;
}
let dtype = read_as
.iter()
.find(|(column, _)| column == name)
.map(|(_, dtype)| dtype)
.unwrap_or(dtype);
reduced.with_column(name.clone(), dtype.clone());
}
Arc::new(reduced)
};
let options = polars::io::parquet::read::ParquetOptions {
schema: Some(target),
..Default::default()
};
let args = UnifiedScanArgs {
cloud_options,
hive_options: polars::io::HiveOptions::new_enabled(),
glob: false,
cast_columns_policy: CastColumnsPolicy {
integer_upcast: true,
integer_to_float_cast: true,
float_upcast: true,
datetime_nanoseconds_downcast: true,
datetime_microseconds_downcast: true,
datetime_milliseconds_upcast: true,
datetime_microseconds_upcast: true,
null_upcast: true,
missing_struct_fields: MissingColumnsPolicy::Insert,
extra_struct_fields: ExtraColumnsPolicy::Ignore,
..CastColumnsPolicy::ERROR_ON_MISMATCH
},
missing_columns_policy: MissingColumnsPolicy::Insert,
extra_columns_policy: ExtraColumnsPolicy::Ignore,
row_index: first_row.map(|first| polars::io::RowIndex {
name: DRIFT_COLUMN.into(),
offset: first as polars::prelude::IdxSize,
}),
..Default::default()
};
let mut lf: LazyFrame = DslBuilder::scan_parquet(sources, options, args)?
.build()
.into();
if !omit.is_empty() {
let nulls: Vec<Expr> = omit
.iter()
.filter_map(|name| {
let dtype = schema.get(name)?;
Some(lit(NULL).cast(dtype.clone()).alias(name.clone()))
})
.collect();
lf = lf.with_columns(nulls);
}
if !read_as.is_empty() {
let texts: Vec<Expr> = read_as
.iter()
.map(|(name, _)| {
crate::past_calendar::text_expr(col(name.clone()), CastOptions::NonStrict)
.alias(name.clone())
})
.collect();
lf = lf.with_columns(texts);
}
if first_row.is_some() {
let mut ordered: Vec<Expr> = schema.iter_names().map(|name| col(name.clone())).collect();
ordered.push(col(DRIFT_COLUMN));
lf = lf.select(ordered);
}
Ok(lf)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_sample_reads_a_file_named_like_a_glob() {
let dir = tempfile::tempdir().unwrap();
let names = |name: &str, format| {
column_schema_of(&dir.path().join(name), format, &ReadAs::default())
.unwrap()
.into_iter()
.map(|(n, _)| n)
.collect::<Vec<_>>()
};
std::fs::write(dir.path().join("d[1].jsonl"), "{\"own\": 1}\n").unwrap();
std::fs::write(dir.path().join("d1.jsonl"), "{\"other\": 1}\n").unwrap();
std::fs::write(dir.path().join("d[1].csv"), "own\n1\n").unwrap();
std::fs::write(dir.path().join("d1.csv"), "other\n1\n").unwrap();
assert_eq!(names("d[1].jsonl", crate::FileFormat::Jsonl), ["own"]);
assert_eq!(names("d[1].csv", crate::FileFormat::Csv), ["own"]);
}
#[test]
fn the_sample_splits_on_the_separator_the_open_uses() {
let dir = tempfile::tempdir().unwrap();
let names = |name: &str, body: &str, format, delimiter| {
let path = dir.path().join(name);
std::fs::write(&path, body).unwrap();
let as_read = ReadAs {
delimiter,
..ReadAs::default()
};
column_schema_of(&path, format, &as_read)
.unwrap()
.into_iter()
.map(|(n, _)| n)
.collect::<Vec<_>>()
};
use crate::FileFormat::{Csv, Psv, Tsv};
assert_eq!(names("a.tsv", "a\tb\n1\t2\n", Tsv, None), ["a", "b"]);
assert_eq!(names("a.psv", "a|b\n1|2\n", Psv, None), ["a", "b"]);
assert_eq!(names("a.csv", "a|b\n1|2\n", Csv, None), ["a|b"]);
assert_eq!(names("b.csv", "a|b\n1|2\n", Csv, Some(b'|')), ["a", "b"]);
}
fn file(columns: &[(&str, DataType)], rows: usize) -> Option<FileFooter> {
let mut schema = Schema::with_capacity(columns.len());
for (name, dtype) in columns {
schema.with_column((*name).into(), dtype.clone());
}
Some(FileFooter {
schema: Arc::new(schema),
row_group_rows: vec![rows],
file_bytes: 0,
row_group_bytes: Vec::new(),
column_bytes: Vec::new(),
})
}
#[test]
fn column_widths_average_over_the_rows_read() {
let with = |rows, bytes: &[(&str, usize)]| {
let mut footer = file(&[], rows)?;
footer.column_bytes = bytes.iter().map(|(n, b)| (n.to_string(), *b)).collect();
Some(footer)
};
let footers = [
with(3, &[("blob", 3_000), ("id", 24)]),
None,
with(1, &[("id", 8)]),
];
assert_eq!(
column_bytes_per_row(&footers),
[("blob".to_string(), 750), ("id".to_string(), 8)]
);
assert!(column_bytes_per_row(&[with(0, &[("id", 0)])]).is_empty());
}
fn cols(files: &[&[&str]]) -> Vec<Vec<String>> {
files
.iter()
.map(|f| f.iter().map(|c| (*c).to_string()).collect())
.collect()
}
fn grew(files: usize, from: usize, to: usize) -> Vec<Vec<String>> {
let names = |n: usize| (0..n).map(|i| format!("c{i}")).collect::<Vec<_>>();
let mut out: Vec<Vec<String>> = (0..files - 1).map(|_| names(from)).collect();
out.push(names(to));
out
}
#[test]
fn one_table_whatever_its_files_did_over_time() {
for (what, files) in [
(
"identical part files",
cols(&[&["a", "b", "c"], &["a", "b", "c"], &["a", "b", "c"]]),
),
(
"a column only one file has",
cols(&[&["id"], &["id", "oops"], &["id"]]),
),
(
"a column that starts",
cols(&[&["id", "ts"], &["id", "ts"], &["id", "ts", "fee"]]),
),
(
"a column that stops",
cols(&[&["id", "ts", "fee"], &["id", "ts"], &["id", "ts"]]),
),
(
"a file truncated to one column",
cols(&[&["a", "b", "c", "d"], &["a"], &["a", "b", "c", "d"]]),
),
("five columns grown to fifty", grew(10, 5, 50)),
("one file", cols(&[&["a", "b"]])),
("no files", Vec::new()),
] {
assert!(is_nested(&files), "{what} should read as one table");
}
}
#[test]
fn a_column_each_way_is_not_nesting_and_costs_a_keystroke() {
for (what, files) in [
(
"one column each way",
cols(&[&["a", "b", "c", "d", "e"], &["a", "b", "c", "d", "f"]]),
),
(
"a column renamed",
cols(&[&["id", "ts", "amount"], &["id", "ts", "amt"]]),
),
] {
assert!(
!is_nested(&files),
"{what} brings a column the widest file cannot account for"
);
}
}
#[test]
fn separate_tables_are_not_one_table() {
for (what, files) in [
("two tables", cols(&[&["a", "b", "c"], &["x", "y", "z"]])),
(
"tables sharing a key",
cols(&[&["id", "a", "b"], &["id", "x", "y"], &["id", "p", "q"]]),
),
(
"a season of six tables",
cols(&[
&[
"season",
"circuit_id",
"url",
"circuit_name",
"lat",
"lng",
"locality",
"country",
],
&["season", "constructor_id", "url", "name", "nationality"],
&[
"season",
"round",
"driver_id",
"position",
"points",
"wins",
"constructor_id",
],
&[
"season",
"driver_id",
"permanent_number",
"code",
"url",
"given_name",
"family_name",
"date_of_birth",
"nationality",
],
&[
"season",
"round",
"race_name",
"circuit_id",
"race_date",
"driver_id",
"constructor_id",
"number",
"grid",
"position",
"position_text",
"points",
"laps",
"status",
"time_millis",
"time_text",
"fastest_lap_rank",
"fastest_lap_number",
"fastest_lap_time",
"fastest_lap_avg_speed",
],
&[
"season",
"round",
"race_name",
"circuit_id",
"circuit_name",
"locality",
"country",
"lat",
"lng",
"date",
"time",
"qualifying_date",
"qualifying_time",
"sprint_date",
"sprint_time",
"sprint_shootout_date",
"sprint_shootout_time",
"url",
],
]),
),
] {
assert!(!is_nested(&files), "{what} should be separate tables");
}
}
#[test]
fn growth_nests_where_unrelated_tables_do_not() {
assert!(is_nested(&grew(10, 5, 50)));
let unrelated = cols(&[&["id", "a", "b"], &["id", "x", "y"], &["id", "p", "q"]]);
assert!(!is_nested(&unrelated));
}
#[test]
fn two_tables_sharing_a_key_do_not_nest() {
let files = cols(&[&["id", "name"], &["id", "customer_id", "amount"]]);
assert!(!is_nested(&files));
}
#[test]
fn a_nested_column_is_one_column_however_it_was_written() {
let old_writer = vec![
"id".to_string(),
"tags.array".to_string(),
"refs.array".to_string(),
];
let new_writer = vec![
"id".to_string(),
"tags.list.element".to_string(),
"refs.list.element".to_string(),
];
let files = vec![
top_level_columns(&old_writer),
top_level_columns(&new_writer),
];
assert_eq!(files[0], vec!["id", "tags", "refs"]);
assert!(is_nested(&files), "the same three columns, written twice");
assert!(!is_nested(&[old_writer, new_writer]));
}
#[test]
fn an_empty_file_does_not_decide_the_directory() {
assert!(is_nested(&cols(&[&["a", "b"], &[], &["a", "b"]])));
assert!(is_nested(&cols(&[&[], &[]])), "nothing to disagree about");
}
fn union(files: &[Option<FileFooter>]) -> DatasetSchema {
union_file_schemas(files, SchemaOrigin::AllFooters(files.len()))
}
fn omitted_names(union: &DatasetSchema, index: usize) -> Vec<String> {
union.omitted[index]
.iter()
.map(|(name, _)| name.to_string())
.collect()
}
fn names(schema: &Schema) -> Vec<String> {
schema.iter_names().map(|n| n.to_string()).collect()
}
#[test]
fn a_conflicting_column_read_as_text_shows_every_file_s_values() {
use polars::prelude::{ParquetWriter, df};
let dir = tempfile::tempdir().unwrap();
let mut paths = Vec::new();
let mut write = |name: &str, mut frame: polars::prelude::DataFrame| {
let path = dir.path().join(name);
let file = std::fs::File::create(&path).unwrap();
ParquetWriter::new(file).finish(&mut frame).unwrap();
paths.push(path.to_string_lossy().to_string());
};
write(
"a.parquet",
df!("id" => &[0i64, 1, 2], "n" => &[10i64, 20, 30]).unwrap(),
);
write("b.parquet", df!("id" => &[3i64], "n" => &["x"]).unwrap());
write("c.parquet", df!("id" => &[4i64], "n" => &[true]).unwrap());
let footers: Vec<Option<FileFooter>> = vec![
file(&[("id", DataType::Int64), ("n", DataType::Int64)], 3),
file(&[("id", DataType::Int64), ("n", DataType::String)], 1),
file(&[("id", DataType::Int64), ("n", DataType::Boolean)], 1),
];
let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(3));
assert_eq!(
dataset.schema.get("n"),
Some(&DataType::Int64),
"the integer file has the most rows"
);
let drift = ScanDrift::new(&paths, &dataset, &[3, 1, 1]).expect("the files disagree");
let plain = lenient_scan(&paths, dataset.schema.clone(), None, Some(&drift), &[])
.unwrap()
.collect()
.unwrap();
let n = plain.column("n").unwrap();
assert_eq!(
(0..n.len())
.map(|i| n.get(i).unwrap().to_string())
.collect::<Vec<_>>(),
["10", "20", "30", "null", "null"],
"the text and boolean files hold a value, and it is not one this column \
can carry"
);
let as_text = [PlSmallStr::from("n")];
let text = lenient_scan(&paths, dataset.schema.clone(), None, Some(&drift), &as_text)
.unwrap()
.collect()
.unwrap();
assert_eq!(
text.column("n").unwrap().dtype(),
&DataType::String,
"the column is text now"
);
let n = text.column("n").unwrap().str().unwrap();
assert_eq!(
n.iter().collect::<Vec<_>>(),
[Some("10"), Some("20"), Some("30"), Some("x"), Some("true")],
"and holds what each file wrote, spelled as that file's own type prints"
);
let ids = text.column("id").unwrap().i64().unwrap();
assert_eq!(
ids.into_no_null_iter().collect::<Vec<_>>(),
[0, 1, 2, 3, 4],
"in dataset order, with the rows still lined up against their ids"
);
assert_eq!(
text.column(DRIFT_COLUMN)
.unwrap()
.u32()
.unwrap()
.into_no_null_iter()
.collect::<Vec<_>>(),
[0, 1, 2, 3, 4],
"and each row still knows its place in the dataset"
);
}
#[test]
fn a_date_past_the_calendar_read_as_text_is_its_stored_number() {
use polars::prelude::{NamedFrom, ParquetWriter, Series, TimeZone};
let dir = tempfile::tempdir().unwrap();
let mut paths = Vec::new();
let mut write = |name: &str, n: Series| {
let path = dir.path().join(name);
let file = std::fs::File::create(&path).unwrap();
let mut frame = polars::prelude::DataFrame::new_infer_height(vec![n.into()]).unwrap();
ParquetWriter::new(file).finish(&mut frame).unwrap();
paths.push(path.to_string_lossy().to_string());
};
let paris = TimeZone::opt_try_new(Some("Europe/Paris")).unwrap();
let stamps = |dtype: DataType| {
Series::new("n".into(), [0, i64::MIN + 1])
.cast(&dtype)
.unwrap()
};
let types = [
DataType::Date,
DataType::Datetime(TimeUnit::Milliseconds, None),
DataType::Datetime(TimeUnit::Microseconds, paris),
];
write("a.parquet", Series::new("n".into(), ["x", "y", "z"]));
write(
"b.parquet",
Series::new("n".into(), [0, i32::MAX])
.cast(&types[0])
.unwrap(),
);
write("c.parquet", stamps(types[1].clone()));
write("d.parquet", stamps(types[2].clone()));
let footers: Vec<Option<FileFooter>> = [DataType::String]
.into_iter()
.chain(types)
.enumerate()
.map(|(i, dtype)| file(&[("n", dtype)], if i == 0 { 3 } else { 2 }))
.collect();
let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(4));
assert_eq!(dataset.schema.get("n"), Some(&DataType::String));
let drift = ScanDrift::new(&paths, &dataset, &[3, 2, 2, 2]).expect("the files disagree");
let as_text = [PlSmallStr::from("n")];
let text = lenient_scan(&paths, dataset.schema.clone(), None, Some(&drift), &as_text)
.unwrap()
.collect()
.unwrap();
assert_eq!(
text.column("n")
.unwrap()
.str()
.unwrap()
.iter()
.collect::<Vec<_>>(),
[
Some("x"),
Some("y"),
Some("z"),
Some("1970-01-01"),
Some("2147483647 days since 1970-01-01"),
Some("1970-01-01 00:00:00.000"),
Some("-9223372036854775807 ms since 1970-01-01 UTC"),
Some("1970-01-01 01:00:00.000000+01:00"),
Some("-9223372036854775807 us since 1970-01-01 UTC"),
]
);
}
#[test]
fn reading_as_text_leaves_a_file_without_the_column_alone() {
use polars::prelude::{ParquetWriter, df};
let dir = tempfile::tempdir().unwrap();
let mut paths = Vec::new();
let mut write = |name: &str, mut frame: polars::prelude::DataFrame| {
let path = dir.path().join(name);
let f = std::fs::File::create(&path).unwrap();
ParquetWriter::new(f).finish(&mut frame).unwrap();
paths.push(path.to_string_lossy().to_string());
};
write(
"a.parquet",
df!("id" => &[0i64, 1], "n" => &[10i64, 20]).unwrap(),
);
write("b.parquet", df!("id" => &[2i64]).unwrap());
write("c.parquet", df!("id" => &[3i64], "n" => &["x"]).unwrap());
let footers: Vec<Option<FileFooter>> = vec![
file(&[("id", DataType::Int64), ("n", DataType::Int64)], 2),
file(&[("id", DataType::Int64)], 1),
file(&[("id", DataType::Int64), ("n", DataType::String)], 1),
];
let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(3));
let drift = ScanDrift::new(&paths, &dataset, &[2, 1, 1]).expect("the files disagree");
let as_text = [PlSmallStr::from("n")];
let text = lenient_scan(&paths, dataset.schema.clone(), None, Some(&drift), &as_text)
.unwrap()
.collect()
.unwrap();
assert_eq!(
text.column("n")
.unwrap()
.str()
.unwrap()
.iter()
.collect::<Vec<_>>(),
[Some("10"), Some("20"), None, Some("x")],
"the file with no `n` has none to show"
);
}
#[test]
fn types_the_cast_agrees_with_are_exactly_the_ones_offered() {
use polars::prelude::*;
let mk = |dtype: DataType| -> Column {
Series::new("x".into(), [1i64, 2])
.cast(&dtype)
.unwrap_or_else(|e| panic!("cannot build a {dtype:?} column: {e}"))
.into()
};
let mut cases: Vec<(DataType, Column)> = vec![
DataType::Int64,
DataType::Float64,
DataType::Boolean,
DataType::Date,
DataType::Time,
DataType::Datetime(TimeUnit::Microseconds, None),
DataType::Duration(TimeUnit::Milliseconds),
DataType::Decimal(10, 2),
DataType::List(Box::new(DataType::Int64)),
]
.into_iter()
.map(|dtype| (dtype.clone(), mk(dtype)))
.collect();
cases.push((DataType::String, Series::new("x".into(), ["a", "b"]).into()));
cases.push((
DataType::Binary,
Series::new("x".into(), [&[0xffu8, 0xfe][..], &[0x41][..]]).into(),
));
let plain =
StructChunked::from_series("x".into(), 2, [Series::new("a".into(), [1i64, 2])].iter())
.unwrap()
.into_series();
cases.push((plain.dtype().clone(), plain.into()));
let inners: [Series; 3] = [
Series::new("a".into(), [1i64, 2])
.cast(&DataType::Duration(TimeUnit::Milliseconds))
.unwrap(),
Series::new("a".into(), [1i64, 2])
.cast(&DataType::List(Box::new(DataType::Int64)))
.unwrap(),
Series::new("a".into(), [&[0xffu8, 0xfe][..], &[0x41][..]]),
];
for inner in inners {
let nested = StructChunked::from_series("x".into(), 2, [inner].iter())
.unwrap()
.into_series();
cases.push((nested.dtype().clone(), nested.into()));
}
for (dtype, column) in cases {
let cast_works = DataFrame::new(2, vec![column])
.unwrap()
.lazy()
.select([col("x").cast(DataType::String)])
.collect()
.is_ok();
assert_eq!(
can_read_as_text(&dtype),
cast_works,
"{dtype:?}: the predicate and the cast must agree"
);
}
}
#[test]
fn a_column_the_cast_refuses_is_read_as_it_was() {
use polars::prelude::{ParquetWriter, df};
let dir = tempfile::tempdir().unwrap();
let mut paths = Vec::new();
let mut write = |name: &str, mut frame: polars::prelude::DataFrame| {
let path = dir.path().join(name);
let f = std::fs::File::create(&path).unwrap();
ParquetWriter::new(f).finish(&mut frame).unwrap();
paths.push(path.to_string_lossy().to_string());
};
write(
"a.parquet",
df!("id" => &[0i64, 1], "n" => &[&[0xffu8, 0xfe][..], &[0x41][..]]).unwrap(),
);
write("b.parquet", df!("id" => &[2i64], "n" => &["x"]).unwrap());
let footers: Vec<Option<FileFooter>> = vec![
file(&[("id", DataType::Int64), ("n", DataType::Binary)], 2),
file(&[("id", DataType::Int64), ("n", DataType::String)], 1),
];
let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(2));
let drifting = dataset
.columns
.iter()
.find(|column| column.name == "n")
.unwrap();
assert!(
!drifting.can_read_as_text(),
"so the Notes tab never offers it"
);
let drift = ScanDrift::new(&paths, &dataset, &[2, 1]).expect("the files disagree");
let as_text = [PlSmallStr::from("n")];
let frame = lenient_scan(&paths, dataset.schema.clone(), None, Some(&drift), &as_text)
.unwrap()
.collect()
.expect("the read still succeeds, which is the point");
assert_eq!(
frame
.column("id")
.unwrap()
.i64()
.unwrap()
.into_no_null_iter()
.collect::<Vec<_>>(),
[0, 1, 2],
"every file is still read, the agreeing one included"
);
assert_ne!(
frame.column("n").unwrap().dtype(),
&DataType::String,
"and the column is as it was, not half-cast"
);
}
#[test]
fn a_type_only_one_file_holds_can_rule_the_column_out() {
use polars::prelude::{IntoLazy, ParquetWriter, col, df};
let dir = tempfile::tempdir().unwrap();
let mut paths = Vec::new();
let mut write = |name: &str, mut frame: polars::prelude::DataFrame| {
let path = dir.path().join(name);
let f = std::fs::File::create(&path).unwrap();
ParquetWriter::new(f).finish(&mut frame).unwrap();
paths.push(path.to_string_lossy().to_string());
};
write(
"a.parquet",
df!("id" => &[0i64, 1, 2], "n" => &[10i64, 20, 30]).unwrap(),
);
write(
"b.parquet",
df!("id" => &[3i64], "n" => &[9i64])
.unwrap()
.lazy()
.group_by([col("id")])
.agg([col("n")])
.collect()
.unwrap(),
);
let footers: Vec<Option<FileFooter>> = vec![
file(&[("id", DataType::Int64), ("n", DataType::Int64)], 3),
file(
&[
("id", DataType::Int64),
("n", DataType::List(Box::new(DataType::Int64))),
],
1,
),
];
let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(2));
assert_eq!(
dataset.schema.get("n"),
Some(&DataType::Int64),
"read as the integer the three rows have"
);
let drifting = dataset
.columns
.iter()
.find(|column| column.name == "n")
.unwrap();
assert!(
can_read_as_text(&drifting.dtype),
"an integer column casts to text on its own account"
);
assert!(
!drifting.can_read_as_text(),
"but one file holds a list, and that file's cast is the one that fails"
);
let drift = ScanDrift::new(&paths, &dataset, &[3, 1]).expect("the files disagree");
let as_text = [PlSmallStr::from("n")];
let frame = lenient_scan(&paths, dataset.schema.clone(), None, Some(&drift), &as_text)
.unwrap()
.collect()
.expect("asking anyway must not cost the read");
assert_eq!(
frame
.column("id")
.unwrap()
.i64()
.unwrap()
.into_no_null_iter()
.collect::<Vec<_>>(),
[0, 1, 2, 3],
"every file is read, the three that agreed included"
);
assert_eq!(
frame.column("n").unwrap().dtype(),
&DataType::Int64,
"and the column is as it was"
);
}
#[test]
fn text_schema_respells_without_reordering() {
let mut schema = Schema::with_capacity(3);
schema.with_column("a".into(), DataType::Int64);
schema.with_column("n".into(), DataType::Int64);
schema.with_column("z".into(), DataType::Float64);
let schema = Arc::new(schema);
let text = text_schema(&schema, &[PlSmallStr::from("n")]);
assert_eq!(
names(&text),
["a", "n", "z"],
"a column read differently does not move"
);
assert_eq!(text.get("n"), Some(&DataType::String));
assert_eq!(
text.get("a"),
Some(&DataType::Int64),
"nor do its neighbours change"
);
assert_eq!(text.get("z"), Some(&DataType::Float64));
assert!(
Arc::ptr_eq(&schema, &text_schema(&schema, &[])),
"asking for nothing is the schema itself"
);
assert_eq!(
names(&text_schema(&schema, &[PlSmallStr::from("ghost")])),
["a", "n", "z"],
"a name the schema does not have adds nothing"
);
}
#[test]
fn a_sampled_dataset_counts_what_it_read_and_not_what_it_did_not() {
let footers = vec![
file(&[("id", DataType::Int64)], 0),
file(&[("id", DataType::Int64), ("x", DataType::String)], 5),
None,
];
let read = [0usize, 250, 499];
let union = union_sampled(500, &read, &footers);
assert_eq!(
union.files, 3,
"the population is the footers read, not the files there are"
);
assert_eq!(union.empty_files, 1, "one of the three held nothing");
assert_eq!(
union.origin,
SchemaOrigin::FooterSample {
read: 3,
total: 500
}
);
assert_eq!(
union.unreadable,
[499],
"and the footer that would not parse is named by its place among the files"
);
assert_eq!(union.file_group.len(), 500);
assert_eq!(union.omitted.len(), 500);
}
#[test]
fn row_groups_are_noted_by_their_middle_size_and_only_when_it_is_large() {
const MIB: usize = 1024 * 1024;
let note = |groups: &[&[usize]]| -> Option<String> {
let files: Vec<Option<FileFooter>> = groups
.iter()
.map(|sizes| {
Some(FileFooter {
schema: Arc::new(Schema::with_capacity(0)),
row_group_rows: vec![1],
file_bytes: 0,
row_group_bytes: sizes.to_vec(),
column_bytes: Vec::new(),
})
})
.collect();
let dataset = union_file_schemas(&files, SchemaOrigin::AllFooters(files.len()));
crate::notes::from_dataset(&dataset)
.into_iter()
.find(|note| note.summary.starts_with("median row group"))
.map(|note| note.summary)
};
assert_eq!(note(&[&[MIB], &[2 * MIB]]), None, "ordinary row groups");
assert_eq!(
note(&[&[64 * MIB]]),
None,
"the threshold itself is not past it"
);
assert_eq!(
note(&[&[65 * MIB]]).as_deref(),
Some("median row group 65.0 MiB, each read whole"),
);
assert_eq!(
note(&[&[MIB, MIB, 4096 * MIB]]),
None,
"one huge row group among small ones does not describe the dataset"
);
assert_eq!(
note(&[&[100 * MIB, 100 * MIB], &[MIB]]).as_deref(),
Some("median row group 100.0 MiB, each read whole"),
"the middle of every row group of every file, not the middle of the files"
);
assert_eq!(
note(&[&[MIB], &[100 * MIB, 100 * MIB]]).as_deref(),
Some("median row group 100.0 MiB, each read whole"),
"including when the large ones are not in the first file"
);
assert_eq!(
note(&[&[100 * MIB], &[MIB], &[100 * MIB]]).as_deref(),
Some("median row group 100.0 MiB, each read whole"),
"and when they arrive out of order"
);
assert_eq!(
note(&[&[MIB], &[100 * MIB], &[MIB]]),
None,
"which cuts both ways: one big group between two small ones is not the middle"
);
assert_eq!(note(&[&[]]), None, "a file with no row groups says nothing");
assert_eq!(
note(&[&[64 * MIB, 65 * MIB]]),
None,
"two row groups either side of the line: the lower one decides"
);
assert_eq!(
note(&[&[65 * MIB, 66 * MIB]]).as_deref(),
Some("median row group 65.0 MiB, each read whole"),
"and when it decides the other way it is still the lower one"
);
}
#[test]
fn a_real_footer_reports_the_compressed_size_of_each_row_group() {
use polars::prelude::{ParquetWriter, df};
let dir = tempfile::tempdir().unwrap();
let rows: Vec<String> = (0..20_000)
.map(|i| format!("{i:0>6}{}", "abcdefghij".repeat(19)))
.collect();
let mut frame = df!("s" => rows).unwrap();
let file = std::fs::File::create(dir.path().join("wide.parquet")).unwrap();
ParquetWriter::new(file)
.with_row_group_size(Some(20_000))
.finish(&mut frame)
.unwrap();
let footer = crate::dataset_files::local_footer(&dir.path().join("wide.parquet"))
.expect("the footer reads");
assert_eq!(footer.rows(), 20_000);
assert_eq!(footer.row_group_bytes.len(), 1, "one row group");
let on_disk = std::fs::metadata(dir.path().join("wide.parquet"))
.unwrap()
.len();
assert_eq!(
footer.file_bytes as u64, on_disk,
"the file's size, as the filesystem reports it"
);
let size = footer.row_group_bytes[0];
assert!(size > 0, "a size is reported");
assert!(
size < 1_000_000,
"and it is the compressed size: 20,000 distinct strings of 200 characters \
are about 4 MiB decoded and a small fraction of that on disk, so {size} \
bytes is the decoded figure"
);
}
#[test]
fn many_files_are_noted_only_when_they_are_also_small() {
const KIB: usize = 1024;
const MIB: usize = 1024 * KIB;
let note = |files: usize, read: usize, sizes: &[usize]| -> Option<String> {
let footers: Vec<Option<FileFooter>> = sizes
.iter()
.cycle()
.take(if sizes.is_empty() { 0 } else { read })
.map(|bytes| {
Some(FileFooter {
schema: Arc::new(Schema::with_capacity(0)),
row_group_rows: vec![1],
file_bytes: *bytes,
row_group_bytes: Vec::new(),
column_bytes: Vec::new(),
})
})
.collect();
let origin = if read == files {
SchemaOrigin::AllFooters(files)
} else {
SchemaOrigin::FooterSample { read, total: files }
};
crate::notes::from_dataset(&union_file_schemas(&footers, origin))
.into_iter()
.find(|note| note.summary.contains("files, median"))
.map(|note| note.summary)
};
assert_eq!(
note(10_000, 10_000, &[40 * KIB]),
None,
"a year of hourly partitions, and more, is an ordinary shape"
);
assert_eq!(
note(10_001, 10_001, &[40 * KIB]).as_deref(),
Some("10,001 files, median 40.0 KiB; every footer read before any row"),
"one more is not"
);
assert_eq!(
note(50_000, 50_000, &[MIB]),
None,
"a megabyte is not small by this measure"
);
assert!(
note(50_000, 50_000, &[MIB - 1]).is_some(),
"a byte under it is"
);
assert_eq!(
note(50_000, 50_000, &[40 * KIB, 40 * KIB, 900 * MIB]).as_deref(),
Some("50,000 files, median 40.0 KiB; every footer read before any row"),
"a large minority does not move the middle"
);
assert_eq!(
note(500_000, 2, &[40 * KIB, 40 * KIB]).as_deref(),
Some("500,000 files, median 40.0 KiB; 2 footers read before any row"),
"the count is the listing's; the footers read are their own number"
);
assert_eq!(note(50_000, 0, &[]), None, "no footer read, nothing to say");
let sampled = union_file_schemas(
&[
Some(FileFooter {
schema: Arc::new(Schema::with_capacity(0)),
row_group_rows: vec![1],
file_bytes: 40 * KIB,
row_group_bytes: Vec::new(),
column_bytes: Vec::new(),
}),
Some(FileFooter {
schema: Arc::new(Schema::with_capacity(0)),
row_group_rows: vec![1],
file_bytes: 40 * KIB,
row_group_bytes: Vec::new(),
column_bytes: Vec::new(),
}),
],
SchemaOrigin::FooterSample {
read: 2,
total: 500_000,
},
);
let sampled_note = crate::notes::from_dataset(&sampled)
.into_iter()
.find(|note| note.summary.contains("files, median"))
.expect("the note is made");
assert_eq!(sampled_note.scope, "in 2 of 500,000 footers (sample)");
assert_eq!(
note(50_000, 2, &[0, 0]),
None,
"and a size of nothing means the size is not known, not that it is small"
);
}
#[test]
fn partition_keys_are_the_key_equals_segments_above_the_file() {
let keys = |path: &str| partition_keys_of(path);
assert_eq!(keys("data/date=2024-01-01/a.parquet"), ["date"]);
assert_eq!(keys("data/y=2024/m=05/a.parquet"), ["m", "y"]);
assert_eq!(
keys("data/m=05/y=2024/a.parquet"),
keys("data/y=2024/m=05/a.parquet"),
"the same two partitions, written in two orders"
);
assert_eq!(
keys("data/x=1/x=2/a.parquet"),
["x"],
"and a key repeated down the tree is one key"
);
assert_eq!(keys("data/a.parquet"), Vec::<String>::new());
assert_eq!(
keys("data/x=1/2024=05.parquet"),
["x"],
"the file's own name is not a partition, whatever it looks like"
);
assert_eq!(
keys("data/=2024/a.parquet"),
Vec::<String>::new(),
"nor is a segment with nothing before the equals"
);
#[cfg(windows)]
assert_eq!(
keys(r"data\date=2024-01-01\a.parquet"),
["date"],
"a path written the other way round is the same path"
);
#[cfg(not(windows))]
assert_eq!(
keys(r"data/we\ird=1/f.parquet"),
["we\\ird"],
"a backslash here is part of the name, not a separator"
);
#[cfg(not(windows))]
assert_eq!(
keys(r"data/x=1\y=2/f.parquet"),
["x"],
"so one directory is one partition, however it is spelled"
);
}
#[test]
fn directories_that_partition_differently_are_counted_each_way() {
let note = |root: &str, paths: &[&str]| -> Option<crate::notes::Note> {
let footers = vec![
Some(FileFooter {
schema: Arc::new(Schema::with_capacity(0)),
row_group_rows: vec![1],
file_bytes: 1,
row_group_bytes: Vec::new(),
column_bytes: Vec::new(),
});
paths.len()
];
let owned: Vec<String> = paths.iter().map(|p| p.to_string()).collect();
let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(paths.len()))
.with_partition_layouts(root, &owned);
crate::notes::from_dataset(&dataset)
.into_iter()
.find(|note| note.summary.contains("mixed partition keys"))
};
assert_eq!(
note("d", &["d/date=1/a.parquet", "d/date=2/b.parquet"]),
None,
"directories that agree have nothing to say"
);
assert_eq!(
note("d", &["d/a.parquet", "d/b.parquet"]),
None,
"nor has a dataset with no partitions at all"
);
assert_eq!(
note("d", &["d/y=1/m=1/a.parquet", "d/m=2/y=2/b.parquet"]),
None,
"nor two orders of the same two keys: hive matches columns by name, so \
that dataset reads perfectly well and has nothing in dispute"
);
assert_eq!(
note(
"d/run=7",
&["d/run=7/loose.parquet", "d/run=7/date=1/a.parquet"]
),
None,
"a key=value directory above the dataset as it was opened is not one of the \
things its directories disagree about — and these two files are where that \
matters, since counting `run` would make the one without a key of its \
own a second layout"
);
assert_eq!(
note(
"s3://b//data/",
&["s3://b/data/date=1/a.parquet", "s3://b/data/dt=2/b.parquet"]
),
None,
"and a path the root is not a prefix of — a typed URL with a doubled \
slash rebuilds without it — is one this cannot place, so it is left out \
rather than read from the top"
);
let renamed = note(
"d",
&[
"d/date=1/a.parquet",
"d/date=2/b.parquet",
"d/date=3/c.parquet",
"d/dt=4/e.parquet",
],
)
.expect("the directories disagree");
assert_eq!(
renamed.summary,
"mixed partition keys: 3 files by date, \
1 file by dt"
);
assert_eq!(
renamed.scope, "in the names of 4 files",
"read off every name, not off the footers datui opened"
);
let loose = note(
"d",
&[
"d/y=1/m=1/a.parquet",
"d/y=1/m=2/b.parquet",
"d/date=3/c.parquet",
"d/loose.parquet",
],
)
.expect("the directories disagree");
assert_eq!(
loose.summary,
"mixed partition keys: 2 files by m/y, \
1 file by date"
);
assert_eq!(
loose.scope, "in the names of 4 files",
"the unpartitioned file is one of the names read"
);
}
#[test]
fn the_layouts_a_note_names_are_the_commonest_of_them() {
let layouts = |paths: &[&str]| -> Vec<(Vec<String>, usize)> {
let owned: Vec<String> = paths.iter().map(|p| p.to_string()).collect();
union_file_schemas(&[], SchemaOrigin::AllFooters(0))
.with_partition_layouts("d", &owned)
.partition_layouts
};
let note = |paths: &[&str]| -> String {
let footers = vec![
Some(FileFooter {
schema: Arc::new(Schema::with_capacity(0)),
row_group_rows: vec![1],
file_bytes: 1,
row_group_bytes: Vec::new(),
column_bytes: Vec::new(),
});
paths.len()
];
let owned: Vec<String> = paths.iter().map(|p| p.to_string()).collect();
let dataset = union_file_schemas(&footers, SchemaOrigin::AllFooters(paths.len()))
.with_partition_layouts("d", &owned);
crate::notes::from_dataset(&dataset)
.into_iter()
.find(|note| note.summary.contains("mixed partition keys"))
.expect("the directories disagree")
.summary
};
assert_eq!(
layouts(&[
"d/zzz=1/b.parquet",
"d/aaa=1/a.parquet",
"d/zzz=2/c.parquet",
"d/zzz=3/e.parquet",
]),
vec![(vec!["zzz".to_string()], 3), (vec!["aaa".to_string()], 1)],
"commonest first, though the rare one sorts first and arrived first"
);
assert_eq!(
layouts(&[
"d/zz=1/a.parquet",
"d/aa=1/b.parquet",
"d/mm=1/c.parquet",
"d/qq=1/e.parquet"
]),
vec![
(vec!["aa".to_string()], 1),
(vec!["mm".to_string()], 1),
(vec!["qq".to_string()], 1),
(vec!["zz".to_string()], 1)
],
"and equally common ones by their keys, so the same dataset reads the \
same way every time it is opened"
);
assert_eq!(
note(&[
"d/aa=1/a.parquet",
"d/bb=1/b.parquet",
"d/cc=1/c.parquet",
"d/dd=1/e.parquet",
]),
"mixed partition keys: 1 file by aa, \
1 file by bb, 2 files by 2 other ways"
);
assert_eq!(
note(&["d/aa=1/a.parquet", "d/bb=1/b.parquet", "d/cc=1/c.parquet"]),
"mixed partition keys: 1 file by aa, \
1 file by bb, 1 file by 1 other way",
"and one of them is one way, not one ways"
);
let many: Vec<String> = (0..100)
.map(|i| format!("d/k{i:0>3}=1/f.parquet"))
.collect();
let many: Vec<&str> = many.iter().map(String::as_str).collect();
assert_eq!(
note(&many),
"mixed partition keys: 1 file by k000, \
1 file by k001, 98 files by 98 other ways"
);
let owned: Vec<String> = many.iter().map(|p| p.to_string()).collect();
let dataset = union_file_schemas(&[], SchemaOrigin::AllFooters(0))
.with_partition_layouts("d", &owned);
assert!(
dataset.partition_layouts.len() <= 64,
"and it is not holding a hundred of them to say so: {}",
dataset.partition_layouts.len()
);
}
#[test]
fn the_footer_count_speaks_only_while_a_pass_is_running() {
let progress = FooterProgress::default();
assert_eq!(progress.reading(), None, "nothing has begun");
progress.begin(3);
assert_eq!(progress.reading(), Some((0, 3)), "none read yet");
progress.advance();
progress.advance();
assert_eq!(progress.reading(), Some((2, 3)));
progress.done();
assert_eq!(progress.reading(), None, "and nothing once it has landed");
progress.begin(2);
assert_eq!(progress.reading(), Some((0, 2)));
}
#[test]
fn the_footer_count_never_passes_its_total() {
let progress = FooterProgress::default();
progress.begin(2);
for _ in 0..5 {
progress.advance();
}
assert_eq!(progress.reading(), Some((2, 2)));
}
#[test]
fn a_pass_that_panics_still_says_it_has_finished() {
let progress = FooterProgress::default();
let caught = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let pass = progress.pass(3);
pass.advance();
panic!("a footer reader gave up");
}));
assert!(caught.is_err(), "the panic happened");
assert_eq!(
progress.reading(),
None,
"and the count went with it rather than sitting there"
);
assert_eq!(
progress.last_pass().read,
1,
"with what it managed still readable"
);
}
#[test]
fn absent_columns_alone_never_split_the_scan() {
let files = 64;
let per_file: Vec<Option<FileFooter>> = (0..files)
.map(|i| {
let mut s = Schema::with_capacity(2);
s.with_column("id".into(), DataType::Int64);
if i % 2 == 1 {
s.with_column("extra".into(), DataType::String);
}
Some(FileFooter {
schema: Arc::new(s),
row_group_rows: vec![1],
file_bytes: 0,
row_group_bytes: Vec::new(),
column_bytes: Vec::new(),
})
})
.collect();
let paths: Vec<String> = (0..files).map(|i| format!("part-{i:05}.parquet")).collect();
let read: Vec<usize> = (0..files).collect();
let dataset = union_sampled(files, &read, &per_file);
let rows = vec![1usize; files];
let drift = ScanDrift::new(&paths, &dataset, &rows).expect("this dataset drifts");
assert!(dataset.drifts());
assert_eq!(
runs_of(&paths, &drift),
1,
"absent columns need no split, however they alternate"
);
let mut with_conflict = per_file.clone();
let mut odd = Schema::with_capacity(2);
odd.with_column("id".into(), DataType::String);
with_conflict[7] = Some(FileFooter {
schema: Arc::new(odd),
row_group_rows: vec![1],
file_bytes: 0,
row_group_bytes: Vec::new(),
column_bytes: Vec::new(),
});
let dataset = union_sampled(files, &read, &with_conflict);
let drift = ScanDrift::new(&paths, &dataset, &rows).unwrap();
assert_eq!(runs_of(&paths, &drift), 3, "before it, it, and after it");
}
fn runs_of(paths: &[String], drift: &ScanDrift) -> usize {
let mut runs = 1;
for pair in paths.windows(2) {
if drift.unread(&pair[0]) != drift.unread(&pair[1]) {
runs += 1;
}
}
runs
}
#[test]
fn a_reader_puts_part_2_before_part_10() {
use std::cmp::Ordering;
let cmp = |a: &str, b: &str| natural_cmp(a, b);
assert_eq!(
cmp("part=2", "part=10"),
Ordering::Less,
"which bytes do not"
);
assert_eq!(cmp("part=10", "part=2"), Ordering::Greater);
assert_eq!(cmp("date=2024-01-02", "date=2024-01-03"), Ordering::Less);
assert_eq!(cmp("date=2024-01-02", "date=2024-01-02"), Ordering::Equal);
assert_eq!(
cmp("m=03", "m=3"),
Ordering::Equal,
"the same number written two ways is neither before nor after itself — a \
dataset that spells one month both ways is past helping, and this at \
least does not invent an order for it"
);
assert_eq!(cmp("a=1/b=2", "a=1/b=10"), Ordering::Less);
assert_eq!(
cmp("x=a", "x=b"),
Ordering::Less,
"and letters are still letters"
);
}
#[test]
fn a_partition_holds_the_ones_below_it() {
assert!(partition_holds("y=2024", "y=2024/m=03"));
assert!(partition_holds("y=2024", "y=2024"));
assert!(!partition_holds("y=2024", "y=2025"));
assert!(
!partition_holds("y=202", "y=2024"),
"a prefix of the spelling is not a directory above it"
);
assert!(!partition_holds("y=2024/m=03", "y=2024"));
}
#[test]
fn a_file_name_is_not_a_partition_and_neither_is_a_bare_segment() {
assert_eq!(partition_values_of("/x=1/f.parquet"), ["x=1"]);
assert_eq!(
partition_values_of("/x=1/2024=05.parquet"),
["x=1"],
"the file's own name is never a partition, whatever it is called"
);
assert_eq!(
partition_values_of("/raw/x=1/f.parquet"),
["x=1"],
"and a segment with no key before the `=` is not one either"
);
assert_eq!(
partition_values_of("/=1/f.parquet"),
Vec::<String>::new(),
"an empty key is no key"
);
assert_eq!(
partition_values_of("/y=2024/m=03/f.parquet"),
["y=2024", "m=03"],
"in the order written, because a partition is a place"
);
assert_eq!(
partition_values_of("/date=2024=05/f.parquet"),
["date=2024=05"],
"and a value may hold an `=` of its own"
);
}
#[test]
fn a_column_only_a_middle_file_has_is_kept() {
let files = [
file(&[("id", DataType::Int64)], 10),
file(&[("id", DataType::Int64), ("oops", DataType::String)], 10),
file(&[("id", DataType::Int64)], 10),
];
let union = union(&files);
assert_eq!(names(&union.schema), ["id", "oops"]);
let oops = union.columns.iter().find(|c| c.name == "oops").unwrap();
assert_eq!(oops.present_in, 1);
}
#[test]
fn the_newest_files_order_leads_and_older_columns_follow() {
let files = [
file(&[("a", DataType::Int64), ("gone", DataType::Int64)], 1),
file(&[("b", DataType::Int64), ("a", DataType::Int64)], 1),
];
assert_eq!(names(&union(&files).schema), ["b", "a", "gone"]);
}
#[test]
fn integer_widths_widen_losslessly() {
let files = [
file(&[("n", DataType::Int32)], 100),
file(&[("n", DataType::Int64)], 1),
];
let union = union(&files);
assert_eq!(union.schema.get("n"), Some(&DataType::Int64));
assert!(union.columns[0].widened);
assert_eq!(union.columns[0].conflicting_files, 0);
assert!(union.omitted.iter().all(|o| o.is_empty()));
}
#[test]
fn an_integer_and_a_float_meet_at_float64() {
let files = [
file(&[("n", DataType::Int32)], 1),
file(&[("n", DataType::Float32)], 1),
];
assert_eq!(union(&files).schema.get("n"), Some(&DataType::Float64));
}
#[test]
fn datetime_units_widen_to_the_finer_one() {
let ms = DataType::Datetime(TimeUnit::Milliseconds, None);
let ns = DataType::Datetime(TimeUnit::Nanoseconds, None);
let files = [file(&[("t", ms)], 1), file(&[("t", ns.clone())], 1)];
assert_eq!(union(&files).schema.get("t"), Some(&ns));
}
#[test]
fn a_struct_has_every_field_either_file_has() {
let old = DataType::Struct(vec![Field::new("a".into(), DataType::Int32)]);
let new = DataType::Struct(vec![
Field::new("a".into(), DataType::Int64),
Field::new("b".into(), DataType::String),
]);
let files = [file(&[("s", old)], 1), file(&[("s", new.clone())], 1)];
assert_eq!(union(&files).schema.get("s"), Some(&new));
}
#[test]
fn a_type_conflict_goes_to_the_majority_of_rows() {
let files = [
file(&[("price", DataType::String)], 10),
file(&[("price", DataType::Int64)], 90),
];
let union = union(&files);
assert_eq!(union.schema.get("price"), Some(&DataType::Int64));
assert_eq!(union.columns[0].conflicting_files, 1);
assert_eq!(union.columns[0].conflicting_types, [DataType::String]);
assert_eq!(omitted_names(&union, 0), ["price"]);
assert!(union.omitted[1].is_empty());
}
#[test]
fn the_majority_can_be_the_text_files() {
let files = [
file(&[("price", DataType::String)], 90),
file(&[("price", DataType::Int64)], 10),
];
let union = union(&files);
assert_eq!(union.schema.get("price"), Some(&DataType::String));
assert_eq!(omitted_names(&union, 1), ["price"]);
}
#[test]
fn a_type_that_covers_more_files_wins_over_one_that_covers_none_extra() {
let files = [
file(&[("n", DataType::Int32)], 30),
file(&[("n", DataType::Float32)], 30),
file(&[("n", DataType::String)], 50),
];
let union = union(&files);
assert_eq!(union.schema.get("n"), Some(&DataType::Float64));
assert_eq!(omitted_names(&union, 2), ["n"]);
}
#[test]
fn names_differing_only_by_case_stay_two_columns() {
let files = [file(
&[("Price", DataType::Int64), ("price", DataType::Int64)],
1,
)];
assert_eq!(names(&union(&files).schema), ["Price", "price"]);
}
#[test]
fn an_unreadable_footer_is_recorded_and_left_out() {
let files = [
file(&[("id", DataType::Int64)], 1),
None,
file(&[("id", DataType::Int64), ("late", DataType::Int64)], 1),
];
let union = union(&files);
assert_eq!(union.unreadable, [1]);
assert_eq!(names(&union.schema), ["id", "late"]);
assert!(union.omitted[1].is_empty());
}
#[test]
fn a_column_of_nulls_takes_the_other_files_type() {
let files = [
file(&[("x", DataType::Null)], 1),
file(&[("x", DataType::Int64)], 1),
];
let union = union(&files);
assert_eq!(union.schema.get("x"), Some(&DataType::Int64));
assert_eq!(union.columns[0].conflicting_files, 0);
}
#[test]
fn unsigned_and_signed_meet_in_a_wider_signed_type() {
assert_eq!(
widen(&DataType::UInt32, &DataType::Int32),
Some(DataType::Int64)
);
assert_eq!(widen(&DataType::UInt64, &DataType::Int64), None);
}
#[test]
fn origins_read_as_sentences() {
assert_eq!(
SchemaOrigin::AllFooters(6541).to_string(),
"all 6,541 footers"
);
assert_eq!(
SchemaOrigin::FooterSample {
read: 5000,
total: 200_000
}
.to_string(),
"5,000 of 200,000 footers (sample)"
);
}
#[test]
fn a_sample_spans_the_files_and_keeps_the_first_and_newest() {
assert_eq!(footers_to_read(3), [0, 1, 2]);
assert_eq!(footers_to_read(MAX_FOOTER_READS).len(), MAX_FOOTER_READS);
let sample = footers_to_read(MAX_FOOTER_READS * 10);
assert_eq!(sample.len(), MAX_FOOTER_READS);
assert_eq!(sample.first(), Some(&0));
assert_eq!(sample.last(), Some(&(MAX_FOOTER_READS * 10 - 1)));
assert!(sample.windows(2).all(|w| w[0] < w[1]), "ascending");
}
#[test]
fn types_the_scan_cannot_cast_are_not_widened() {
let ms = DataType::Duration(TimeUnit::Milliseconds);
let us = DataType::Duration(TimeUnit::Microseconds);
assert_eq!(widen(&ms, &us), None);
assert_eq!(widen(&DataType::Binary, &DataType::String), None);
assert_eq!(widen(&DataType::Date, &ms), None);
}
#[test]
fn drifting_counts_against_the_files_read_not_the_busiest_column() {
let files = [
file(&[("a", DataType::Int64)], 1),
file(&[("b", DataType::Int64)], 1),
];
let union = union(&files);
let drifting: Vec<_> = union.drifting().map(|c| c.name.to_string()).collect();
assert_eq!(drifting, ["b", "a"]);
}
#[test]
fn drifting_names_only_the_columns_worth_a_note() {
let files = [
file(&[("id", DataType::Int64)], 1),
file(&[("id", DataType::Int64), ("oops", DataType::String)], 1),
];
let union = union(&files);
let drifting: Vec<_> = union.drifting().map(|c| c.name.to_string()).collect();
assert_eq!(drifting, ["oops"]);
}
}