1use crate::matrix::common_io::{open_buf_reader, unquote_field};
24use parquet::file::reader::{FileReader, SerializedFileReader};
25use parquet::record::reader::RowIter;
26use parquet::record::Field;
27use parquet::schema::types::Type as SchemaType;
28use std::fs::File;
29use std::io::BufRead;
30use std::path::Path;
31use std::sync::Arc;
32
33#[derive(Clone, Copy, Debug, PartialEq, Eq)]
35pub enum Delimiter {
36 Tab,
38 Comma,
40 Whitespace,
42}
43
44impl Delimiter {
45 pub fn detect(header: &str) -> Self {
48 if header.contains('\t') {
49 Delimiter::Tab
50 } else if header.contains(',') {
51 Delimiter::Comma
52 } else {
53 Delimiter::Whitespace
54 }
55 }
56
57 pub fn split<'a>(&self, line: &'a str) -> Vec<&'a str> {
59 match self {
60 Delimiter::Tab => line.split('\t').map(|s| unquote_field(s.trim())).collect(),
61 Delimiter::Comma => line.split(',').map(|s| unquote_field(s.trim())).collect(),
62 Delimiter::Whitespace => line.split_whitespace().map(unquote_field).collect(),
63 }
64 }
65}
66
67enum Source {
68 Parquet,
69 Text { delim: Delimiter },
70}
71
72pub struct TableReader {
75 path: Box<str>,
76 header: Vec<Box<str>>,
77 source: Source,
78}
79
80pub fn is_parquet_path(path: &str) -> bool {
82 matches!(
83 Path::new(path)
84 .extension()
85 .and_then(|e| e.to_str())
86 .map(str::to_ascii_lowercase)
87 .as_deref(),
88 Some("parquet") | Some("pq")
89 )
90}
91
92impl TableReader {
93 pub fn open(path: &str) -> anyhow::Result<Self> {
99 if is_parquet_path(path) {
100 let file = File::open(path).map_err(|e| anyhow::anyhow!("opening {path}: {e}"))?;
101 let reader = SerializedFileReader::new(file)?;
102 let header = reader
103 .metadata()
104 .file_metadata()
105 .schema()
106 .get_fields()
107 .iter()
108 .map(|f| Box::from(f.name()))
109 .collect();
110 return Ok(Self {
111 path: path.into(),
112 header,
113 source: Source::Parquet,
114 });
115 }
116 let mut reader =
117 open_buf_reader(path).map_err(|e| anyhow::anyhow!("opening {path}: {e}"))?;
118 let mut line = String::new();
119 anyhow::ensure!(
120 read_header_line(&mut reader, &mut line)?,
121 "{path}: no header line"
122 );
123 let l = line.trim_end_matches(['\n', '\r']);
124 let l = l.strip_prefix('#').unwrap_or(l);
125 let delim = Delimiter::detect(l);
126 let header = delim.split(l).into_iter().map(Box::from).collect();
127 Ok(Self {
128 path: path.into(),
129 header,
130 source: Source::Text { delim },
131 })
132 }
133
134 pub fn path(&self) -> &str {
135 &self.path
136 }
137
138 pub fn header(&self) -> &[Box<str>] {
140 &self.header
141 }
142
143 pub fn column_index(&self, name: &str) -> Option<usize> {
145 self.header.iter().position(|h| h.as_ref() == name)
146 }
147
148 pub fn has_columns(&self, names: &[&str]) -> bool {
150 names.iter().all(|n| self.column_index(n).is_some())
151 }
152
153 pub fn find_column(&self, candidates: &[&str]) -> Option<usize> {
156 candidates.iter().find_map(|c| self.column_index(c))
157 }
158
159 pub fn select(&self, names: &[&str]) -> anyhow::Result<Vec<usize>> {
162 names
163 .iter()
164 .map(|n| {
165 self.column_index(n).ok_or_else(|| {
166 anyhow::anyhow!(
167 "{}: no column `{n}`; the header has {}",
168 self.path,
169 self.header.join(", ")
170 )
171 })
172 })
173 .collect()
174 }
175
176 pub fn rows(&self, cols: &[usize]) -> anyhow::Result<TableRows> {
180 if let Some(&j) = cols.iter().find(|&&j| j >= self.header.len()) {
181 anyhow::bail!(
182 "{}: column index {j} out of range ({} columns)",
183 self.path,
184 self.header.len()
185 );
186 }
187 match &self.source {
188 Source::Parquet => self.parquet_rows(cols),
189 Source::Text { delim } => {
190 let mut reader = open_buf_reader(&self.path)?;
191 let mut line = String::new();
192 read_header_line(&mut reader, &mut line)?;
193 Ok(TableRows::Text {
194 reader,
195 delim: *delim,
196 cols: cols.to_vec(),
197 line: String::new(),
198 })
199 }
200 }
201 }
202
203 fn parquet_rows(&self, cols: &[usize]) -> anyhow::Result<TableRows> {
204 let file = File::open(&*self.path)?;
205 let reader = SerializedFileReader::new(file)?;
206 let root = reader.metadata().file_metadata().schema();
207 let root_name = root.name().to_string();
208 let fields = root.get_fields().to_vec();
209 let mut distinct: Vec<usize> = cols.to_vec();
212 distinct.sort_unstable();
213 distinct.dedup();
214 let proj_fields: Vec<Arc<SchemaType>> =
215 distinct.iter().map(|&j| fields[j].clone()).collect();
216 let proj = SchemaType::group_type_builder(&root_name)
217 .with_fields(proj_fields)
218 .build()?;
219 let pos: Vec<usize> = cols
220 .iter()
221 .map(|j| distinct.binary_search(j).expect("in the projection"))
222 .collect();
223 let iter = RowIter::from_file_into(Box::new(reader)).project(Some(proj))?;
224 Ok(TableRows::Parquet {
225 iter: Box::new(iter),
226 pos,
227 })
228 }
229}
230
231fn read_header_line(reader: &mut dyn BufRead, line: &mut String) -> anyhow::Result<bool> {
234 loop {
235 line.clear();
236 if reader.read_line(line)? == 0 {
237 return Ok(false);
238 }
239 let l = line.trim_end_matches(['\n', '\r']);
240 if !l.starts_with("##") && !l.trim().is_empty() {
241 return Ok(true);
242 }
243 }
244}
245
246pub enum TableRows {
248 Parquet {
249 iter: Box<RowIter<'static>>,
250 pos: Vec<usize>,
251 },
252 Text {
253 reader: Box<dyn BufRead>,
254 delim: Delimiter,
255 cols: Vec<usize>,
256 line: String,
257 },
258}
259
260impl Iterator for TableRows {
261 type Item = anyhow::Result<Vec<Box<str>>>;
262
263 fn next(&mut self) -> Option<Self::Item> {
264 match self {
265 TableRows::Parquet { iter, pos } => {
266 let row = match iter.next()? {
267 Ok(r) => r,
268 Err(e) => return Some(Err(e.into())),
269 };
270 let cells: Vec<&Field> = row.get_column_iter().map(|(_, f)| f).collect();
271 Some(Ok(pos.iter().map(|&p| field_to_string(cells[p])).collect()))
272 }
273 TableRows::Text {
274 reader,
275 delim,
276 cols,
277 line,
278 } => loop {
279 line.clear();
280 match reader.read_line(line) {
281 Ok(0) => return None,
282 Ok(_) => {}
283 Err(e) => return Some(Err(e.into())),
284 }
285 let l = line.trim_end_matches(['\n', '\r']);
286 if l.trim().is_empty() || l.starts_with('#') {
287 continue;
288 }
289 let cells = delim.split(l);
290 return Some(Ok(cols
291 .iter()
292 .map(|&j| Box::from(cells.get(j).copied().unwrap_or("")))
293 .collect()));
294 },
295 }
296 }
297}
298
299pub fn field_to_string(f: &Field) -> Box<str> {
301 match f {
302 Field::Null => "".into(),
303 Field::Str(s) => s.as_str().into(),
304 Field::Bool(b) => b.to_string().into(),
305 Field::Byte(x) => x.to_string().into(),
306 Field::Short(x) => x.to_string().into(),
307 Field::Int(x) => x.to_string().into(),
308 Field::Long(x) => x.to_string().into(),
309 Field::UByte(x) => x.to_string().into(),
310 Field::UShort(x) => x.to_string().into(),
311 Field::UInt(x) => x.to_string().into(),
312 Field::ULong(x) => x.to_string().into(),
313 Field::Float(x) => x.to_string().into(),
314 Field::Double(x) => x.to_string().into(),
315 Field::Float16(x) => x.to_string().into(),
316 Field::Bytes(b) => String::from_utf8_lossy(b.data()).into_owned().into(),
317 other => other.to_string().into(),
318 }
319}
320
321#[cfg(test)]
322#[path = "table_tests.rs"]
323mod tests;