"""Row-group sizing: frd_downloader writes 65 536-row groups for intraday files (1 M for daily), and scripts/rewrite_row_groups.py rewrites existing files in place, resumably.""" from __future__ import annotations import importlib.util import json import sys from pathlib import Path import duckdb import pandas as pd import pytest ROOT = Path(__file__).resolve().parents[1] def _load(name: str, path: Path): spec = importlib.util.spec_from_file_location(name, path) mod = importlib.util.module_from_spec(spec) sys.modules[name] = mod spec.loader.exec_module(mod) return mod @pytest.fixture(scope="module") def frd(): if importlib.util.find_spec("requests") is None: # downloader-only dependency, not in the API venv import types stub = types.ModuleType("requests") stub.Session, stub.HTTPError, stub.RequestException = object, Exception, Exception # type: ignore[attr-defined] sys.modules.setdefault("requests", stub) return _load("frd_downloader_under_test", ROOT / "frd_downloader.py") @pytest.fixture(scope="module") def rewriter(): return _load("rewrite_row_groups_under_test", ROOT / "scripts" / "rewrite_row_groups.py") def test_row_group_size_by_timeframe(frd): pq = Path("/lake/parquet") assert frd.row_group_size(pq / "stock" / "1min" / "UNADJUSTED" / "AAPL_1min.parquet") == 65_536 assert frd.row_group_size(pq / "futures_contracts" / "1hour" / "update" / "ES_Z25_1hour.parquet") == 65_536 assert frd.row_group_size(pq / "stock" / "1day" / "adj_splitdiv" / "AAPL_1day.parquet") == 1_000_000 assert frd.row_group_size(pq / "options" / "2025_q2" / "AAPL_month_option_chain.parquet") == 1_000_000 def _big_file(path: Path, rows: int = 200_000) -> None: idx = pd.date_range("2024-01-02 09:30", periods=rows, freq="min") df = pd.DataFrame({"ticker": "AAA", "datetime": idx, "open": 1.0, "high": 1.0, "low": 1.0, "close": 1.0, "volume": 1.0}) path.parent.mkdir(parents=True, exist_ok=True) con = duckdb.connect() con.register("df", df) con.execute(f"COPY df TO '{path}' (FORMAT PARQUET, COMPRESSION ZSTD, ROW_GROUP_SIZE 1000000)") def test_rewrite_is_resumable_and_atomic(rewriter, tmp_path, monkeypatch, capsys): root = tmp_path / "lake" f1 = root / "parquet" / "stock" / "1min" / "UNADJUSTED" / "AAA_1min.parquet" f2 = root / "parquet" / "stock" / "5min" / "adj_splitdiv" / "AAA_5min.parquet" daily = root / "parquet" / "stock" / "1day" / "adj_splitdiv" / "AAA_1day.parquet" _big_file(f1) _big_file(f2, rows=100_000) _big_file(daily, rows=5_000) con = duckdb.connect() assert rewriter.parquet_layout(con, f1)[0] == 1 # dry run touches nothing monkeypatch.setattr(sys, "argv", ["x", "--data-root", str(root), "--dry-run"]) assert rewriter.main() == 0 assert rewriter.parquet_layout(con, f1)[0] == 1 and not (root / "state" / "row_groups.json").exists() assert "would rewrite" in capsys.readouterr().out # one file only (--max-files 1), then resume for the rest monkeypatch.setattr(sys, "argv", ["x", "--data-root", str(root), "--max-files", "1", "--workers", "1"]) assert rewriter.main() == 0 manifest = json.loads((root / "state" / "row_groups.json").read_text()) done = [k for k, v in manifest["files"].items() if v["status"] == "done"] assert len(done) == 1 monkeypatch.setattr(sys, "argv", ["x", "--data-root", str(root), "--workers", "2"]) assert rewriter.main() == 0 manifest = json.loads((root / "state" / "row_groups.json").read_text()) assert {v["status"] for v in manifest["files"].values()} == {"done"} assert len(manifest["files"]) == 2 # daily file never a candidate groups, rows = rewriter.parquet_layout(con, f1) assert rows == 200_000 and groups == -(-200_000 // 65_536) assert rewriter.parquet_layout(con, daily)[0] == 1 assert not list(f1.parent.glob(".*.tmp")) # no temp files left behind # idempotent: a third run processes nothing assert rewriter.main() == 0 assert json.loads(capsys.readouterr().out.strip().splitlines()[-1])["processed"] == 0