Part of the MuMDIA developer documentation (see docs/README.md).
The mumdia-io crate is the on-disk contract layer for the whole engine. Every
stage reads its path-addressable inputs and writes its outputs through this
crate, so no stage hand-rolls Arrow RecordBatches or touches the Parquet
reader/writer directly. The crate provides five things:
- A small typed column model (
Col) and a read-back table (Table) over Arrow + Parquet, so a stage declares its schema as aVec<Col>and reads it back by column name with typed getters (table.rs). - Parquet write (
write_table) and read (Table::read), SNAPPY-compressed, the open self-describing interstage format (see the interstage contract indocs/18_findings_and_decisions.md, B1). - The per-artifact
<artifact>.report.jsonsidecar (ArtifactReportinreport.rs) so a stage can be evaluated (row counts, resolved params, key distributions, model identity, timing) without loading the full table. - blake3 content hashing of files and strings (
hash.rs), which feeds thecontent_hashon every artifact record and theconfig_hashrecorded on each manifestArtifactRecord(manifest.rs:19) and on the manifest header (manifest.rs:27). These hashes are provenance only: nothing in the engine reads them back to invalidate or skip a downstream artifact, andrunalways recomputes the full chain. (report.jsonitself carries noconfig_hashfield; see the report fields below.) - The
mumdia inspect <artifact>implementation (inspectinlib.rs): schema + head sample + row count for any Parquet file.
This crate reads no Config fields of its own. It is a mechanism layer; the
concrete per-artifact column schemas live in each stage's write_table call,
and the frozen schema identifiers live in mumdia-core (schema.rs). The only
"configuration" here is compile-time (SNAPPY compression, a 64 KiB hash buffer,
the schema-id tuples).
| path | role |
|---|---|
rust/mumdia/crates/mumdia-io/src/lib.rs |
crate root: init_logging, record_artifact, inspect; re-exports the modules |
rust/mumdia/crates/mumdia-io/src/table.rs |
Col enum (write side), write_table, Table (read side) and the typed getters |
rust/mumdia/crates/mumdia-io/src/span_cache.rs |
SpanCache, the coalescing ChunkReader behind TableFile::scan with ScanOptions::coalesced() |
rust/mumdia/crates/mumdia-io/src/codec.rs |
ColumnEncoder (parallel column encode), the codec pool (CodecPool, claim, codec_threads, set_codec_threads) shared by encode and parallel decode |
rust/mumdia/crates/mumdia-io/src/report.rs |
ArtifactReport struct + write_for (the .report.json sidecar) |
rust/mumdia/crates/mumdia-io/src/hash.rs |
blake3_file, blake3_str, HashingWrite (hash on write) |
rust/mumdia/crates/mumdia-io/src/json.rs |
write_json, read_json (pretty JSON via serde) |
rust/mumdia/crates/mumdia-core/src/schema.rs |
frozen (logical name, schema version) tuples for every artifact |
rust/mumdia/crates/mumdia-core/src/manifest.rs |
ArtifactRecord / Manifest (populated from record_artifact) |
The IO layer is schema-agnostic: it does not itself consume or produce a fixed set of named artifacts. It is the read/write mechanism that every stage uses. What is fixed here is the artifact identity registry and the shape of the two JSON sidecars.
Each artifact carries a logical schema name and version. The intent is that a
stage could validate its inputs and refuse to apply a model under a mismatched
schema, but no such read-side check is implemented: schema_version is only
ever written into the report and the manifest, never read back and compared.
They are pub const tuples in the artifact submodule, referenced as
mumdia_core::schema::artifact::PEPTIDES and so on. The tuples are
(name, version):
| logical name | version | producing stage |
|---|---|---|
spectra_ms1 |
1 | convert |
spectra_ms2 |
1 | convert |
isolation_windows |
1 | convert |
ms2_to_ms1 |
1 | convert |
peptides |
1 | digest |
peptidoforms |
1 | peptidoforms |
fragment_library_precursors |
1 | predict-frag / library-input |
fragment_library_fragments |
1 | predict-frag / library-input |
seed_psms |
1 | search-seed |
run_windows |
1 | rt-im-train |
psms_extracted |
1 | extract |
chromatograms |
1 | extract |
chromatograms |
2 (CHROMATOGRAMS_V2, only under extract.chromatogram_schema = 2) |
extract |
features |
1 | features |
psms_competed |
2 | compete |
psms_scored |
3 | rescore |
peptide_quant |
2 | quant |
protein_group_quant |
2 | quant |
fragment_quant |
1 | quant |
There are 18 registered artifacts. Four have been version-bumped past 1:
psms_competed, peptide_quant, and protein_group_quant are at version 2, and
psms_scored is at version 3; every other artifact is still version 1. The
version is recorded in the report and the manifest but is not consumed anywhere:
no stage reads a prior artifact's version back, so there is no schema-mismatch
gate on read. A reader could compare the recorded version by hand, but the engine
does not.
The IO crate stores no schema definitions; a stage's write_table(path, vec![ Col::… ]) is the schema. Two examples read from the actual code:
isolation_windows (stages/convert.rs:215-226), one row per distinct window:
window_id: u32, target: f64, lower: f64, upper: f64.
ms2_to_ms1 (stages/convert.rs:228-234): ms2_scan_index: u32,
ms1_scan_index: i32.
peptides (stages/digest.rs:286-298): id: u32, peptide: utf8,
protein: utf8, start: i32, end: i32, label: utf8, target_id: i32,
decoy_strategy: utf8.
To see any artifact's real schema, run mumdia inspect <artifact.parquet> (see
below) rather than trusting a doc; the schema is authoritative on disk.
<artifact>.report.json(written for most primary stage Parquets viaArtifactReport::write_for,report.rs:28). Coverage is not universal: schema/PIN files, report TSVs, and several optional or Python-written artifacts have no report. Fields below.manifest.json(written once by therunorchestrator from the collectedArtifactRecords,stages/run.rs). This crate supplies the per-record builderrecord_artifact(lib.rs:20); it does not write the manifest file.
Col (table.rs:23-42) is an enum with one variant per supported column type.
Each variant carries (String name, Vec<values>). The scalar variants are
I64, I32, U32, F64, F32, Bool, Str; the nullable variants are
OptF64, OptF32, OptI32, OptStr (each a Vec<Option<T>>); the list
variants are ListF32, ListF64, and LargeListF32. LargeListF32
(table.rs:41) is encoded as an Arrow LargeList with 64-bit offsets, needed
when the total list-value count across all rows can exceed the ~2.1 billion
limit of a 32-bit ListArray offset buffer (for example per-fragment
chromatograms when extraction accepts a very large candidate set). The Opt*
variants exist to back conditional and ion-mobility columns under the
missing-value policy (table.rs:5-6; see docs/18_findings_and_decisions.md);
ion-mobility columns are written null throughout the 3D MVP.
Col has four private helpers used by the writer (all non-pub):
name()(table.rs:45): the column name, via a match over every variant.len()(table.rs:64): the row count of the innerVec.field()(table.rs:83): the ArrowField. Scalar variants arenullable = false; theOpt*and all list variants arenullable = true. List inner items are declaredField::new("item", Float32/64, true).into_array()(table.rs:107): the consuming conversion into anArrayRef. It moves theVecinto the Arrow array instead of cloning, so the column data is copied only once during a write. ScalarVec<T>andVec<Option<T>>go straight throughPrimitiveArray::from; list variants are built with aListBuilder/LargeListBuilder, appending each row's slice and thenappend(true)(a present, possibly empty, list).
write_table(path, cols) (table.rs:151) is the single write entry point and
returns the row count as u64:
- Reject an empty column set (
table.rs:152). - Reject duplicate column names via a
HashSet(table.rs:157-165). Arrow allows duplicate names but readers resolve a name to the first match, which would silently hide the second column, so this is a hard error. - Take
nrowsfrom column 0 and require every column to match (table.rs:166-176); a mismatch is a hard error naming the offending column. - Build the
Schemafromfield()over all columns (table.rs:177-178), then consume the columns intoArrayRefs (table.rs:181). Thefieldsvector is captured before the consume so the schema still has everything. RecordBatch::try_new(table.rs:182), create parent dirs (create_dir_all(...).ok(), best-effort,table.rs:185-187), create the file, buildWriterPropertieswithCompression::SNAPPY(table.rs:189-191), and write a single batch throughArrowWriter, thenclose()(table.rs:192-194). Onewrite_tablecall produces exactly one row group / one logical batch.
The list above describes the original single-batch write. write_table now hands
the writer the rows in WRITE_TABLE_CHUNK_ROWS (65,536-row) chunks through a
TableWriter, which keeps one chunk of Arrow arrays resident instead of a second
copy of the whole table; the row groups still fall at the writer's 1,048,576-row
default. write_table_chunked(path, nrows, chunk) is the same write for a caller
that produces its rows chunk by chunk: chunk(start..end) is called for exactly the
ranges write_table would cut (one empty range for an empty table), so the file is
byte-identical to write_table over the concatenated columns, including the page
framing of an all-null column, which a different chunk sequence would move
(write_table_chunked_writes_the_write_table_file). write_table_chunked_hashed
is the same write, hashed as it is written ("Hash on write" below). rt-im-train
writes run_windows.parquet this way, so the seven whole columns (76 bytes per
candidate) are never resident.
A writer with a row-group cap (TableWriter::with_row_group_rows,
BatchWriter::with_row_group_rows, WriteOptions::row_group_rows) cuts its data
pages by size only: writer_props sets the data page row limit to the cap
(at most MAX_DATA_PAGE_ROWS, 131,072), where parquet-rs cuts every 20,000
rows. A scalar column of a 65,536-row group is then one page instead of four,
and the plain reader, which fetches every page with its own seek, reads it with
one seek. Pages still end at parquet's 1 MB data page size, so list leaves
barely change. Measured on the AIF artifacts at their own row-group sizes
(bench_rewrite_a_real_artifact): features 2,003 -> 1,039 data pages at +0.5%
bytes, psms_competed 1,592 -> 399 at +0.5%, chromatograms 1,141 -> 947 at
+0.05%. The writer's in-progress buffers grow with the page: 551.9 -> 620.2 MB
for one 131,072-row group of the 398-column competed table. Values are
unchanged; file bytes and content hashes of capped writers change; uncapped
writers (write_table, write_batches, BatchWriter::new) are byte-identical
to parquet-rs's defaults. This is step 3 of R1 in the 2026-09-25 performance
survey.
A capped writer does not decide its float dictionaries blind. Before it encodes
a byte it holds its first cap / 4 rows (plan_sample_rows, whatever chunks
they arrive in), and EncodingPlan::of looks at every FLOAT and DOUBLE leaf,
scalar or list item:
- a leaf with at least 4,096 sampled values of which more than 80% are distinct
(
PLAN_PLAIN_ABOVE_DISTINCT, counted by bit pattern as parquet's dictionary interns them) is written PLAIN from its first page; - every other float leaf keeps the c = 0.5 dictionary limit of
writer_props, sized from the leaf's values per row group (sampled values per row times the cap) instead of its rows, so a list leaf keeps parquet's 1 MB default where the row-sized limit cut the chromatogram traces' dictionaries at 128 KB.
The held batches are then encoded in the order and the chunks they arrived in.
The unplanned rule fell back to PLAIN only after the dictionary filled, and the
pages before the fallback kept their dictionary; the plan does not pay for that
prefix. Measured on the AIF artifacts at their own row-group sizes
(bench_rewrite_a_real_artifact, against the unplanned layout): features
-8.6%, psms_competed -14.9%, chromatograms -6.3%, spectra_ms2 -0.2%. Encoding
is also cheaper where a column skips dictionary interning: the competed
rewrite took 0.54 s against 1.31 s and its writer peak memory_size fell from
620 to 297 MB; the features rewrite 0.65 s against 1.00 s, 316 to 269 MB.
Values are unchanged; bytes and content hashes of capped writers change;
uncapped writers do not plan and are byte-identical to parquet-rs's defaults.
The same rows plan the same encodings in any chunking
(the_plan_depends_on_the_rows_not_the_chunks), so the output stays
deterministic. MUMDIA_PARQUET_PLAN=0 restores the unplanned layout, for a
byte comparison against a binary from before the plan. A pooled table spliced
from band artifacts carries each band's own plan in its row groups, which is
legal parquet. This is F2 of the 2026-09-25 performance survey.
A writer can also name columns to write without any dictionary
(WriteOptions::plain_column, TableWriter::with_plain_column), whatever the
plan says. The distinct-fraction test cannot see the one case where PLAIN wins
on a low-cardinality column: a list whose rows repeat whole runs of values,
which snappy shortens in PLAIN form and cannot find in bit-packed dictionary
indices. extract names the chromatogram rt axis, which each fragment row of
a candidate repeats: the AIF chromatograms came out 12.0% smaller than with
the planned dictionary (160.9 against 182.8 MB; 17.5% against the unplanned
layout), written in 5.4 against 7.2 s. This is X6 of the survey. An unknown
column name is an error when the writer opens.
parquet-rs's ArrowWriter encodes a row group's columns one after another on
the calling thread. Every writer here (TableWriter, BatchWriter,
write_batches) encodes through codec::ColumnEncoder instead, which builds
the same parts (ArrowWriter::try_new(..).into_serialized_writer(), then one
ArrowColumnWriter per leaf from the ArrowRowGroupWriterFactory) and encodes
the root columns of each batch concurrently. It splits a batch at the row-group
cap where ArrowWriter::write splits it, skips empty batches, gives each column
writer one write per (split) batch so the page checks fall where they fell,
closes the column chunks concurrently and appends them in schema order, and
inherits the ARROW:schema metadata from try_new. The file is the serial
writer's file byte for byte: the_parallel_encoder_writes_the_serial_writers_bytes
and a_wide_table_on_a_pool_is_the_serial_writers_file compare them over several
row groups, list columns, nulls, empty and straddling batches, pools of 1, 4
and 6 threads and no pool. A byte row-group cap or content-defined chunking is
refused rather than reproduced; no writer here sets either.
The work runs on a dedicated rayon pool, never the global one. Its size is
MUMDIA_PARQUET_THREADS when set (0 or 1 is serial), else --threads capped at
8 (set_codec_threads, called by the CLI; --threads 1 is serial), else the
machine's parallelism capped at 8. --threads N therefore bounds two pools
separately, the global pool at N and the codec pool at min(N, 8), and the codec
pool works for plain writer and loader threads while the global pool is busy
(extract's chromatogram writer encodes while the candidate loop scores, the
features writer while the chunk workers compute). Up to N + min(N, 8) threads
can then be busy at once. MUMDIA_PARQUET_THREADS=1 makes the codec serial,
which brings a run back to N threads plus its few plain writer and loader
threads, as before the codec pool. A writer called from inside any rayon pool
encodes on its own thread instead: rayon lets a worker that waits on another
pool steal jobs of its own pool meanwhile, and such a job could take a lock the
writer's caller holds. The parallel path is therefore used from plain threads
(the main thread, extract's chromatogram writer thread, the features writer
thread, rescore's handoff and psms_scored writers); a band or a run that already
runs on the global pool is parallel at that level.
The pool is shared by every plain-thread writer and decoder of the process, and
it counts the callers waiting on it (codec::claim). When as many callers wait
as the pool has threads, the next caller encodes or decodes on its own thread
instead of queueing behind them, which writes the same bytes. Under
groups.parallel every band in flight has its own chromatogram writer thread,
and features and compete writers can run beside them. Without the count, more
than eight such writers would share eight codec threads where each used to have
a core of its own, and a band's candidate loop would wait on its writer's
channel. With it, encode capacity never falls below one thread per concurrent
writer (a_saturated_pool_sends_the_next_caller_to_its_own_thread,
concurrent_writers_beyond_the_pool_size_write_the_serial_writers_file). The
timings below are for one writer at a time; a banded run with
groups.parallel >= 8 has not been timed against the serial codec.
Measured on the AIF artifacts re-chunked to 65,536-row batches under the shipped
properties (bench_parallel_encode_a_real_artifact, median of 3 rounds, every
arm byte-identical): features 0.45 s serial, 0.26 s on 2 threads, 0.19 s on 4,
0.14 s on 8; psms_competed 0.42 / 0.27 / 0.19 / 0.16 s; chromatograms 8.07 /
4.57 / 4.40 / 4.41 s (two list columns carry nearly all of it, so two threads
take the whole gain); spectra_ms2 0.42 / 0.25 / 0.26 / 0.27 s. This is R2 of
the 2026-09-25 performance survey.
Table (table.rs:200-204) holds the Arc<Schema>, the Vec<RecordBatch>,
and nrows. Table::read(path) (table.rs:207) opens the file, builds a
ParquetRecordBatchReaderBuilder, captures the schema, then iterates the reader
collecting every batch and summing num_rows(). It errors with context
"opening {path}" if the file cannot be opened and "reading parquet {path}" if
the builder cannot parse the Parquet footer. The whole file is materialized
into memory; there is no streaming or predicate pushdown.
Column access is by name. idx(name) (table.rs:381) resolves a name to a
column index via schema.index_of, returning a descriptive error listing all
column names if the name is absent. The typed getters each downcast every
batch's column to the concrete Arrow array type and concatenate across batches
into one Vec. They error if the downcast fails, so the type is checked at read
time. The message wording is per getter: "column '<name>' is not f64|f32|i64|i32|u32|bool" for the scalar getters (and opt_f64 reuses the f64
message), "column '<name>' is not utf8" for str (note: utf8, not str),
and for list_f32 either "column '<name>' is not a list" when the column is
neither a List nor a LargeList (table.rs:572) or "list '<name>' inner is not f32" when the inner element array is not f32 (table.rs:549). The str
getter downcasts to StringArray only (table.rs:503-511), that is Arrow
Utf8 with 32-bit offsets. An Arrow LargeUtf8 (large_utf8) column fails
with the same "is not utf8" message even though it holds strings; unlike
list_f32, the string path has no large-offset fallback.
The getters and their exact null behaviour:
f64(table.rs:387),f32(table.rs:407): fast path whennull_count() == 0usesextend_from_slice(a.values()); otherwise iterate and map a null tof64::NAN/f32::NAN. Nulls become NaN.f64_widening(onTableandTableFile): an f64 column exactly asf64reads it, or an f32 column widened byf64::from, which is exact; a null is NaN either way, and any other type is thef64error. It exists for columns whose stored width differs between artifact versions: the feature columns offeatures.parquetv2 andpsms_competed.parquetv4 are Float32 where v1 and v3 stored Float64 (docs/15_data_dictionary.md).i64(table.rs:427),i32(table.rs:447),u32(table.rs:467): fast path on no nulls; otherwise iterate pushinga.value(k)without checkingis_null. A null therefore comes through as the underlying buffer value (typically 0), not as a sentinel. See the gotcha below.bool(table.rs:487): always iteratesa.value(k); never checks null.str(table.rs:503): iterates; a null maps to an emptyString.opt_f64(table.rs:523): the only null-preserving getter. ReturnsVec<Option<f64>>, mapping a null toNone.list_f32(table.rs:542): reads an f32 list column and accepts bothList(32-bit offsets) andLargeList(64-bit offsets) encodings, so a chromatogram artifact written byCol::ListF32orCol::LargeListF32reads back through the same call. An outer null list becomes an emptyVec; a present list is materialized viaf.values().to_vec().
column_names() (table.rs:373) returns the schema field names in order.
parquet-rs's synchronous reader fetches every page on its own, a header read and a body read through a fresh seek each time, and the arrow reader advances every projected column in lockstep, one batch at a time. A batch of a 398-column features or competed table therefore issues about 400 page reads, each one column chunk away from the previous one, and a row group is covered in several strided passes. On an SSD or from the page cache that costs nothing measurable. On a spinning array it is seek-bound: one reader of the immunopeptidomics competed tables measured 44 MB/s against a 133 MB/s sequential ceiling, and seven concurrent readers only 47-53 MB/s together.
TableFile::scan(columns, batch_size, &ScanOptions) is batches with read
options. ScanOptions::default() is exactly batches. With
ScanOptions::coalesced() (or coalesce: Some(SpanReadOptions { .. })) the
reader is given a span_cache::SpanCache instead of the File:
- for each selected row group, in reading order, the projected column chunks'
byte ranges (
ColumnChunkMetaData::byte_range) are coalesced into spans. Two chunks at mostmax_gap_bytes(1 MB) apart share a span, so a narrow projection never reads the columns it skipped wholesale. ARowSpanhandle (open_rows,span) plans only its own row groups; - a span is read with one sequential read on first use and every page request
inside it is sliced from memory. Requests outside a span (the footer, an
unplanned column) and spans over
max_span_bytes(512 MB) are read directly, exactly as the plain reader reads them; - a span is released as soon as every column chunk in it has been read to its
end, so a forward scan holds the one or two row groups a batch straddles, plus
one prefetched span when
prefetchis on: a helper thread with its own file handle reads span i+1 while the decoder works through span i, withinmax_resident_bytes(1 GB). The decoder never waits on the budget; - the decoder receives the same bytes through the same metadata, projection,
selection and batch size, so the batches are identical to the plain reader's
(
coalesced_reads_yield_the_plain_readers_batchescompares them over projections, span handles, straddling batch sizes and every option corner;each_span_is_read_once_and_released_when_its_chunks_are_donecounts the reads). Adebugline on drop reports spans read, bytes, prefetches and direct reads.
Two full scans can use it, through stages::wide_scan_options: rescore's feature
stream (for_each_feature_batch, all ~390 feature columns of every row of every
competed input) and compete's pass-through copy (copy_kept_rows, every column of
the features table). Those are the wide scans of the run on spinning storage, where
the memory the cache holds (two or three row groups of the projection, about 0.9 GB
on a 131,072-row competed group) buys a forward read. It does not help a reader
that skips pages through the page index, because a span holds whole column chunks,
so every other reader keeps the plain File.
MUMDIA_WIDE_SCAN chooses the reader of those two scans (stages::WideScan):
| value | reader |
|---|---|
unset, plain (default) |
the plain reader with its parallel decode, at the scans' own batch sizes |
rowgroup |
the plain reader, with rescore's feature stream decoding one row group a batch |
coalesced |
the span cache above, one sequential read per row group |
The plain reader stays the default because it is the only one measured not to lose.
A coalesced scan decodes with one reader (see "Parallel decode" below), so from
the page cache or an SSD it gives up the automatic column groups. Measured from
the page cache on the HYE competed table (879,018 rows, 387 features, 131,072-row
groups, bench_feature_stream_a_real_artifact, minimum of 4 rounds on a loaded
host): the feature stream took 1.54 s plain, 1.58 s plain at one row group a batch,
and 2.50 s coalesced. One reader decodes about 1.1 GB/s of f64, far above the
44-133 MB/s of the spinning array the change is for, so the single reader would not
limit a seek-bound read, but neither rowgroup nor coalesced has been measured on
that array yet (the iostat check of the 2026-09-25 survey). Promote one once it has.
The output is the same under all three, because the batches are.
The arrow reader decodes every projected column of a batch on one thread. With
ScanOptions::decode_threads a scan instead builds one reader per contiguous
group of projected root columns (ReadSpec::column_groups, balanced by the
compressed bytes of the selected row groups), advances the groups together on
the codec pool (codec.rs) and joins each batch column-wise. Every group
reader has the same row groups, row selection and batch size, so each yields
the same rows per batch; the joined batches are the single reader's batches
exactly, schema included (parallel_decode_yields_the_single_readers_batches
covers projections, row spans trimmed at both ends, straddling batch sizes, more
groups than columns and coalesced reads). A group that ends early or returns a
different row count is an error, not a short batch.
Because the batches are identical, this is on by default: decode_threads: None
uses codec_threads() groups (at most 8, fewer when the projection has fewer
root columns or less than 4 MB of compressed data per group,
MIN_DECODE_GROUP_BYTES; automatic_decode_groups). It uses the single reader
from inside a rayon pool, for the reason the encoder avoids the pool there, and
for a coalesced scan. A coalesced scan exists to give seek-bound storage one
forward read per row group. Column groups would each hold their own spans, read
their own slices of every row group and run their own prefetcher, so the disk
would again see up to eight interleaved streams. Some(1) is the single reader;
Some(k) asks for k groups whatever the size, coalesced or not, and
MUMDIA_PARQUET_DECODE_THREADS=k does the same for every automatic scan of a
process, which is how a whole run, small artifacts included, is checked end to
end (the smoke with MUMDIA_PARQUET_DECODE_THREADS=3 writes every artifact
byte for byte as without it). Under coalesce with explicit groups each group
holds its own spans, and the resident budget is divided between them.
BatchReader::decode_groups reports what a scan uses. A getter that reads one
column is unchanged.
Automatic groups stay on for uncoalesced scans, including on spinning storage.
The plain reader already fetches each page on its own with a seek, from every
projected column chunk in turn (previous section). Column groups issue the same
page reads, of the same sizes, from several threads at once, so they change the
queue depth and not the number of seeks. On the spinning array where the plain
reader was measured, seven concurrent strided readers reached 47-53 MB/s
together against 44 MB/s for one, so concurrency did not lower throughput there.
The groups themselves have not been timed on that array. If a scan on rotational
storage is slower with them, MUMDIA_PARQUET_DECODE_THREADS=1 restores the
single reader for the whole process, and an iostat A/B on the server is the
check.
Measured on the AIF artifacts (bench_scan_a_real_artifact, full scans in
16,384-row batches, median of 5 interleaved rounds on a loaded 32-thread host):
features 1.18 s with one reader, 0.80 s with 2 groups, 0.65 s with 4, 0.60 s
with 8, 0.51 s automatic; psms_competed 1.22 / 0.86 / 0.55 / 0.47, 0.45 s
automatic; chromatograms 5.14 s against 2.99 s with 2 groups and 3.29 s
automatic (7 groups: the two list columns carry the work); spectra_ms2 0.85
against 0.44 s. Coalesced reads split into the same automatic group counts,
which a coalesced scan now uses only when decode_threads asks for them: 0.43,
0.47, 2.97 and 0.46 s.
TableFile::open parses the footer only. TableFile::open_with_offset_index
also loads the offset index, the location and first row of every data page of
every column chunk, where the file carries one (the engine's writer writes it;
pyarrow does not by default, and such a file opens as with open). A reader
built on that footer and given a row selection skips a page that holds no
selected row without reading or decompressing it. TableFile::batches_runs
takes a selection as (rows, keep) runs over the handle's own rows, on a whole
file or a span, and reuses the handle's footer, so a caller that selects row
group by row group parses the footer once. It always materialises the selection
as a queue of selectors (RowSelectionPolicy::Selectors): parquet-rs's
automatic choice may pick a bitmask, which decodes every row of a chunk and
filters afterwards, and so reads every page.
Skipping rows inside a page is not free. parquet-rs steps over the rows of a
list column one repetition level at a time, which on long lists costs more than
decoding them: on the AIF chromatograms rewritten in the current layout (90
points per trace on average, 65,536-row groups) selecting 11% of one group's
rt rows took 101 ms against 78 ms for the whole group, while on a run of the
six-file Astral LFQ experiment (21 points per trace) selecting 7% took 4.2 ms
against 12.3 ms. A caller that wants a
selection to cost no more than a full read therefore cuts it at page boundaries
(TableFile::page_starts gives the rows at which a column's pages begin) and
reads each column on its own selection, which is what quant does
(docs/12_quant_lfq_align_mbr_report_audit.md). What a skipped page holds is
never checked, so a corrupt page there goes unnoticed where a full read would
have refused it.
Several tables the engine reads are produced by a Python helper rather than by
write_table: an imported spectral library, an augmented library, a decoy
build. The reader is deliberately narrow, so those files must match it exactly.
Matching the column names and types is not sufficient.
- Codec. The workspace pins
parquetwithdefault-features = false, features = ["arrow", "snap"], so the reader decodes SNAPPY and uncompressed files only. Any other codec fails at read with a message of the form"Parquet error: Disabled feature at compile time: zstd". Polars writes zstd by default, so an externally produced table must set snappy explicitly. - String encoding. Write string columns as arrow
utf8. Polars defaults tolarge_utf8, which thestrgetter rejects with"column 'peptidoform' is not utf8"(see the read side above). - Library ids and row order. A library precursor table has two further
preconditions, both checked when the fragment index loads the library in the
mumdiacrate and both hard errors:candidate_idmust be the contiguous, row-aligned range0..ncand(index.rs:112-125), and rows must be ascending byprecursor_mz(index.rs:215-231, the precondition for thepartition_pointcandidate-range search; an unsorted import silently returns wrong candidate windows if the check is removed). The fragment table is grouped by a counting sort overcandidate_id(index.rs:132-153), so its rows need valid ids but no particular order, and each candidate's fragments keep their stored order.
mumdia inspect prints the schema but not the codec, so it does not by itself
confirm a file is readable. Read the row-group metadata (for example with
pyarrow.parquet.ParquetFile) to check compression, or rewrite the file with
snappy before feeding it to the engine.
blake3_file(path) (hash.rs:14) streams the file in 64 KiB chunks
([0u8; 1 << 16]) through a blake3::Hasher and returns the hex digest. This
is the artifact content_hash. It is fallible: it returns Result and errors
with context "hashing {path}" if the file cannot be opened or a read fails.
blake3_str(s) (hash.rs:29) is a one-shot hex digest of a string, used for the
config_hash; it is infallible and returns a plain String rather than a
Result. The engine derives the
config hash from Config::canonical_json() (config.rs:1125, a plain
serde_json::to_string), for example at main.rs:417. The convert command is
a deliberate exception: because --max-spectra, --top-peaks-ms2, and
--top-peaks-ms1 are CLI caps that change the spectra output but are not part
of Config, they are folded into the hash with a unit-separator (\u{1f})
alongside the canonical config JSON so two different caps do not collapse to the
same config_hash (main.rs:402-404).
Every stage used to publish its output and then read the whole file back
through blake3_file for the content_hash in its report, so each artifact byte
crossed the disk or the page cache twice. The parquet writers in table.rs only
ever append to their sink and never seek, so the digest of the byte stream is the
digest of the finished file. A writer can therefore hash while it writes:
hash::HashingWrite<W>forwards every byte toWand feeds the bytesWaccepted to ablake3::Hasher;- in
table.rsthe output file is aSink:File->HashingWrite-> a 1 MBBufWriter, so blake3 sees large contiguous inputs rather than parquet's page headers one at a time. An unhashed sink has a zero-capacity buffer, which passes every write straight through, so it issues the same writes as before; - hashing is opt-in per writer:
WriteOptions::content_hash()forBatchWriter::with_options,TableWriter::with_content_hash(),SpliceWriter::create_hashed,write_table_hashed,write_table_chunked_hashedandwrite_batches_hashed. The digest is finalised after parquet has written the footer (ColumnEncoder::into_inner, which ends inSerializedFileWriter::into_inner, forTableWriter,BatchWriterandwrite_batches;SerializedFileWriter::into_innerdirectly forSpliceWriter; both write the same footer bytes asclose) and returned as areport::Written { rows, content_hash }by the writer'sclose_hashed; - it is off for files nobody records: the pool's per-row-group rewrite temp files, compete's splice scratch files and the rescore sidecar handoff.
The stages that write through mumdia-io and then hashed the same file now use
the streamed digest: convert (both spectra tables, the isolation windows and the
MS2-to-MS1 map), predict-frag (both library tables), search-seed, rt-im-train,
extract (psms_extracted and chromatograms), features, compete (the spliced or
rewritten table; a hard-linked one carries the features hash), the pool's three
spliced tables in run_groups, and rescore's psms_scored. The small tables
(digest, peptidoforms, quant, prescan, align) still hash by reading back.
The digest is the same blake3 over the same bytes, so hashing while writing
changes no byte and no content_hash by itself: with only this change, the
smoke run's manifests and reports were identical to the previous binary's apart
from the git SHA and argv (hash_on_write_tests in table.rs pins digest ==
blake3_file for every writer type: empty, one block, many row groups, list
columns). The capped-writer layout changes in this release do change bytes and
hashes: the page cut ("Page layout of capped writers"), the float plan ("Float
encodings planned from the first rows") and the PLAIN chromatogram rt axis.
Each recorded hash is still the blake3 of the file as written. blake3_file itself is
still not memoised: features uses it in its tests as an independent integrity
check of a published file.
Both JSON sidecars (<artifact>.report.json and manifest.json) and every JSON
scalar the engine persists go through this module.
write_json<T: Serialize>(path, value) (json.rs:7) best-effort-creates the
parent directory (create_dir_all(...).ok(), json.rs:8-9, same best-effort
policy as write_table), serializes with serde_json::to_string_pretty
(pretty-printed, human-diffable), and writes the file, erroring with context
"writing json {path}" on an I/O failure.
read_json<T: DeserializeOwned>(path) (json.rs:16) is the counterpart used to
load configs and JSON sidecars back; it errors with context "reading json {path}" if the file cannot be read and "parsing json {path}" if
deserialization fails. Serialized key order follows the type's serde field order,
which is why ArtifactReport.stats uses a BTreeMap (below) to keep key order
deterministic.
ArtifactReport (report.rs:11-24) is the per-artifact JSON summary; it derives
Clone, Debug, Serialize, Deserialize (report.rs:10). A stage
constructs it and calls write_for(artifact_path) (report.rs:28), which
appends .report.json to the artifact path and writes it via
json::write_json (pretty-printed). Fields:
| field | type | meaning |
|---|---|---|
logical_name |
String | logical artifact name (usually the schema name) |
schema_name |
String | schema name from the schema.rs tuple |
schema_version |
u32 | schema version from the tuple |
stage |
String | producing stage, e.g. "digest" |
rows |
u64 | row count returned by write_table |
content_hash |
String | blake3_file of the artifact just written |
params |
serde_json::Value |
the resolved parameters the stage actually used |
stats |
BTreeMap<String, Value> |
summary key distributions / metrics (ordered) |
model_identity |
Option<String> |
sidecar / predictor identity, else None |
elapsed_ms |
u128 | wall-clock ms for the stage |
stats is a BTreeMap (not a HashMap) so the JSON key order is
deterministic. Concrete example: digest records params with enzyme, missed
cleavages, length bounds, decoy strategy, rng seed, and max_decoy_attempts,
and stats with n_targets, n_decoys, decoy_collision_retries, and
dropped_target_decoy_pairs (stages/digest.rs:301-331; the collision-retry and
dropped-pair counters were added with the collision-safe decoy resolver in
e7d7fa5). convert has a
stage-local helper write_reports (stages/convert.rs:269-290) that writes one
report per artifact with a shared params and empty stats; note this
write_reports lives in convert.rs, it is not part of the mumdia-io public
API. Stages with report coverage build their ArtifactReport directly (see the
write_for call sites in align, compete, extract, search_seed,
predict_frag, peptidoforms, rt_im_train, rescore, features, quant).
record_artifact(logical_name, schema, path, rows, stage, config_hash)
(lib.rs:20) builds an ArtifactRecord (manifest.rs:10-20, in
mumdia_core::manifest; derives Clone/Debug/Serialize/Deserialize) for
the run manifest. The record has nine fields: logical_name, path, and rows
are copied straight from the arguments; format is hard-coded to "parquet"
(lib.rs:31); schema_name/schema_version come from the schema tuple;
content_hash is blake3_file(path) of the file just written; producing_stage
comes from the stage argument (the struct field and the argument are named
differently); and config_hash is the argument. The
run orchestrator calls it for selected primary Parquet artifacts and records
those into the Manifest (stages/run.rs, many call sites around
record_artifact(...)), which is then serialized to manifest.json.
Calibration JSON, PIN/schema companions, TSVs, and some diagnostics are not
manifest records. Standalone single-stage invocations write their normal
sidecars, when implemented, but not a manifest.
record_artifact reads and hashes the whole file. Every stage has already done
that once for its own <artifact>.report.json, so calling it after the stage
read each artifact twice; on a grouped run the band artifacts are tens of GB.
The orchestrators therefore take the hash from the stage instead:
- each stage an orchestrator calls has a
run_hashedsibling ofrun(digest,peptidoforms,predict_frag,search_seed,rt_im_train,extract,featuresandfeatures::run_with_bounds_hashed,compete,rescore,quant) that returns areport::Written { rows, content_hash }per output, taken from the report it just wrote.convert::runreturns the four hashes inConvertOutputs::hashes.runitself is unchanged and returns the row counts, so the CLI and the tests call it as before; - the caller records with
Written::record, a thin wrapper ofrecord_artifact_with_hash, so the manifest carries the identical hash; - in library-input mode the two library records reuse the hash
runcomputed when it recorded the same files as inputs; - the grouped path (
run_groups) builds no band record at all when it has no manifest to put it in.run-experimentpasses none, and before this change every band artifact was hashed a second time there for a record that was then dropped; run-experimentreuses the rescore and quant hashes for its experiment manifest, and the by-source split hashes its per-run tables while it writes them. The MBR worker's table, the LFQ matrix, the pooled seed and the DeepLC library tables have no Rust report hash, so they are still hashed byrecord_artifact.
Reusing the stage hashes changes no hash value in manifest.json,
experiment_manifest.json or the *.report.json files; only the second read
is gone. hash::blake3_file itself is not memoised: features uses it as an
independent integrity check. The stage's own hash no longer reads the file back
either for the large artifacts ("Hash on write" above).
One hash does change in the same release, for a different reason: compete
now publishes psms_competed.parquet as the features file's own bytes when
it removes no row (docs/11 "compete: how the competed table is published").
Its content_hash then equals the features hash, and it differs from the
hash of the 131,072-row-group rewrite an earlier binary wrote whenever the
table holds more than one row group. On a grouped run the pooled
psms_competed hash differs for the same reason (docs/33 section 5).
psms_scored.parquet and every artifact after it are byte-identical.
inspect(path) reads the whole table via Table::read, then builds a string:
artifact: <path>, rows: <n>, a schema: block listing each field as
<name>: <DataType> with a (nullable) suffix when the field is nullable,
and a head: block. The head is the first up-to-10 rows sliced from the
first batch only (first.slice(0, min(num_rows, 10)), lib.rs:57-59),
formatted with arrow::util::pretty::pretty_format_batches. The result is
appended only inside an if let Ok(p) = ... (lib.rs:60), so if pretty-print
fails the head is silently omitted; there is no else branch. The CLI
command Cmd::Inspect { artifact } (main.rs:712) simply prints the returned
string (main.rs:713).
init_logging() (lib.rs:13) initializes tracing_subscriber once, honouring
RUST_LOG and defaulting to info, with with_target(false). It uses
try_init so a second call is a no-op rather than a panic.
| name | file:line | what it does |
|---|---|---|
Col (enum) |
table.rs:23 |
typed write-side column: scalar, Opt*, and list variants |
Col::field |
table.rs:83 |
Arrow Field; scalars non-nullable, Opt*/lists nullable |
Col::into_array |
table.rs:107 |
consuming move of the Vec into an ArrayRef (copy once) |
write_table |
table.rs:151 |
validate + write one SNAPPY Parquet batch; returns row count |
write_table_chunked |
table.rs |
the write_table file, built from caller-produced 65,536-row chunks |
write_table_chunked_hashed |
table.rs |
write_table_chunked, returning the rows and the content hash taken while writing |
Table (struct) |
table.rs:200 |
read-back table: schema, batches, nrows |
Table::read |
table.rs:207 |
read a Parquet file fully into memory |
Table::column_names |
table.rs:227 |
schema field names, in order |
Table::f64 / f32 |
table.rs:241 / 261 |
float getters; null -> NaN |
Table::f64_widening / TableFile::f64_widening |
push_f64_widening |
an f64 or f32 column as f64 (f32 widened exactly); null -> NaN |
Table::i64/i32/u32 |
table.rs:281/301/321 |
integer getters; null NOT checked (-> buffer value) |
Table::bool |
table.rs:341 |
bool getter; null NOT checked |
Table::str |
table.rs:357 |
string getter; null -> "" |
Table::opt_f64 |
table.rs:377 |
only null-preserving getter; -> Vec<Option<f64>> |
Table::list_f32 |
table.rs:396 |
f32 list getter; reads List and LargeList; null row -> empty Vec |
TableFile / TableFile::open / open_rows |
table.rs:1673 / :1742 / :1765 |
footer-only handle whose getters stream one column batch by batch; open_rows is a row span of the file that behaves as a smaller file |
TableFile::row_parts |
table.rs:1868 |
cut a handle into at most max_parts row-contiguous parts, in order: whole row groups, merged, and page-aligned ranges inside a group that has an offset index (a group without one is never split) |
TableFile::batches / batches_dict |
table.rs:2038 / :2050 |
streaming batch reader over the named columns; batches_dict reads the named Utf8 columns as Dictionary(Int32, Utf8), with the same row values |
TableFile::batches_selected |
table.rs:2100 |
stream only the rows of (rows, keep) runs, skipping pages with no kept row where the file has an offset index (whole-file handles only) |
TableFile::open_with_offset_index / offset_indexed |
table.rs |
open with the page locations of every column chunk loaded as well (opt-in); whether a handle holds them for every row group it covers |
TableFile::batches_runs |
table.rs |
stream only the rows of (rows, keep) runs of this handle, a whole file or a span, through the footer it already holds, always as a queue of selectors; a page with no kept row is skipped unread on an offset-indexed handle |
TableFile::row_group_parts / page_starts |
table.rs |
a handle cut at the file's row-group boundaries; the handle rows at which a column's data pages begin in every one of its leaves, from the offset index |
WriteOptions::metadata / TableWriter::with_metadata / TableFile::metadata_value |
table.rs |
record a small fact about the whole table under a key in the footer's key-value metadata, next to the arrow schema; read it back from the footer alone (the overlap loser table names its band tables this way) |
TableFile::str_flat_rows |
table.rs |
str_flat at given rows only: one text arena plus offsets for the picked rows, the same null policy |
TableFile::str_interned / str_flat |
table.rs:2309 / :2358 |
a string column as one id per row plus its distinct values (first appearance), or as one text arena plus offsets; both refuse a NULL with the row |
StrBatch / StrInterner |
table.rs:1517 / :1550 |
one batch of a string column, plain or through its dictionary; first-appearance interning with a per-batch key memo |
ListF32 |
table.rs:1422 |
borrowed view of a batch's f32 list column (row slices of the batch's own buffer) |
require_no_nulls |
table.rs:1203 |
refuse a NULL in a required column of a batch a reader walks itself, naming the absolute row |
ArtifactReport |
report.rs:11 |
per-artifact JSON summary struct |
ArtifactReport::write_for |
report.rs:28 |
write <artifact>.report.json |
blake3_file |
hash.rs:8 |
streamed blake3 hex digest of a file (content_hash) |
blake3_str |
hash.rs:23 |
one-shot blake3 hex digest of a string (config_hash) |
write_json / read_json |
json.rs:7 / 16 |
pretty serde JSON write / read, creating parent dirs on write |
record_artifact |
lib.rs:20 |
build an ArtifactRecord (format hard-coded "parquet") |
record_artifact_with_hash |
lib.rs |
the same record from a hash the caller already has |
report::Written |
report.rs |
an output's row count and report content hash, returned by each stage's run_hashed; Written::record builds the manifest record |
inspect |
lib.rs:43 |
schema + head(<=10, first batch) + row count as a string |
init_logging |
lib.rs:13 |
tracing init honouring RUST_LOG, default info |
This subsystem reads no Config fields. It has no config surface of its own, so
the recent pruning of dead config fields in mumdia-core::config did not touch
it. Its behaviour is fixed at compile time:
- Compression is
SNAPPY, hard-coded inwrite_table(table.rs:204-206); it is not configurable and there is no other codec path. The read side is equally narrow, so this is a contract on externally written inputs too, not only on this crate's outputs (see "Parquet written outside this crate"). - The hash read buffer is 64 KiB (
hash.rs:11). - Schema identifiers are the constants in
mumdia-core/src/schema.rs:7-24; a new artifact requires a new tuple there, not a config change. - Dependency features are pinned in the workspace
Cargo.toml:arrowv59 with["prettyprint"](needed byinspect),parquetv59 withdefault-features = false, features = ["arrow", "snap"](pure-Rust snap, no cmake/C toolchain),blake3v1. Do not re-enable Parquet default features; they pull in a C compression backend that breaks the pure-Rust build constraint (see CLAUDE.md environment gotchas). The same pin also bounds what the engine can read: a Parquet file compressed with any other codec is a read error, whoever wrote it.
- Determinism.
write_tablewrites columns in the order given and a single batch, so byte layout is stable for stable input.ArtifactReport.statsis aBTreeMap, so JSON key order is fixed.config_hashcomes fromConfig::canonical_json(serde field order) so it is reproducible. blake3 is deterministic. All of this is required by the determinism contract (seedocs/18_findings_and_decisions.md, B2). - Integer/bool null gotcha.
i64/i32/u32/boolgetters do not checkis_null; in the presence of nulls they return the raw buffer value (usually 0 / false), silently. This is safe today only because the MVP integer/bool columns are written non-nullable (theI*/Boolvariants setnullable = false). If you write a nullable integer column viaOptI32, reading it back withi32()will erase the null distinction. Onlyopt_f64is null-preserving on the read side. This is a known gap (CLAUDE.md "Correctness": add null-aware getters). - Read-side coverage asymmetry. The write side has
OptF32,OptI32,OptStr,ListF64variants with no matching null-aware / typed reader.OptF32reads back viaf32()(null -> NaN, acceptable),OptStrviastr()(null ->""),OptI32viai32()(null erased), and there is nolist_f64reader at all. Add the corresponding getter before relying on a round-trip of those variants. - Inner-list nulls are not represented. List item fields are declared
nullable, but the builders only ever
append(true)present lists and never append inner-element nulls;list_f32reads inner values withvalues().to_vec()and cannot surface an inner null. Only the outer list-level null (emptyVec) is modeled. - Duplicate column names are rejected at write time (
table.rs:157-165) because Arrow readers resolve to the first match and would hide the rest. - Length equality is enforced: all columns must have identical length or
write_tableerrors (table.rs:166-176). - Everything is loaded into memory.
Table::readcollects all batches andinspectreads the full table just to print 10 rows; there is no streaming path. For very large artifacts this is a real memory cost.inspect's head is taken from the first batch only, so a file with tiny leading batches shows few rows even when later batches are large. inspectswallows pretty-print errors (lib.rs:60, anif let Ok(p)with noelse): a formatting failure omits the head silently rather than erroring, so absence of a head block is not proof of an empty table.create_dir_allon write is best-effort (.ok(),table.rs:186); a real permission failure surfaces later atFile::create, not at the mkdir.record_artifacthard-codesformat = "parquet"(lib.rs:31); it is not suitable for a non-Parquet artifact without a change there.- Externally written Parquet must match the reader, not just the schema. A
table produced outside
mumdia-iomust be SNAPPY-compressed with arrowutf8string columns. zstd andlarge_utf8are both Polars defaults and both are hard read errors ("Parquet error: Disabled feature at compile time: zstd","column 'peptidoform' is not utf8"). A library table must additionally satisfy thecandidate_idcontiguity (index.rs:112-125) andprecursor_mzordering (index.rs:215-231) preconditions. See "Parquet written outside this crate". - LargeList transparency.
list_f32intentionally reads bothListandLargeList, so an artifact whose encoding differs between two builds (32- vs 64-bit offsets) still reads identically; do not assume a fixed offset width when consuming chromatograms. - Test coverage is one round-trip. The crate's only unit test is
roundtrip_mixed_columns(table.rs:438), which writes then reads backU32/F64/Str/OptF64/ListF32/LargeListF32, asserting among other things that aLargeListF32column cross-reads throughlist_f32.hash.rs,json.rs,report.rs, andlib.rs(inspect,record_artifact,init_logging) have no unit tests, and the integer/bool null gotcha and theOpt*/ListF64reader gaps above are consequently not exercised by the suite.
- Add a new column type. Add a variant to
Col(table.rs:23) and extend the four matches:name,len,field,into_array. Decide nullability infield. Then add the matching read-side getter onTablefollowing the downcast + concatenate pattern of the existing getters, and be explicit about null handling (prefer anopt_*return over a silent sentinel for anything that can be null in practice). - Add a null-aware integer/string getter. This is the standing correctness
item. Mirror
opt_f64(table.rs:377): iterateis_null(k)and returnVec<Option<T>>. Do not change the existing non-optional getters' signatures; add new ones so current call sites are unaffected. - Register a new artifact. Add a
(name, version)tuple tomumdia-core/src/schema.rsand pass it towrite_table(schema = theVec<Col>you write) plus theArtifactReport/record_artifactcalls. Bump the version only on a breaking column-schema change (as was done forpsms_scored, now at version 3, and forpsms_competed/peptide_quant/protein_group_quant, now at version 2) and document the change. There is no read-side version check today, so the bump is provenance only; if you need a hard gate, add the check on the read path (it does not exist yet). - Change compression. It is a one-line change in
write_table(table.rs:205); keep it to a codec available under the pinned pure-Rust Parquet features, and re-measure round-trip and file size. It is not a write-side-only change: every externally written input has to move to the same codec, since the reader decodes only what the pin enables. - Extend the report. Add a field to
ArtifactReport(report.rs:11). KeepstatsaBTreeMapfor deterministic key order. Since it isSerialize/Deserialize, adding a field is backward compatible only if it isOptionor has a serde default; otherwise old.report.jsonfiles fail to deserialize. - Reuse the inspector. Any tool that needs a schema/head view should call
mumdia_io::inspectrather than re-opening Parquet, so the output format stays consistent with themumdia inspectCLI command.