"""Parquet snapshot store for Yahoo day_gainers ticks. Layout: data/cache/yahoo_gainers/{YYYY-MM-DD}/{HH-MM-SS}.parquet Each file contains one tick's worth of GainerQuote rows. The SQLite tgtc_snapshots table mirrors this for fast UI polling. """ from __future__ import annotations import logging import os from pathlib import Path log = logging.getLogger(__name__) _DEFAULT_CACHE_DIR = "data/cache/yahoo_gainers" def save_snapshot_parquet( quotes: list, # list[GainerQuote] date_str: str, captured_at_iso: str, cache_dir: str = _DEFAULT_CACHE_DIR, ) -> None: """Write one snapshot tick to Parquet (atomic).""" try: import pyarrow as pa import pyarrow.parquet as pq except ImportError: log.warning("TGTC snapshot_store: pyarrow not available, skipping Parquet write") return if not quotes: return # captured_at_iso like "2026-05-05T09:35:00" safe_ts = captured_at_iso.replace(":", "-").replace("T", "_")[:16] # "2026-05-05_09-35" dir_path = Path(cache_dir) / date_str dir_path.mkdir(parents=True, exist_ok=True) file_path = dir_path / f"{safe_ts}.parquet" schema = pa.schema([ pa.field("captured_at", pa.string()), pa.field("symbol", pa.string()), pa.field("rank", pa.int32()), pa.field("price", pa.float64()), pa.field("pct_change", pa.float64()), pa.field("volume", pa.float64()), pa.field("market_cap", pa.float64()), ]) table = pa.table( { "captured_at": [captured_at_iso] * len(quotes), "symbol": [q.symbol for q in quotes], "rank": [q.rank for q in quotes], "price": [q.price for q in quotes], "pct_change": [q.pct_change for q in quotes], "volume": [q.volume for q in quotes], "market_cap": [q.market_cap for q in quotes], }, schema=schema, ) # Atomic write via temp file tmp_path = file_path.with_suffix(f".{os.urandom(4).hex()}.tmp") try: pq.write_table(table, str(tmp_path), compression="snappy") os.replace(str(tmp_path), str(file_path)) except Exception as exc: log.error("TGTC snapshot_store: write failed %s: %s", file_path, exc) if tmp_path.exists(): tmp_path.unlink(missing_ok=True) def load_snapshots_for_date(date_str: str, cache_dir: str = _DEFAULT_CACHE_DIR) -> list[dict]: """Load all snapshot ticks for a date from Parquet files, sorted by captured_at.""" try: import pyarrow.parquet # noqa: F401 except ImportError: return [] dir_path = Path(cache_dir) / date_str if not dir_path.exists(): return [] rows: list[dict] = [] for parquet_file in sorted(dir_path.glob("*.parquet")): try: import pyarrow.parquet as pq table = pq.read_table(str(parquet_file)) for batch in table.to_batches(): rows.extend(dict(zip(table.schema.names, row)) for row in zip(*[batch.column(i).to_pylist() for i in range(len(table.schema))])) except Exception as exc: log.warning("TGTC snapshot_store: failed to read %s: %s", parquet_file, exc) return rows