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::formats::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::formats::csv_dialect::head(path, &options, None, false)
.ok()?
.names;
let reader = crate::formats::readers::csv::configure_csv_reader(
crate::formats::readers::csv::csv_reader_of(path).ok()?,
&options,
None,
);
crate::formats::csv_dialect::name_columns(reader.finish().ok()?, header.as_deref())
.ok()?
}
Some(crate::cli::Lines::Json) => {
LazyJsonLineReader::new(crate::cloud::source::polars_literal_path(path).ok()?)
.finish()
.ok()?
}
Some(crate::cli::Lines::Text) => {
return Some(
crate::formats::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 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::formats::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::formats::readers::hive::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;