use std::fs::File;
use std::io::Cursor;
use std::num::NonZeroUsize;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use polars::prelude::*;
use super::Spool;
const SCAN_NAME: &str = "NDJSON";
pub(crate) struct LinesScan {
path: PathBuf,
schema: SchemaRef,
ignore_errors: bool,
spool: Option<Arc<Spool>>,
held: Option<Arc<File>>,
}
pub(crate) fn complete(bytes: &[u8], ended: bool) -> usize {
if ended {
return bytes.len();
}
memchr::memrchr(b'\n', bytes).map_or(0, |at| at + 1)
}
fn after_rows(bytes: &[u8], rows: usize) -> usize {
let mut start = 0;
let mut seen = 0;
while start < bytes.len() && seen < rows {
let end = memchr::memchr(b'\n', &bytes[start..]).map_or(bytes.len(), |at| start + at);
if polars::io::ndjson::core::is_json_line(&bytes[start..end]) {
seen += 1;
}
start = end + 1;
}
start.min(bytes.len())
}
const RUN: usize = 64 << 20;
fn run_end(bytes: &[u8], start: usize, run: usize) -> usize {
let at = start.saturating_add(run);
if at >= bytes.len() {
return bytes.len();
}
memchr::memchr(b'\n', &bytes[at..]).map_or(bytes.len(), |n| at + n + 1)
}
fn map(file: &File) -> std::io::Result<Option<memmap2::Mmap>> {
if file.metadata()?.len() == 0 {
return Ok(None);
}
Ok(Some(unsafe { memmap2::Mmap::map(file)? }))
}
impl LinesScan {
pub(crate) fn open(
path: &Path,
infer: Option<NonZeroUsize>,
ignore_errors: bool,
spool: Option<Arc<Spool>>,
) -> PolarsResult<LinesScan> {
let ended = spool.as_ref().is_some_and(|s| s.ended().is_some());
let file = File::open(path)?;
let map = map(&file)?;
let bytes = map.as_deref().unwrap_or_default();
let bytes = &bytes[..complete(bytes, ended)];
let bytes = bytes.strip_prefix(b"\xef\xbb\xbf").unwrap_or(bytes);
let schema = match polars::io::ndjson::infer_schema(&mut Cursor::new(bytes), infer) {
Err(_) if ignore_errors => {
let mut objects = Vec::new();
let wanted = infer.map_or(usize::MAX, NonZeroUsize::get);
let lines = bytes
.split(|&b| b == b'\n')
.filter(|line| {
matches!(
serde_json::from_slice::<serde_json::Value>(line),
Ok(serde_json::Value::Object(_))
)
})
.take(wanted);
for line in lines {
objects.extend_from_slice(line);
objects.push(b'\n');
}
polars::io::ndjson::infer_schema(&mut Cursor::new(objects), infer)?
}
inferred => inferred?,
};
Ok(LinesScan {
path: path.to_path_buf(),
schema: Arc::new(schema),
ignore_errors,
spool,
held: None,
})
}
pub(crate) fn lazy(self) -> PolarsResult<LazyFrame> {
let schema = self.schema.clone();
LazyFrame::anonymous_scan(
Arc::new(self),
ScanArgsAnonymous {
schema: Some(schema),
name: SCAN_NAME,
..Default::default()
},
)
}
pub(super) fn in_plan(scan: &polars::lazy::dsl::DslPlan) -> Option<&LinesScan> {
use polars::lazy::dsl::{DslPlan, FileScanDsl};
let DslPlan::Scan { scan_type, .. } = scan else {
return None;
};
let FileScanDsl::Anonymous { function, .. } = &**scan_type else {
return None;
};
function.as_any().downcast_ref::<LinesScan>()
}
pub(super) fn of<'a>(
scan: &'a polars::lazy::dsl::DslPlan,
path: &str,
) -> Option<&'a LinesScan> {
Self::in_plan(scan)
.filter(|s| s.held.is_none() && super::same_file(&s.path.to_string_lossy(), path))
}
pub(super) fn ignore_errors(&self) -> bool {
self.ignore_errors
}
pub(crate) fn schema(&self) -> &SchemaRef {
&self.schema
}
pub(super) fn held(&self, file: &File) -> Option<LinesScan> {
Some(LinesScan {
held: Some(Arc::new(file.try_clone().ok()?)),
..self.with_schema(self.schema.clone())
})
}
pub(crate) fn with_schema(&self, schema: SchemaRef) -> LinesScan {
LinesScan {
path: self.path.clone(),
schema,
ignore_errors: self.ignore_errors,
spool: self.spool.clone(),
held: self.held.clone(),
}
}
}
impl AnonymousScan for LinesScan {
fn as_any(&self) -> &dyn std::any::Any {
self
}
fn schema(&self, _infer_schema_length: Option<usize>) -> PolarsResult<SchemaRef> {
Ok(self.schema.clone())
}
fn allows_predicate_pushdown(&self) -> bool {
true
}
fn allows_projection_pushdown(&self) -> bool {
true
}
fn scan(&self, mut args: AnonymousScanArgs) -> PolarsResult<DataFrame> {
args.predicate = crate::formats::pushdown::evaluable(args.predicate.take());
let ended = self.spool.as_ref().is_some_and(|s| s.ended().is_some());
let file = match &self.held {
Some(held) => held.try_clone()?,
None => File::open(&self.path)?,
};
let columns: Schema = match args.with_columns.as_deref() {
Some(names) => names
.iter()
.filter_map(|name| self.schema.get_field(name))
.collect(),
None => (*self.schema).clone(),
};
let map = map(&file)?;
let bytes = map.as_deref().unwrap_or_default();
let mut bytes = &bytes[..complete(bytes, ended)];
if let Some(rows) = args.n_rows {
bytes = &bytes[..after_rows(bytes, rows)];
}
if columns.is_empty() {
let df = DataFrame::empty_with_height(polars::io::ndjson::count_rows(bytes));
return match args.predicate {
Some(predicate) => df.lazy().filter(predicate).collect(),
None => Ok(df),
};
}
let columns = Arc::new(columns);
let mut out = DataFrame::empty_with_schema(&columns);
let mut start = 0;
while start < bytes.len() {
let end = run_end(bytes, start, RUN);
let mut df = parse_run(&bytes[start..end], &columns, self.ignore_errors)?;
if let Some(predicate) = &args.predicate {
df = df.lazy().filter(predicate.clone()).collect()?;
}
out.vstack_mut(&df)?;
start = end;
}
out.rechunk_mut();
Ok(out)
}
}
pub(super) fn parse_run(
run: &[u8],
columns: &Schema,
ignore_errors: bool,
) -> PolarsResult<DataFrame> {
let run = run.strip_prefix(b"\xef\xbb\xbf").unwrap_or(run);
match polars::io::ndjson::core::parse_ndjson(run, None, columns, ignore_errors) {
Err(_) if ignore_errors => {}
read => return read,
}
let mut out = DataFrame::empty_with_schema(columns);
for line in run.split(|&b| b == b'\n') {
if !polars::io::ndjson::core::is_json_line(line) {
continue;
}
let row = polars::io::ndjson::core::parse_ndjson(line, Some(1), columns, true)
.unwrap_or_else(|_| DataFrame::full_null(columns, 1));
out.vstack_mut(&row)?;
}
Ok(out)
}
#[cfg(test)]
mod tests;