use rayon::prelude::*;
use super::window::count_newlines;
pub(super) fn tokens(line: &str) -> impl Iterator<Item = &str> {
line.split(|c: char| c == ',' || c.is_whitespace())
.filter(|t| !t.is_empty())
}
pub(super) fn is_skipped(line: &str) -> bool {
let trimmed = line.trim_start();
trimmed.is_empty() || trimmed.starts_with('#')
}
#[inline]
pub(super) fn is_separator(b: u8) -> bool {
matches!(b, b' ' | b'\t' | b',' | b'\r' | 0x0b | 0x0c)
}
fn token_end(bytes: &[u8], start: usize) -> usize {
let mut end = start;
while end < bytes.len() && !is_separator(bytes[end]) {
end += 1;
}
end
}
pub(super) fn parse_numbers(line: &str, push: impl FnMut(f64)) -> std::result::Result<(), &str> {
if !line.is_ascii() {
return parse_unicode_numbers(line, push);
}
parse_ascii_numbers(line, push)
}
fn parse_unicode_numbers(line: &str, mut push: impl FnMut(f64)) -> std::result::Result<(), &str> {
for token in tokens(line) {
push(token.parse::<f64>().map_err(|_| token)?);
}
Ok(())
}
fn parse_ascii_numbers(line: &str, mut push: impl FnMut(f64)) -> std::result::Result<(), &str> {
let bytes = line.as_bytes();
let n = bytes.len();
let mut i = 0;
while i < n {
while i < n && is_separator(bytes[i]) {
i += 1;
}
if i == n {
break;
}
match fast_float2::parse_partial::<f64, _>(&bytes[i..]) {
Ok((value, used)) if i + used == n || is_separator(bytes[i + used]) => {
push(value);
i += used;
}
_ => return Err(&line[i..token_end(bytes, i)]),
}
}
Ok(())
}
pub(super) const MIN_PARALLEL_BYTES: usize = 1 << 20;
pub(super) fn parse_workers(text: &str, threads: usize, min_bytes: usize) -> usize {
if threads > 1 && text.len() >= min_bytes {
threads
} else {
1
}
}
pub(super) fn line_chunks(text: &str, workers: usize) -> Vec<&str> {
let bytes = text.as_bytes();
let mut chunks = Vec::with_capacity(workers);
let mut start = 0;
for k in 1..workers {
let mut end = (bytes.len() / workers * k).max(start);
while end < bytes.len() && bytes[end] != b'\n' {
end += 1;
}
if end < bytes.len() {
end += 1;
}
chunks.push(&text[start..end]);
start = end;
}
chunks.push(&text[start..]);
chunks
}
pub(super) fn map_line_chunks<T: Send>(
pool: Option<&rayon::ThreadPool>,
text: &str,
first_line: usize,
workers: usize,
parse_chunk: impl Fn(usize, &str) -> T + Send + Sync,
) -> (Vec<T>, usize) {
let chunks = line_chunks(text, workers);
let newlines = |chunk: &&str| count_newlines(chunk.as_bytes());
let counts: Vec<usize> = match pool {
Some(pool) => pool.install(|| chunks.par_iter().map(newlines).collect()),
None => chunks.iter().map(newlines).collect(),
};
let mut lineno = first_line;
let first_lines: Vec<usize> = counts
.iter()
.map(|count| {
let first = lineno;
lineno += count;
first
})
.collect();
let results = match pool {
Some(pool) => pool.install(|| {
chunks
.par_iter()
.zip(&first_lines)
.map(|(chunk, &first)| parse_chunk(first, chunk))
.collect()
}),
None => chunks
.iter()
.zip(&first_lines)
.map(|(chunk, &first)| parse_chunk(first, chunk))
.collect(),
};
(results, lineno - first_line)
}
pub(super) struct Parse {
threads: usize,
pub(super) min_bytes: usize,
pool: Option<rayon::ThreadPool>,
}
impl Parse {
pub(super) fn new(threads: usize) -> Self {
Self {
threads,
min_bytes: MIN_PARALLEL_BYTES,
pool: None,
}
}
pub(super) fn map<T: Send>(
&mut self,
text: &str,
first_line: usize,
parse_chunk: impl Fn(usize, &str) -> T + Send + Sync,
) -> (Vec<T>, usize) {
let workers = parse_workers(text, self.threads, self.min_bytes);
if workers <= 1 {
let result = parse_chunk(first_line, text);
return (vec![result], count_newlines(text.as_bytes()));
}
if self.pool.is_none() {
self.pool = rayon::ThreadPoolBuilder::new()
.num_threads(workers)
.build()
.ok();
}
map_line_chunks(self.pool.as_ref(), text, first_line, workers, parse_chunk)
}
}