SPB Git forge

spb/hfmarketdata

Public

Open high-frequency market data platform — FirstRate full-history downloader, DuckDB/Parquet lake, open REST API and React docs platform (www.hfmarketdata.io)

127commits 1branches 0releases
24.7 MBsize
maindefault branch
11 days agolast push
JavaScript 53.7% Python 38.3% CSS 4.6% TypeScript 3.1%
10.2 KB · 215 lines python
Raw Blame History
1#!/usr/bin/env python32"""Rewrite the intraday Parquet files of the lake with small row groups (65 536 rows) — resumable, file by file.34Why: the files written by frd_downloader.py before 2026-09 have 1 000 000-row groups (3 groups for a 2.5 M-row51-minute file). A 5-minute window forces DuckDB to decode a whole third of the file; with 65 536-row groups6(≈ one week of 1-minute bars) the min/max statistics let it skip everything but one group.78Each file is rewritten with DuckDB (`COPY (SELECT * FROM read_parquet(src)) TO tmp (FORMAT PARQUET, COMPRESSION9ZSTD, ROW_GROUP_SIZE 65536)`), verified (row count identical) and swapped in atomically with `os.replace`.10Progress is kept in a JSON manifest so the job can be stopped and resumed (nightly PM2 task).1112Usage (production node M3U96b — do NOT run while frd_downloader.py rewrites the same directory):13    cd ~/hfmarketdata && venv/bin/python scripts/rewrite_row_groups.py --data-root ~/firstratedata --dry-run14    … --asset stock --timeframe 1min                 one directory family15    … --workers 3 --row-group-size 65536             parallel files (each worker uses 2 DuckDB threads)16    … --max-files 500                                bounded nightly batch17    … --manifest ~/firstratedata/state/row_groups.json1819Author: Simon-Pierre Boucher <contact@spboucher.ai>20"""21from __future__ import annotations2223import argparse24import json25import logging26import os27import sys28import tempfile29import threading30import time31from concurrent.futures import ThreadPoolExecutor, as_completed32from datetime import datetime, timezone33from pathlib import Path3435import duckdb3637log = logging.getLogger("rewrite_row_groups")3839ASSETS_WITH_ADJ = ("stock", "etf", "crypto", "index", "fx", "futures", "futures_contracts")40INTRADAY = ("1min", "5min", "30min", "1hour")41DEFAULT_ROW_GROUP = 65_536424344# ---- manifest -------------------------------------------------------------------------------------------4546class Manifest:47    """{"files": {"<relative path>": {"status": "done"|"failed"|"skipped", "row_groups": n, "rows": n, "at": iso, ...}}}"""4849    def __init__(self, path: Path) -> None:50        self.path = path51        self.lock = threading.Lock()52        self.data: dict = {"version": 1, "row_group_size": None, "files": {}}53        if path.is_file():54            try:55                self.data = json.loads(path.read_text())56            except json.JSONDecodeError:57                log.warning("manifest %s unreadable, starting over", path)58        self.data.setdefault("files", {})59        self._dirty = 06061    def status(self, rel: str) -> str | None:62        return (self.data["files"].get(rel) or {}).get("status")6364    def mark(self, rel: str, status: str, **extra) -> None:65        with self.lock:66            self.data["files"][rel] = {"status": status, "at": datetime.now(timezone.utc).isoformat(timespec="seconds"), **extra}67            self._dirty += 168            if self._dirty >= 20:69                self._flush()7071    def _flush(self) -> None:72        self.path.parent.mkdir(parents=True, exist_ok=True)73        tmp = self.path.with_suffix(".tmp")74        tmp.write_text(json.dumps(self.data, indent=1, sort_keys=True))75        os.replace(tmp, self.path)76        self._dirty = 07778    def flush(self) -> None:79        with self.lock:80            self._flush()818283# ---- discovery --------------------------------------------------------------------------------------------8485def candidate_files(parquet_root: Path, assets: list[str] | None, timeframes: list[str]) -> list[Path]:86    out: list[Path] = []87    for asset in assets or ASSETS_WITH_ADJ:88        for tf in timeframes:89            base = parquet_root / asset / tf90            if not base.is_dir():91                continue92            for adj_dir in sorted(p for p in base.iterdir() if p.is_dir()):93                with os.scandir(adj_dir) as it:94                    out.extend(Path(e.path) for e in it if e.name.endswith(".parquet") and e.is_file())95    return sorted(out)969798def parquet_layout(con: duckdb.DuckDBPyConnection, path: Path) -> tuple[int, int]:99    """(row_groups, rows) from the file footer — cheap."""100    p = str(path).replace("'", "''")101    rg = con.execute(f"SELECT count(DISTINCT row_group_id), coalesce(sum(row_group_num_rows), 0) "102                     f"FROM parquet_metadata('{p}') WHERE column_id = 0").fetchone()103    return int(rg[0]), int(rg[1])104105106# ---- rewrite ----------------------------------------------------------------------------------------------107108def rewrite_one(path: Path, row_group_size: int, threads: int, dry_run: bool) -> dict:109    con = duckdb.connect()110    con.execute(f"SET threads TO {threads}")111    con.execute("SET memory_limit = '2GB'")112    try:113        groups, rows = parquet_layout(con, path)114        if groups and rows and rows / groups <= row_group_size * 1.5:115            return {"status": "skipped", "reason": "row groups already small", "row_groups": groups, "rows": rows}116        if dry_run:117            return {"status": "dry-run", "row_groups": groups, "rows": rows,118                    "target_groups": -(-rows // row_group_size) if rows else 0}119        src = str(path).replace("'", "''")120        fd, tmp_name = tempfile.mkstemp(prefix=f".{path.stem}.", suffix=".parquet.tmp", dir=str(path.parent))121        os.close(fd)122        tmp = Path(tmp_name)123        try:124            t0 = time.perf_counter()125            dst = str(tmp).replace("'", "''")126            con.execute(f"COPY (SELECT * FROM read_parquet('{src}')) TO '{dst}' "127                        f"(FORMAT PARQUET, COMPRESSION ZSTD, ROW_GROUP_SIZE {int(row_group_size)})")128            new_groups, new_rows = parquet_layout(con, tmp)129            if new_rows != rows:130                raise RuntimeError(f"row count mismatch after rewrite: {rows} → {new_rows}")131            old_size, new_size = path.stat().st_size, tmp.stat().st_size132            os.replace(tmp, path)          # atomic on the same filesystem; readers see old or new, never partial133            return {"status": "done", "row_groups": new_groups, "rows": rows, "old_row_groups": groups,134                    "bytes_before": old_size, "bytes_after": new_size, "seconds": round(time.perf_counter() - t0, 2)}135        finally:136            if tmp.exists():137                tmp.unlink(missing_ok=True)138    finally:139        con.close()140141142def main() -> int:143    ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)144    ap.add_argument("--data-root", type=Path, default=Path(os.environ.get("HFMD_DATA_ROOT", "/Volumes/ssd/firstratedata")))145    ap.add_argument("--asset", action="append", help="restrict to an asset directory (repeatable): stock, etf, futures, futures_contracts, crypto, index, fx")146    ap.add_argument("--timeframe", action="append", help="restrict to a timeframe (repeatable): 1min, 5min, 30min, 1hour (default: all intraday)")147    ap.add_argument("--row-group-size", type=int, default=DEFAULT_ROW_GROUP)148    ap.add_argument("--workers", type=int, default=2, help="files rewritten in parallel (default 2)")149    ap.add_argument("--threads-per-worker", type=int, default=2)150    ap.add_argument("--max-files", type=int, default=0, help="stop after N files rewritten (0 = no bound)")151    ap.add_argument("--manifest", type=Path, help="progress file (default <data-root>/state/row_groups.json)")152    ap.add_argument("--retry-failed", action="store_true", help="re-attempt files marked failed in the manifest")153    ap.add_argument("--dry-run", action="store_true", help="list what would be rewritten, touch nothing")154    ap.add_argument("--verbose", "-v", action="store_true")155    args = ap.parse_args()156    logging.basicConfig(level=logging.INFO if args.verbose or args.dry_run else logging.WARNING,157                        format="%(asctime)s %(levelname)s %(message)s", stream=sys.stdout)158159    parquet_root = args.data_root / "parquet"160    if not parquet_root.is_dir():161        log.error("no parquet directory under %s", args.data_root)162        return 2163    tfs = args.timeframe or list(INTRADAY)164    bad = [t for t in tfs if t not in INTRADAY]165    if bad:166        log.error("only intraday timeframes are rewritten (%s); got %s", ", ".join(INTRADAY), ", ".join(bad))167        return 2168    manifest = Manifest(args.manifest or (args.data_root / "state" / "row_groups.json"))169    manifest.data["row_group_size"] = args.row_group_size170171    files = candidate_files(parquet_root, args.asset, tfs)172    todo = []173    for f in files:174        rel = str(f.relative_to(parquet_root))175        st = manifest.status(rel)176        if st in ("done", "skipped") or (st == "failed" and not args.retry_failed):177            continue178        todo.append((f, rel))179    if args.max_files:180        todo = todo[:args.max_files]181    log.info("%d candidate files, %d to process (%s)", len(files), len(todo), "dry run" if args.dry_run else "rewrite")182183    counts = {"done": 0, "skipped": 0, "failed": 0, "dry-run": 0}184    t0 = time.time()185186    def job(item):187        f, rel = item188        try:189            return rel, rewrite_one(f, args.row_group_size, args.threads_per_worker, args.dry_run)190        except Exception as e:  # noqa: BLE001 — one bad file must not stop the batch191            return rel, {"status": "failed", "error": f"{e.__class__.__name__}: {str(e)[:300]}"}192193    try:194        with ThreadPoolExecutor(max_workers=max(1, args.workers)) as ex:195            for rel, res in (fut.result() for fut in as_completed([ex.submit(job, it) for it in todo])):196                counts[res["status"]] = counts.get(res["status"], 0) + 1197                if res["status"] == "dry-run":198                    print(f"would rewrite {rel}: {res['row_groups']} groups → {res['target_groups']} ({res['rows']:,} rows)")199                    continue200                manifest.mark(rel, **res)201                if res["status"] == "failed":202                    log.error("FAILED %s: %s", rel, res["error"])203                else:204                    log.info("%s %s (%s)", res["status"], rel, ", ".join(f"{k}={v}" for k, v in res.items() if k not in ("status",)))205    finally:206        if not args.dry_run:207            manifest.flush()208    summary = {"processed": len(todo), **counts, "seconds": round(time.time() - t0, 1), "manifest": str(manifest.path)}209    print(json.dumps(summary))210    return 1 if counts.get("failed") else 0211212213if __name__ == "__main__":214    sys.exit(main())215