Files
obikmer/docmd/implementation/chunkreader.md
T
Eric Coissac 00b4b1fa51 docs: add future work section for parallel gzip decompression
Proposes replacing the single-threaded niffler/flate2 pipeline in `obiread::xopen` with `rapidgzip-rs` for local `.gz` files. Details constraints such as path dependencies, non-seekable streams, C++ toolchain requirements, and binding maturity. Marks the optimization as parked pending throughput and correctness validation.
2026-07-07 13:43:01 +02:00

6.8 KiB
Raw Blame History

Chunk reader — implementation

obiread exposes two distinct sequence reading paths, each optimised for a different use case.

Two reading paths

Path API Output unit Per-record identity Use case
Record path read_sequence_chunksparse_chunk SeqRecord (id + raw sequence + normalised rope) yes query — must read complete records
Stream path open_nuc_stream NucPage (flat normalised byte buffer) no index, superkmer — bulk throughput

The record path uses Rope-backed chunks and is described in detail below. The stream path (NucStream / NucPage) is described in the scatter section of pipeline.


Record path: chunk reader

The chunk reader reads FASTA or FASTQ files in fixed-size blocks and yields self-contained chunks, each ending on a complete sequence record boundary. parse_chunk then converts each chunk into a Vec<SeqRecord>, where each record carries its identifier, raw sequence bytes, and a normalised rope ready for superkmer building.

This path is mandatory for query, where superkmers must be tracked back to their originating sequence (id, kmer offset) for output annotation.

Output type: Rope

Each chunk is a Rope — a segmented byte sequence: a Vec of blocks, where each block is a Vec<Cell<u8>>. The consumer iterates over the blocks via a forward or backward cursor.

Rope::split_off(pos) splits at an absolute byte offset in O(log n) (binary search over block-start index). If pos falls inside a block, that block is split in two via Vec::split_off — no memcpy in the common case.

SeqChunkIter

pub struct SeqChunkIter<R: Read> { /* private */ }

impl<R: Read> Iterator for SeqChunkIter<R> {
    type Item = io::Result<Rope>;
}

pub fn fasta_chunks<R: Read>(source: R) -> SeqChunkIter<R>
pub fn fastq_chunks<R: Read>(source: R) -> SeqChunkIter<R>

next() loop:

1. read one block of block_size bytes → push onto Rope
2. call splitter(rope) → Option<abs_offset>
   if Some(pos):
       tail = rope.split_off(pos)    ← O(log n), may split one block
       chunk = mem::replace(&mut rope, tail)
       return Some(Ok(chunk))
3. if EOF and rope non-empty: return Some(Ok(rope)) as final chunk
4. if EOF and rope empty: return None

The Splitter function signature is fn(&Rope) -> Option<usize>. It returns the absolute byte offset of the start of the last complete record, or None if no boundary was found in the accumulated rope (need more data).

Boundary detection — FASTA

Backward scan with a 2-state machine. Searches (right to left) for > followed by \n or \r (i.e., a > that is preceded by a newline in forward order):

stateDiagram-v2
    direction LR
    [*]      --> Scanning
    Scanning --> FoundGt  : '>'
    FoundGt  --> Scanning : other
    FoundGt  --> [*]      : '\\n' / '\\r' ✓

Returns the byte offset of the > that starts the last complete record. Returns None if only one > is found (cannot confirm there is a prior complete record).

Boundary detection — FASTQ

FASTQ records have a rigid 4-line structure (@header, sequence, +, quality). The @ character (ASCII 64, Phred score 31) can appear legitimately in quality lines, making any forward heuristic unreliable. The backward scanner verifies the full structural context before accepting a candidate @.

7-state machine (states 06), scanning from right to left. Each time a + is found, its position is saved as restart; any state mismatch resets the scan to that position.

stateDiagram-v2
    direction LR

    [*]          --> Scanning

    Scanning     --> FoundPlus    : '+' (save restart)
    FoundPlus    --> AfterNlPlus  : '\\n' / '\\r'
    FoundPlus    --> Scanning     : other → backtrack

    AfterNlPlus  --> AfterNlPlus  : séparateur
    AfterNlPlus  --> InSequence   : lettre / - / . / [ / ]
    AfterNlPlus  --> Scanning     : other → backtrack

    InSequence   --> AfterSequence : '\\n' / '\\r'
    InSequence   --> InSequence    : lettre / - / . / [ / ]
    InSequence   --> Scanning      : other → backtrack

    AfterSequence --> AfterSequence : '\\n' / '\\r'
    AfterSequence --> InHeader      : other

    InHeader     --> FoundAt    : '@' (save cut)
    InHeader     --> Scanning   : '\\n' / '\\r' → backtrack
    InHeader     --> InHeader   : other

    FoundAt      --> [*]       : '\\n' / '\\r' ✓
    FoundAt      --> InHeader  : other

restart is updated each time a + is found. When any state fails its expected input, the scan jumps back to restart and continues from there — guaranteeing that a @ in a quality line cannot be accepted as a record start, because the \n+\n structure immediately following it (going backward) will not be found.

Returns the byte offset of the @ that starts the last complete record.


Future work — parallel gzip decompression in xopen

obiread::xopen (xopen.rs) decompresses gzip via nifflerflate2, which is single-threaded (standard DEFLATE has no parallel-decodable structure). For large local gzip inputs this single-threaded decompression can become the throughput bottleneck feeding the query/index/superkmer pipelines, since chunk/page production for a given file is serialized ahead of the worker pool.

Candidate: special-case local, on-disk, gzip-magic-detected paths in open_raw/xopen to use rapidgzip-rs (ReaderBuilder::new().parallelism(n).open(path), implements Read + Seek) instead of niffler, keeping niffler for every other case: stdin (-), HTTP(S) sources, and all non-gzip formats (bzip2, xz, zstd — less used in practice here).

Constraints identified so far (not yet validated against real data):

  • Branch point must move earlier than the current decompress() call in open_raw — rapidgzip's fast path needs the file path, not an already-opened generic Read, so the gzip/local-file detection has to happen before the generic File::open + niffler::send::get_reader path is taken.
  • stdin and HTTP sources are not seekable — they stay on niffler regardless; the gain only applies to local on-disk .gz files.
  • rapidgzip-sys vendors a native C++ engine: requires CMake ≥ 3.17, a C++17 compiler, and nasm on x86 targets — a real build-toolchain addition, not just a pure-Rust crate.
  • Low maturity of the Rust binding at review time (2 GitHub stars, ~15 commits, April 2026 latest release) — the underlying C++ engine is validated (HPDC 2023 paper), but the binding itself has limited production track record.

Decision: parked for now. Before adopting, validate on real data: throughput vs. niffler on representative large .gz inputs, and byte-for-byte correctness of decompressed output.