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%
4.1 KB · 92 lines python
Raw Blame History
1"""Row-group sizing: frd_downloader writes 65 536-row groups for intraday files (1 M for daily), and2scripts/rewrite_row_groups.py rewrites existing files in place, resumably."""3from __future__ import annotations45import importlib.util6import json7import sys8from pathlib import Path910import duckdb11import pandas as pd12import pytest1314ROOT = Path(__file__).resolve().parents[1]151617def _load(name: str, path: Path):18    spec = importlib.util.spec_from_file_location(name, path)19    mod = importlib.util.module_from_spec(spec)20    sys.modules[name] = mod21    spec.loader.exec_module(mod)22    return mod232425@pytest.fixture(scope="module")26def frd():27    if importlib.util.find_spec("requests") is None:      # downloader-only dependency, not in the API venv28        import types29        stub = types.ModuleType("requests")30        stub.Session, stub.HTTPError, stub.RequestException = object, Exception, Exception   # type: ignore[attr-defined]31        sys.modules.setdefault("requests", stub)32    return _load("frd_downloader_under_test", ROOT / "frd_downloader.py")333435@pytest.fixture(scope="module")36def rewriter():37    return _load("rewrite_row_groups_under_test", ROOT / "scripts" / "rewrite_row_groups.py")383940def test_row_group_size_by_timeframe(frd):41    pq = Path("/lake/parquet")42    assert frd.row_group_size(pq / "stock" / "1min" / "UNADJUSTED" / "AAPL_1min.parquet") == 65_53643    assert frd.row_group_size(pq / "futures_contracts" / "1hour" / "update" / "ES_Z25_1hour.parquet") == 65_53644    assert frd.row_group_size(pq / "stock" / "1day" / "adj_splitdiv" / "AAPL_1day.parquet") == 1_000_00045    assert frd.row_group_size(pq / "options" / "2025_q2" / "AAPL_month_option_chain.parquet") == 1_000_000464748def _big_file(path: Path, rows: int = 200_000) -> None:49    idx = pd.date_range("2024-01-02 09:30", periods=rows, freq="min")50    df = pd.DataFrame({"ticker": "AAA", "datetime": idx, "open": 1.0, "high": 1.0, "low": 1.0, "close": 1.0, "volume": 1.0})51    path.parent.mkdir(parents=True, exist_ok=True)52    con = duckdb.connect()53    con.register("df", df)54    con.execute(f"COPY df TO '{path}' (FORMAT PARQUET, COMPRESSION ZSTD, ROW_GROUP_SIZE 1000000)")555657def test_rewrite_is_resumable_and_atomic(rewriter, tmp_path, monkeypatch, capsys):58    root = tmp_path / "lake"59    f1 = root / "parquet" / "stock" / "1min" / "UNADJUSTED" / "AAA_1min.parquet"60    f2 = root / "parquet" / "stock" / "5min" / "adj_splitdiv" / "AAA_5min.parquet"61    daily = root / "parquet" / "stock" / "1day" / "adj_splitdiv" / "AAA_1day.parquet"62    _big_file(f1)63    _big_file(f2, rows=100_000)64    _big_file(daily, rows=5_000)65    con = duckdb.connect()66    assert rewriter.parquet_layout(con, f1)[0] == 16768    # dry run touches nothing69    monkeypatch.setattr(sys, "argv", ["x", "--data-root", str(root), "--dry-run"])70    assert rewriter.main() == 071    assert rewriter.parquet_layout(con, f1)[0] == 1 and not (root / "state" / "row_groups.json").exists()72    assert "would rewrite" in capsys.readouterr().out7374    # one file only (--max-files 1), then resume for the rest75    monkeypatch.setattr(sys, "argv", ["x", "--data-root", str(root), "--max-files", "1", "--workers", "1"])76    assert rewriter.main() == 077    manifest = json.loads((root / "state" / "row_groups.json").read_text())78    done = [k for k, v in manifest["files"].items() if v["status"] == "done"]79    assert len(done) == 180    monkeypatch.setattr(sys, "argv", ["x", "--data-root", str(root), "--workers", "2"])81    assert rewriter.main() == 082    manifest = json.loads((root / "state" / "row_groups.json").read_text())83    assert {v["status"] for v in manifest["files"].values()} == {"done"}84    assert len(manifest["files"]) == 2                      # daily file never a candidate85    groups, rows = rewriter.parquet_layout(con, f1)86    assert rows == 200_000 and groups == -(-200_000 // 65_536)87    assert rewriter.parquet_layout(con, daily)[0] == 188    assert not list(f1.parent.glob(".*.tmp"))               # no temp files left behind89    # idempotent: a third run processes nothing90    assert rewriter.main() == 091    assert json.loads(capsys.readouterr().out.strip().splitlines()[-1])["processed"] == 092