use color_eyre::Result;
use super::{BASE, EVERYWHERE, Kind, Reader, ScanIn, Signature, Trusted, Unnamed};
#[cfg(feature = "cloud")]
use crate::error_display::FileError;
use crate::export::export_modal::ExportFormat;
use crate::export::python_script::{self as py, Python};
use crate::loading::scan::Scan;
use std::fs::File;
use std::path::{Path, PathBuf};
use arrow::array::types::{
Date32Type, Date64Type, Float32Type, Float64Type, Int8Type, Int16Type, Int32Type, Int64Type,
TimestampMillisecondType, UInt8Type, UInt16Type, UInt32Type, UInt64Type,
};
use arrow::array::{Array, AsArray};
use arrow::record_batch::RecordBatch;
use orc_rust::ArrowReaderBuilder;
use polars::prelude::*;
use super::Read;
use crate::loading::unfinished::Writer;
use crate::{OpenOptions, ParseStringsTarget};
#[cfg(feature = "cloud")]
fn bucket_csv(input: super::BucketIn<'_>) -> Result<polars::prelude::LazyFrame> {
use polars::prelude::{LazyCsvReader, LazyFileListReader};
let super::BucketIn {
url,
path,
cloud,
glob,
options,
format,
} = input;
let named = std::path::Path::new(url);
let failed = |e: polars::prelude::PolarsError| {
FileError::new(
named,
format!("could not read it as {}: {}", format.name(), said(&e)),
)
};
let reader = || {
LazyCsvReader::new(path.clone())
.with_cloud_options(Some(cloud.clone()))
.with_glob(glob)
};
if options.header_rows().is_some() {
return Err(FileError::new(
named,
"--header-rows reads a file's own lines, so it cannot read these in place. Download the files, or name the header with --skip-lines.",
)
.into());
}
let nv = super::csv::build_null_values_with(options, None, || {
super::csv::csv_schema_for_null_values(reader(), options)
})?;
let lf = super::csv::configure_csv_reader(reader(), options, nv.as_ref())
.finish()
.and_then(|lf| crate::formats::csv_dialect::name_columns(lf, None))
.and_then(|lf| {
if !options.skip_initial_space {
return Ok(lf);
}
crate::formats::csv_dialect::skip_initial_space(lf, |column| {
super::csv::csv_null_values_for(options, column)
})
})
.map_err(failed)?;
super::csv::apply_skip_tail_rows_csv(lf, options)
.map_err(|e| crate::error_display::in_file(named, e))
}
#[cfg(feature = "cloud")]
fn bucket_json_lines(input: super::BucketIn<'_>) -> Result<polars::prelude::LazyFrame> {
use polars::prelude::LazyFileListReader;
polars::prelude::LazyJsonLineReader::new(input.path)
.with_cloud_options(Some(input.cloud))
.finish()
.map_err(|e| {
FileError::new(
std::path::Path::new(input.url),
format!("could not read it as {}: {}", input.format.name(), said(&e)),
)
.into()
})
}
#[cfg(feature = "cloud")]
fn bucket_arrow(input: super::BucketIn<'_>) -> Result<polars::prelude::LazyFrame> {
let url = input.url;
let args = polars::prelude::UnifiedScanArgs {
cloud_options: Some(input.cloud),
glob: input.glob,
..Default::default()
};
polars::prelude::LazyFrame::scan_ipc(input.path, Default::default(), args).map_err(|e| {
let folder = url
.split('*')
.next()
.and_then(|head| head.rsplit_once('/'))
.map_or(url, |(folder, _)| folder);
FileError::new(
std::path::Path::new(url),
format!(
"could not read it as Arrow IPC files: {}. A glob reads IPC files in place; Arrow streams are read by their folder: open {folder}/",
said(&e)
),
)
.into()
})
}
#[cfg(feature = "cloud")]
fn said(e: &polars::prelude::PolarsError) -> String {
let said = crate::error_display::user_message_from_polars(e);
said.trim_end_matches('.').to_string()
}
fn frame(read: Read, input: ScanIn<'_>) -> Result<Scan> {
let lf = resolved(read.lf)?;
input.report.read_python = read.python;
input.report.read_notes = read.notes;
input.report.typing = read.typing;
if let (Some(units), Some(delimited)) = (read.units, input.report.delimited.as_mut()) {
let mut merged = (**delimited).clone();
merged.units = units;
*delimited = std::sync::Arc::new(merged);
}
Ok(lf.into())
}
pub(crate) fn resolved(mut lf: LazyFrame) -> Result<LazyFrame> {
lf.collect_schema()?;
Ok(lf)
}
fn json_frame(lf: LazyFrame, input: ScanIn<'_>) -> Result<Scan> {
apply_parse_dates_to_json_lazyframe(resolved(lf)?, input.options, &mut input.report.read_python)
.map(Scan::from)
}
fn scan_parquet(input: ScanIn<'_>) -> Result<Scan> {
frame(each(input.paths, parquet)?.into(), input)
}
fn scan_csv(input: ScanIn<'_>) -> Result<Scan> {
let read = match input.paths {
[one] => super::csv::read_delimited(one, b',', input.options, &Writer::default())?,
many => super::csv::from_csv_paths(many, input.options)?,
};
frame(read, input)
}
fn scan_delimited(input: ScanIn<'_>) -> Result<Scan> {
let separator = input.format.separator().unwrap_or(b',');
let read =
super::csv::read_delimited(input.path(), separator, input.options, &Writer::default())?;
frame(read, input)
}
fn scan_json(input: ScanIn<'_>) -> Result<Scan> {
let lf = each(input.paths, |p| json(p, JsonFormat::Json))?;
json_frame(lf, input)
}
fn scan_json_lines(input: ScanIn<'_>) -> Result<Scan> {
if input.options.follow {
let path = input.paths[0].clone();
return crate::loading::follow::scan_lines(
&path,
input.options,
false,
&mut input.report.read_python,
)
.map(Scan::from);
}
let lf = each(input.paths, |p| json(p, JsonFormat::JsonLines))?;
json_frame(lf, input)
}
fn scan_arrow(input: ScanIn<'_>) -> Result<Scan> {
let paths = input.paths;
if crate::formats::ipc_stream::starts_with_stream(paths) {
return Ok(Scan::Streams(paths.to_vec()));
}
let lf = match paths {
[one] => ipc(one)?,
many => match each(many, ipc).and_then(resolved) {
Ok(lf) => lf,
Err(_) if crate::formats::ipc_stream::any_stream(many) => {
return Ok(Scan::Streams(many.to_vec()));
}
Err(e) => return Err(e),
},
};
frame(lf.into(), input)
}
fn scan_avro(input: ScanIn<'_>) -> Result<Scan> {
frame(each(input.paths, avro)?.into(), input)
}
fn scan_orc(input: ScanIn<'_>) -> Result<Scan> {
frame(each(input.paths, orc)?.into(), input)
}
fn scan_excel(input: ScanIn<'_>) -> Result<Scan> {
let (lf, detail) = crate::formats::excel::read(input.path(), input.options)?;
input.report.opened = Some(std::sync::Arc::new(crate::formats::members::Opened {
detail: Some(std::sync::Arc::new(detail)),
..Default::default()
}));
frame(lf.into(), input)
}
pub(crate) const PARQUET: Reader = Reader {
preview: Some(super::Preview::RowGroup),
python: Some(Python {
call: "pl.scan_parquet",
eager: false,
glob_flag: true,
arguments: Some(py::parquet_arguments),
}),
scan: scan_parquet,
facts: Some(super::Facts {
read: crate::formats::parquet_footer::facts,
footer: true,
}),
signatures: &[Signature {
says: |head, file| {
head.starts_with(b"PAR1") && file.is_none_or(crate::home::discover::has_parquet_magic)
},
kind: Kind::Magic,
trusted: Trusted {
open: Unnamed::NoExtension,
..EVERYWHERE
},
}],
export: Some(ExportFormat::Parquet),
..BASE
};
pub(crate) const CSV: Reader = Reader {
#[cfg(feature = "cloud")]
bucket_scan: Some(bucket_csv),
preview: Some(super::Preview::Scan),
python: Some(Python {
call: "pl.scan_csv",
eager: false,
glob_flag: true,
arguments: Some(py::csv_arguments),
}),
scan: scan_csv,
export: Some(ExportFormat::Csv),
..BASE
};
pub(crate) const TSV: Reader = Reader {
scan: scan_delimited,
#[cfg(feature = "cloud")]
bucket_scan: None,
export: Some(ExportFormat::Tsv),
..CSV
};
pub(crate) const PSV: Reader = Reader {
export: Some(ExportFormat::Psv),
..TSV
};
pub(crate) const JSON: Reader = Reader {
python: Some(Python {
call: "pl.read_json",
eager: true,
glob_flag: false,
arguments: None,
}),
scan: scan_json,
export: Some(ExportFormat::Json),
..BASE
};
pub(crate) const JSONL: Reader = Reader {
#[cfg(feature = "cloud")]
bucket_scan: Some(bucket_json_lines),
preview: Some(super::Preview::Scan),
python: Some(Python {
call: "pl.scan_ndjson",
eager: false,
glob_flag: false,
arguments: Some(py::ndjson_arguments),
}),
scan: scan_json_lines,
export: Some(ExportFormat::Ndjson),
..BASE
};
pub(crate) const ARROW: Reader = Reader {
#[cfg(feature = "cloud")]
bucket_scan: Some(bucket_arrow),
preview: Some(super::Preview::Scan),
python: Some(Python {
call: "pl.scan_ipc",
eager: false,
glob_flag: true,
arguments: Some(py::arrow_arguments),
}),
scan: scan_arrow,
facts: Some(super::Facts {
read: super::facts::arrow,
footer: false,
}),
signatures: &[
Signature {
says: |head, _| head.starts_with(b"ARROW1"),
kind: Kind::Magic,
trusted: Trusted {
open: Unnamed::Never,
..EVERYWHERE
},
},
Signature {
says: |head, file| match file {
Some(file) => crate::formats::ipc_stream::is_stream_file(file),
None => crate::formats::ipc_stream::is_stream_head(head),
},
kind: Kind::Structure,
trusted: Trusted {
open: Unnamed::NoExtension,
..EVERYWHERE
},
},
],
export: Some(ExportFormat::Ipc),
..BASE
};
pub(crate) const AVRO: Reader = Reader {
python: Some(Python {
call: "pl.read_avro",
eager: true,
glob_flag: false,
arguments: None,
}),
scan: scan_avro,
facts: Some(super::Facts {
read: super::facts::avro,
footer: false,
}),
signatures: &[Signature {
says: |head, _| head.starts_with(b"Obj\x01"),
kind: Kind::Magic,
trusted: Trusted {
open: Unnamed::Never,
..EVERYWHERE
},
}],
export: Some(ExportFormat::Avro),
..BASE
};
pub(crate) const ORC: Reader = Reader {
scan: scan_orc,
facts: Some(super::Facts {
read: super::facts::orc,
footer: false,
}),
signatures: &[Signature {
says: |head, _| head.starts_with(b"ORC"),
kind: Kind::Magic,
trusted: Trusted {
pipe: false,
open: Unnamed::Never,
listing: true,
tables: false,
},
}],
..BASE
};
pub(crate) const EXCEL: Reader = Reader {
python: Some(Python {
call: "pl.read_excel",
eager: true,
glob_flag: false,
arguments: Some(py::excel_arguments),
}),
scan: scan_excel,
signatures: &[Signature {
says: crate::formats::excel::is_listable,
kind: Kind::Magic,
trusted: Trusted {
pipe: false,
open: Unnamed::Any,
listing: false,
tables: true,
},
}],
tables: Some(crate::formats::excel::sheets),
..BASE
};
fn each(paths: &[PathBuf], read: impl Fn(&Path) -> Result<LazyFrame>) -> Result<LazyFrame> {
match paths {
[] => Err(color_eyre::eyre::eyre!("No paths provided")),
[one] => read(one),
many => {
let frames = many.iter().map(|p| read(p)).collect::<Result<Vec<_>>>()?;
Ok(concat(frames.as_slice(), Default::default())?)
}
}
}
pub(super) fn parquet(path: &Path) -> Result<LazyFrame> {
let args = ScanArgsParquet {
glob: crate::cloud::source::expands_as_glob(path),
..Default::default()
};
Ok(LazyFrame::scan_parquet(
PlRefPath::try_from_path(path)?,
args,
)?)
}
pub(super) fn ipc(path: &Path) -> Result<LazyFrame> {
let args = UnifiedScanArgs {
glob: crate::cloud::source::expands_as_glob(path),
..Default::default()
};
Ok(LazyFrame::scan_ipc(
PlRefPath::try_from_path(path)?,
Default::default(),
args,
)?)
}
pub(super) fn avro(path: &Path) -> Result<LazyFrame> {
let file = File::open(path)?;
Ok(polars::io::avro::AvroReader::new(file).finish()?.lazy())
}
pub(super) fn orc(path: &Path) -> Result<LazyFrame> {
let file = File::open(path)?;
let reader = ArrowReaderBuilder::try_new(file)
.map_err(|e| color_eyre::eyre::eyre!("ORC: {}", e))?
.build();
let batches: Vec<RecordBatch> = reader
.collect::<std::result::Result<Vec<_>, _>>()
.map_err(|e| color_eyre::eyre::eyre!("ORC: {}", e))?;
Ok(arrow_record_batches_to_dataframe(&batches)?.lazy())
}
fn json(path: &Path, format: JsonFormat) -> Result<LazyFrame> {
let file = File::open(path)?;
Ok(JsonReader::new(file)
.with_json_format(format)
.finish()?
.lazy())
}
pub(crate) fn union_of_files() -> polars::prelude::UnionArgs {
polars::prelude::UnionArgs {
diagonal: true,
to_supertypes: true,
..Default::default()
}
}
fn arrow_record_batches_to_dataframe(batches: &[RecordBatch]) -> Result<DataFrame> {
if batches.is_empty() {
return Ok(DataFrame::empty());
}
let mut all_dfs = Vec::with_capacity(batches.len());
for batch in batches {
let n_cols = batch.num_columns();
let schema = batch.schema();
let mut series_vec = Vec::with_capacity(n_cols);
for (i, col) in batch.columns().iter().enumerate() {
let name = schema.field(i).name().as_str();
let s = arrow_array_to_polars_series(name, col)?;
series_vec.push(s.into());
}
let df = DataFrame::new_infer_height(series_vec)?;
all_dfs.push(df);
}
let mut out = all_dfs.remove(0);
for df in all_dfs {
out = out.vstack(&df)?;
}
Ok(out)
}
fn arrow_array_to_polars_series(name: &str, array: &dyn Array) -> Result<Series> {
use arrow::datatypes::DataType as ArrowDataType;
let strings = |values: Vec<Option<&str>>| Series::new(name.into(), values);
match array.data_type() {
ArrowDataType::Int8 => primitive::<Int8Type>(name, array, "Int8"),
ArrowDataType::Int16 => primitive::<Int16Type>(name, array, "Int16"),
ArrowDataType::Int32 => primitive::<Int32Type>(name, array, "Int32"),
ArrowDataType::Int64 => primitive::<Int64Type>(name, array, "Int64"),
ArrowDataType::UInt8 => primitive::<UInt8Type>(name, array, "UInt8"),
ArrowDataType::UInt16 => primitive::<UInt16Type>(name, array, "UInt16"),
ArrowDataType::UInt32 => primitive::<UInt32Type>(name, array, "UInt32"),
ArrowDataType::UInt64 => primitive::<UInt64Type>(name, array, "UInt64"),
ArrowDataType::Float32 => primitive::<Float32Type>(name, array, "Float32"),
ArrowDataType::Float64 => primitive::<Float64Type>(name, array, "Float64"),
ArrowDataType::Date32 => primitive::<Date32Type>(name, array, "Date32"),
ArrowDataType::Date64 => primitive::<Date64Type>(name, array, "Date64"),
ArrowDataType::Timestamp(_, _) => {
primitive::<TimestampMillisecondType>(name, array, "Timestamp")
}
ArrowDataType::Boolean => {
let a = array
.as_boolean_opt()
.ok_or_else(|| color_eyre::eyre::eyre!("ORC: expected Boolean array"))?;
Ok(Series::new(name.into(), a.iter().collect::<Vec<_>>()))
}
ArrowDataType::Utf8 => {
let a = array
.as_string_opt::<i32>()
.ok_or_else(|| color_eyre::eyre::eyre!("ORC: expected Utf8 array"))?;
Ok(strings(a.iter().collect()))
}
ArrowDataType::LargeUtf8 => {
let a = array
.as_string_opt::<i64>()
.ok_or_else(|| color_eyre::eyre::eyre!("ORC: expected LargeUtf8 array"))?;
Ok(strings(a.iter().collect()))
}
other => Err(color_eyre::eyre::eyre!(
"ORC: unsupported column type {:?} for column '{}'",
other,
name
)),
}
}
fn primitive<T: arrow::datatypes::ArrowPrimitiveType>(
name: &str,
array: &dyn Array,
type_name: &str,
) -> Result<Series>
where
Series: NamedFrom<Vec<Option<T::Native>>, [Option<T::Native>]>,
{
let a = array
.as_primitive_opt::<T>()
.ok_or_else(|| color_eyre::eyre::eyre!("ORC: expected {type_name} array"))?;
Ok(Series::new(name.into(), a.iter().collect::<Vec<_>>()))
}
pub(crate) fn apply_parse_dates_to_json_lazyframe(
lf: LazyFrame,
options: &OpenOptions,
read: &mut Vec<String>,
) -> Result<LazyFrame> {
if !options.parse_dates {
return Ok(lf);
}
super::csv::type_string_columns(
lf,
&ParseStringsTarget::All,
options.parse_strings_sample_rows,
super::csv::StringTypes {
dates: true,
numbers: false,
},
read,
&[],
&mut Vec::new(),
)
}
#[cfg(test)]
mod reader_errors {
use crate::FileFormat;
use crate::formats::readers::bad_input::each_names_its_file;
#[test]
fn errors_name_the_file() {
let garbage: &[u8] = b"\x00\x01\x02 this is not a file of any format \xff\xfe";
for (format, name) in [
(FileFormat::Parquet, "bad.parquet"),
(FileFormat::Arrow, "bad.arrow"),
(FileFormat::Avro, "bad.avro"),
(FileFormat::Orc, "bad.orc"),
(FileFormat::Excel, "bad.xlsx"),
(FileFormat::Json, "bad.json"),
(FileFormat::Jsonl, "bad.jsonl"),
] {
each_names_its_file(format, &[(name, garbage, "")]);
}
each_names_its_file(FileFormat::Csv, &[("ragged.csv", b"a,b\n1,2\n\"3,4\n", "")]);
}
}