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)
JavaScript 53.7%
Python 38.3%
CSS 4.6%
TypeScript 3.1%
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