Skip to main content

legume_numeric/matrix/
table.rs

1//! Row-streaming reader over a named-column table: parquet, or delimited
2//! text (`.tsv`, `.csv`, `.txt`, gzip or not), with columns picked by
3//! header name.
4//!
5//! Every cell comes back as a string, whatever the storage type: a parquet
6//! `INT64` or `DOUBLE` cell is formatted, a null is `""`. Callers that want a
7//! number parse it; callers that want a label get the value rather than a
8//! blank. Rows stream, so a multi-gigabyte association table is never held
9//! in memory.
10//!
11//! ```no_run
12//! use legume_numeric::matrix::table::TableReader;
13//! let t = TableReader::open("pairs.parquet")?;
14//! let cols = t.select(&["variant_id", "gene_id", "pval_nominal"])?;
15//! for row in t.rows(&cols)? {
16//!     let row = row?;
17//!     let (variant, gene, p) = (&row[0], &row[1], &row[2]);
18//! #   let _ = (variant, gene, p);
19//! }
20//! # Ok::<(), anyhow::Error>(())
21//! ```
22
23use 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/// How the cells of a text line are separated.
34#[derive(Clone, Copy, Debug, PartialEq, Eq)]
35pub enum Delimiter {
36    /// One tab per boundary; empty cells are kept.
37    Tab,
38    /// One comma per boundary; empty cells are kept, quotes stripped.
39    Comma,
40    /// Runs of spaces or tabs; empty cells cannot occur.
41    Whitespace,
42}
43
44impl Delimiter {
45    /// The delimiter a header line uses: a tab if it has one, else a comma,
46    /// else whitespace.
47    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    /// Split one line into cells, trimmed and unquoted.
58    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
72/// A table opened for streaming: its header is read up front, its rows on
73/// demand through [`TableReader::rows`].
74pub struct TableReader {
75    path: Box<str>,
76    header: Vec<Box<str>>,
77    source: Source,
78}
79
80/// `true` when the path names a parquet file (`.parquet` / `.pq`).
81pub 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    /// Open `path` and read its header. Parquet is recognised by extension;
94    /// anything else is delimited text, gzip-decoded when it ends in `.gz`.
95    /// A text header is the first line that does not start with `##` (VCF-
96    /// style meta lines are skipped), with one leading `#` dropped
97    /// (`#chr start end` is a BED-style header).
98    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    /// Column names, in file order.
139    pub fn header(&self) -> &[Box<str>] {
140        &self.header
141    }
142
143    /// Index of a column by exact name.
144    pub fn column_index(&self, name: &str) -> Option<usize> {
145        self.header.iter().position(|h| h.as_ref() == name)
146    }
147
148    /// `true` when every name is a column.
149    pub fn has_columns(&self, names: &[&str]) -> bool {
150        names.iter().all(|n| self.column_index(n).is_some())
151    }
152
153    /// The first of `candidates` that is a column: for a field that
154    /// different releases spell differently.
155    pub fn find_column(&self, candidates: &[&str]) -> Option<usize> {
156        candidates.iter().find_map(|c| self.column_index(c))
157    }
158
159    /// Indices of `names`, in request order; a missing name is an error
160    /// that lists what the file does have.
161    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    /// Stream the rows, each cut down to the columns at `cols` (in that
177    /// order; repeats allowed). A text row too short for a requested column
178    /// yields `""` there; blank lines and `#` lines are skipped.
179    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        // Project onto the distinct requested columns in file order, then map
210        // every request onto its position in the projection.
211        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
231/// Read up to and including the header line: the first non-blank line not
232/// starting with `##`, left in `line`. `false` at end of file.
233fn 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
246/// The row stream of a [`TableReader`].
247pub 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
299/// A parquet cell as text: strings verbatim, numbers formatted, null `""`.
300pub 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;