Flamegraph of precompute_features on 1Q ES showed 62% of CPU time in zstd decompression, 6% in DBN FSM parsing, and only 2% in the actual feature math — single-threaded zstd was the bottleneck, not compute. Two fixes: 1. Per-quarter parallelism on the volume-bar trades loop (was sequential `for file in &trade_files`); brings it in line with the OFI path that already used par_iter. 2. Predecoded sidecar cache in `crates/ml-features/src/predecoded.rs`: first call to a `.dbn.zst` writes a bincode'd Vec<Mbp10Snapshot> or Vec<DbnTrade> under `<output_dir>/predecoded/`. Subsequent calls deserialize the sidecar and skip zstd entirely. An mtime+size header self-invalidates the sidecar when the source changes — no manual flush needed when a quarter is re-downloaded. Local 1Q ES results: - cold (writes sidecar): 40.7s (was 39.3s; +1.4s for write) - warm (HIT): 4.7s (8.7× faster) - zstd in flat perf: 62% → 0% of CPU samples - sidecar disk per Q: ~150MB The sidecar layer also auto-dedupes within a single run: the OFI section re-loads trades, but the second call hits the sidecar that the volume-bar section wrote moments earlier. CLI: `--rebuild-predecoded` purges sidecars for cold-path testing or after a wire-format change to Mbp10Snapshot / DbnTrade. Sidecars also self-invalidate on format-version mismatch so old caches are skipped silently rather than mis-deserializing. Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
282 lines
9.3 KiB
Rust
282 lines
9.3 KiB
Rust
//! Predecoded sidecar cache for parsed DBN files.
|
|
//!
|
|
//! Wraps `load_trades_sync` and `parse_mbp10_streaming` with an on-disk cache.
|
|
//! First call decodes zstd-DBN and writes a sidecar; subsequent calls
|
|
//! deserialize the sidecar directly. zstd decompression was 62% of
|
|
//! precompute_features wall time on the 2026-05-16 1Q ES profile — moving it
|
|
//! to a one-shot cache eliminates that cost for all repeat runs.
|
|
//!
|
|
//! # Cache key
|
|
//!
|
|
//! The sidecar header records the source file's mtime and size. A sidecar is
|
|
//! considered valid only when both match the source on the current call;
|
|
//! re-downloading a quarter (which bumps mtime + usually size) silently
|
|
//! invalidates the sidecar.
|
|
|
|
use std::fs::{self, File};
|
|
use std::io::{BufReader, BufWriter, Read, Write};
|
|
use std::path::{Path, PathBuf};
|
|
use std::time::SystemTime;
|
|
|
|
use bincode;
|
|
use data::providers::databento::dbn_parser::DbnParser;
|
|
use data::providers::databento::mbp10::Mbp10Snapshot;
|
|
use serde::{Deserialize, Serialize};
|
|
|
|
use crate::trades_loader::{load_trades_sync, DbnTrade};
|
|
use crate::MLError;
|
|
|
|
const MAGIC: u32 = u32::from_le_bytes(*b"PCD1");
|
|
const VERSION: u32 = 1;
|
|
|
|
fn sidecar_path(source: &Path, predecoded_dir: &Path, kind: &str) -> PathBuf {
|
|
let basename = source.file_name().unwrap_or_default().to_string_lossy();
|
|
predecoded_dir.join(format!("{basename}.{kind}.predecoded.bin"))
|
|
}
|
|
|
|
fn source_mtime_size(source: &Path) -> std::io::Result<(u128, u64)> {
|
|
let meta = fs::metadata(source)?;
|
|
let mtime = meta
|
|
.modified()?
|
|
.duration_since(SystemTime::UNIX_EPOCH)
|
|
.map(|d| d.as_nanos())
|
|
.unwrap_or(0);
|
|
Ok((mtime, meta.len()))
|
|
}
|
|
|
|
fn write_header<W: Write>(w: &mut W, mtime_nanos: u128, size: u64) -> std::io::Result<()> {
|
|
w.write_all(&MAGIC.to_le_bytes())?;
|
|
w.write_all(&VERSION.to_le_bytes())?;
|
|
w.write_all(&mtime_nanos.to_le_bytes())?;
|
|
w.write_all(&size.to_le_bytes())?;
|
|
Ok(())
|
|
}
|
|
|
|
fn read_header_and_validate<R: Read>(r: &mut R, source: &Path) -> std::io::Result<bool> {
|
|
let mut magic = [0u8; 4];
|
|
r.read_exact(&mut magic)?;
|
|
if u32::from_le_bytes(magic) != MAGIC {
|
|
return Ok(false);
|
|
}
|
|
let mut version = [0u8; 4];
|
|
r.read_exact(&mut version)?;
|
|
if u32::from_le_bytes(version) != VERSION {
|
|
return Ok(false);
|
|
}
|
|
let mut mtime_bytes = [0u8; 16];
|
|
r.read_exact(&mut mtime_bytes)?;
|
|
let mut size_bytes = [0u8; 8];
|
|
r.read_exact(&mut size_bytes)?;
|
|
let stored_mtime = u128::from_le_bytes(mtime_bytes);
|
|
let stored_size = u64::from_le_bytes(size_bytes);
|
|
let (current_mtime, current_size) = source_mtime_size(source)?;
|
|
Ok(stored_mtime == current_mtime && stored_size == current_size)
|
|
}
|
|
|
|
fn try_read_sidecar<T: for<'de> Deserialize<'de>>(
|
|
source: &Path,
|
|
sidecar: &Path,
|
|
) -> std::io::Result<Option<T>> {
|
|
if !sidecar.exists() {
|
|
return Ok(None);
|
|
}
|
|
let f = File::open(sidecar)?;
|
|
let mut r = BufReader::new(f);
|
|
if !read_header_and_validate(&mut r, source)? {
|
|
return Ok(None);
|
|
}
|
|
let payload: T = bincode::deserialize_from(r)
|
|
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
|
|
Ok(Some(payload))
|
|
}
|
|
|
|
fn write_sidecar<T: Serialize>(
|
|
source: &Path,
|
|
sidecar: &Path,
|
|
payload: &T,
|
|
) -> std::io::Result<()> {
|
|
if let Some(parent) = sidecar.parent() {
|
|
fs::create_dir_all(parent)?;
|
|
}
|
|
let (mtime, size) = source_mtime_size(source)?;
|
|
let f = File::create(sidecar)?;
|
|
let mut w = BufWriter::new(f);
|
|
write_header(&mut w, mtime, size)?;
|
|
bincode::serialize_into(&mut w, payload)
|
|
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
|
|
w.flush()?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Load trades from a DBN file, using a predecoded sidecar when available.
|
|
///
|
|
/// On miss (no sidecar / stale sidecar): decode `.dbn.zst` via
|
|
/// `load_trades_sync`, persist `Vec<DbnTrade>` as a sidecar, return trades.
|
|
/// On hit: deserialize sidecar (skips zstd).
|
|
///
|
|
/// Sidecar write failures are logged and ignored — the caller still receives
|
|
/// the freshly-decoded `Vec<DbnTrade>`.
|
|
pub fn load_or_predecode_trades(
|
|
source: &Path,
|
|
predecoded_dir: &Path,
|
|
) -> Result<Vec<DbnTrade>, MLError> {
|
|
let sidecar = sidecar_path(source, predecoded_dir, "trades");
|
|
match try_read_sidecar::<Vec<DbnTrade>>(source, &sidecar) {
|
|
Ok(Some(trades)) => {
|
|
tracing::info!(
|
|
"Predecoded HIT: {} ({} trades)",
|
|
sidecar.display(),
|
|
trades.len()
|
|
);
|
|
return Ok(trades);
|
|
}
|
|
Ok(None) => {}
|
|
Err(e) => {
|
|
tracing::warn!(
|
|
"Predecoded trades sidecar read failed for {}: {} — falling back to zstd",
|
|
sidecar.display(),
|
|
e
|
|
);
|
|
}
|
|
}
|
|
let trades = load_trades_sync(source)?;
|
|
if let Err(e) = write_sidecar(source, &sidecar, &trades) {
|
|
tracing::warn!(
|
|
"Failed to write predecoded trades sidecar {}: {}",
|
|
sidecar.display(),
|
|
e
|
|
);
|
|
} else {
|
|
tracing::info!(
|
|
"Predecoded WRITE: {} ({} trades)",
|
|
sidecar.display(),
|
|
trades.len()
|
|
);
|
|
}
|
|
Ok(trades)
|
|
}
|
|
|
|
/// Load MBP-10 snapshots from a DBN file, using a predecoded sidecar when available.
|
|
///
|
|
/// On miss: streams the file via `DbnParser::parse_mbp10_streaming`, accumulates
|
|
/// `Vec<Mbp10Snapshot>`, persists a sidecar, returns the snapshots. On hit:
|
|
/// deserializes the sidecar (skips zstd).
|
|
pub fn load_or_predecode_mbp10(
|
|
source: &Path,
|
|
predecoded_dir: &Path,
|
|
) -> Result<Vec<Mbp10Snapshot>, MLError> {
|
|
let sidecar = sidecar_path(source, predecoded_dir, "mbp10");
|
|
match try_read_sidecar::<Vec<Mbp10Snapshot>>(source, &sidecar) {
|
|
Ok(Some(snapshots)) => {
|
|
tracing::info!(
|
|
"Predecoded HIT: {} ({} snapshots)",
|
|
sidecar.display(),
|
|
snapshots.len()
|
|
);
|
|
return Ok(snapshots);
|
|
}
|
|
Ok(None) => {}
|
|
Err(e) => {
|
|
tracing::warn!(
|
|
"Predecoded mbp10 sidecar read failed for {}: {} — falling back to zstd",
|
|
sidecar.display(),
|
|
e
|
|
);
|
|
}
|
|
}
|
|
let parser = DbnParser::new()
|
|
.map_err(|e| MLError::InsufficientData(format!("DbnParser init failed: {e}")))?;
|
|
let mut snapshots = Vec::new();
|
|
parser
|
|
.parse_mbp10_streaming(source, 100, |snap| {
|
|
snapshots.push(snap.clone());
|
|
})
|
|
.map_err(|e| MLError::InsufficientData(format!("parse_mbp10_streaming failed: {e}")))?;
|
|
if let Err(e) = write_sidecar(source, &sidecar, &snapshots) {
|
|
tracing::warn!(
|
|
"Failed to write predecoded mbp10 sidecar {}: {}",
|
|
sidecar.display(),
|
|
e
|
|
);
|
|
} else {
|
|
tracing::info!(
|
|
"Predecoded WRITE: {} ({} snapshots)",
|
|
sidecar.display(),
|
|
snapshots.len()
|
|
);
|
|
}
|
|
Ok(snapshots)
|
|
}
|
|
|
|
/// Remove all `*.predecoded.bin` files from `predecoded_dir`. Idempotent.
|
|
///
|
|
/// Used by the `--rebuild-predecoded` CLI flag to force a fresh decode of every
|
|
/// source file. Returns the number of files removed.
|
|
pub fn purge_predecoded_dir(predecoded_dir: &Path) -> std::io::Result<usize> {
|
|
if !predecoded_dir.exists() {
|
|
return Ok(0);
|
|
}
|
|
let mut removed = 0usize;
|
|
for entry in fs::read_dir(predecoded_dir)? {
|
|
let entry = entry?;
|
|
let path = entry.path();
|
|
if path
|
|
.file_name()
|
|
.and_then(|n| n.to_str())
|
|
.is_some_and(|n| n.ends_with(".predecoded.bin"))
|
|
{
|
|
fs::remove_file(&path)?;
|
|
removed += 1;
|
|
}
|
|
}
|
|
Ok(removed)
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use tempfile::TempDir;
|
|
|
|
#[test]
|
|
fn sidecar_path_uses_kind_suffix() {
|
|
let p = sidecar_path(
|
|
Path::new("/data/ES.FUT_2024-Q1.dbn.zst"),
|
|
Path::new("/tmp/predecoded"),
|
|
"trades",
|
|
);
|
|
assert_eq!(
|
|
p,
|
|
PathBuf::from("/tmp/predecoded/ES.FUT_2024-Q1.dbn.zst.trades.predecoded.bin")
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn header_round_trip_validates_source_metadata() {
|
|
let tmp = TempDir::new().unwrap();
|
|
let source = tmp.path().join("source.bin");
|
|
fs::write(&source, b"hello world").unwrap();
|
|
let sidecar = tmp.path().join("source.bin.test.predecoded.bin");
|
|
let payload: Vec<u32> = vec![1, 2, 3];
|
|
write_sidecar(&source, &sidecar, &payload).unwrap();
|
|
|
|
let read: Option<Vec<u32>> = try_read_sidecar(&source, &sidecar).unwrap();
|
|
assert_eq!(read, Some(vec![1, 2, 3]));
|
|
|
|
// Mutate source → sidecar must be invalidated.
|
|
fs::write(&source, b"different content here").unwrap();
|
|
let read: Option<Vec<u32>> = try_read_sidecar(&source, &sidecar).unwrap();
|
|
assert_eq!(read, None);
|
|
}
|
|
|
|
#[test]
|
|
fn purge_removes_only_predecoded_files() {
|
|
let tmp = TempDir::new().unwrap();
|
|
fs::write(tmp.path().join("a.trades.predecoded.bin"), b"x").unwrap();
|
|
fs::write(tmp.path().join("b.mbp10.predecoded.bin"), b"y").unwrap();
|
|
fs::write(tmp.path().join("unrelated.txt"), b"z").unwrap();
|
|
let removed = purge_predecoded_dir(tmp.path()).unwrap();
|
|
assert_eq!(removed, 2);
|
|
assert!(tmp.path().join("unrelated.txt").exists());
|
|
}
|
|
}
|