cc-repackage / scripts /build_examples.py
malteos
Add advanced SQL filter, precomputed Examples, and Max-records limit
c380737 unverified
Raw History Blame Contribute Delete
6.5 kB
#!/usr/bin/env python
"""One-time data-prep: build + host the precomputed Example datasets.
For each example defined in ``src/examples.py`` this:
1. runs the real index scan (``src.index_query`` in cdn mode) with the example's
SQL filter + record cap, writing a capped ``ranges.csv`` locally;
2. runs the real ``cdxt repackage`` (same local real-fetch path as
``scripts/verify_real.sh``) to produce the output WARC(s) locally;
3. uploads the WARC(s) + ranges.csv to the PUBLIC examples bucket;
4. prints the confirmed ``n_records`` / ``total_bytes`` / ``scan_files`` / filenames
to paste back into ``EXAMPLES`` in ``src/examples.py``.
Needs ``HF_TOKEN`` (bucket read of the CC index + write to the examples bucket) and a
positive-balance HF account is NOT required (no HF Jobs — everything runs locally).
Usage:
HF_TOKEN=... scripts/build_examples.py [--bucket ns/bucket] [--only key,key] [--keep-private]
"""
from __future__ import annotations
import argparse
import glob
import io
import os
import re
import subprocess
import sys
from contextlib import redirect_stdout
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from src import config, examples # noqa: E402
CDXT = ".venv/bin/cdxt"
CDXT_SPEC = (
"cdx_toolkit[hf] @ git+https://github.com/commoncrawl/cdx_toolkit.git@feat/warc-range-sources"
)
# Build-time only: cap the parquet files scanned for the broad no-domain examples so the
# build takes minutes, not hours (PDFs / homepages are common enough in any few files).
BUILD_MAX_FILES = {"pdfs": 8, "homepages": 4}
_ESTIMATE_RE = re.compile(r"ESTIMATE\s+n_records=(\d+)\s+total_bytes=(\d+)")
_SCAN_RE = re.compile(r"\[index\]\s+(\d+)\s+parquet file\(s\) to scan")
def _run_index(ex, ranges_csv: str) -> tuple[int, int, int, str]:
"""Run the real index scan for one example; return (n_records, total_bytes, scan_files, log)."""
from src import index_query
env = {
"CC_INDEX_MODE": "cdn",
"CC_CRAWLS": ",".join(ex.crawls),
"CC_HOSTNAMES": ",".join(ex.hostnames),
"CC_DOMAINS": ",".join(ex.domains),
"CC_LANGUAGES": ",".join(ex.languages),
"CC_SQL_WHERE": ex.sql_where,
"CC_MAX_RECORDS": str(ex.max_records or 0),
"CC_MAX_FILES": str(BUILD_MAX_FILES.get(ex.key, 0)),
"CC_RANGES_OUT": ranges_csv,
}
old = {k: os.environ.get(k) for k in env}
os.environ.update(env)
buf = io.StringIO()
try:
with redirect_stdout(buf):
rc = index_query.main()
finally:
for k, v in old.items():
if v is None:
os.environ.pop(k, None)
else:
os.environ[k] = v
log = buf.getvalue()
print(log)
if rc != 0:
raise SystemExit(f"index scan failed for {ex.key} (rc={rc})")
m = _ESTIMATE_RE.search(log)
n_records, total_bytes = (int(m.group(1)), int(m.group(2))) if m else (0, 0)
ms = _SCAN_RE.search(log)
scan_files = int(ms.group(1)) if ms else 0
return n_records, total_bytes, scan_files, log
def _bootstrap_cdxt() -> None:
if not os.access(CDXT, os.X_OK):
print("== installing cdx_toolkit for the fetch step ==")
subprocess.run(
["uv", "pip", "install", "--python", ".venv", CDXT_SPEC, "uvloop"], check=True
)
def _run_fetch(ex, ranges_csv: str, out_dir: str) -> list[str]:
"""Run the real cdxt repackage locally; return the produced WARC file paths."""
_bootstrap_cdxt()
out_prefix = os.path.join(out_dir, ex.name)
cmd = [
CDXT, "-v", "repackage",
"--target-source", "csv", "--csv-path", ranges_csv,
"--warc-download-prefix", f"hf://buckets/{config.CC_BUCKET}", "--hf-reader", "cdn",
"--prefix", out_prefix,
"--processes", "1", "--parallel_readers", str(config.DEFAULT_PARALLEL_READERS),
"--size", str(config.WARC_TARGET_SIZE),
]
env = dict(os.environ, CDXT_UVLOOP="1")
subprocess.run(cmd, check=True, env=env)
return sorted(glob.glob(f"{out_prefix}*.warc.gz"))
def _upload(api, bucket: str, ex, warcs: list[str], ranges_csv: str) -> list[str]:
add = [(w, f"{ex.path}/{os.path.basename(w)}") for w in warcs]
add.append((ranges_csv, f"{ex.path}/{config.RANGES_CSV_NAME}"))
api.batch_bucket_files(bucket, add=add, token=os.environ["HF_TOKEN"])
return [os.path.basename(w) for w in warcs]
def main() -> int:
ap = argparse.ArgumentParser()
ap.add_argument("--bucket", default=examples.EXAMPLES_BUCKET)
ap.add_argument("--only", default="", help="comma-separated example keys to build")
ap.add_argument("--keep-private", action="store_true", help="do NOT make the bucket public")
ap.add_argument("--out-dir", default="/tmp/cc-examples-build")
args = ap.parse_args()
token = os.environ.get("HF_TOKEN")
if not token:
raise SystemExit("set HF_TOKEN")
from huggingface_hub import HfApi
api = HfApi()
api.create_bucket(args.bucket, private=args.keep_private, exist_ok=True, token=token)
print(f"== examples bucket: {args.bucket} (public={not args.keep_private}) ==")
only = {k for k in args.only.split(",") if k}
todo = [ex for ex in examples.EXAMPLES if not only or ex.key in only]
os.makedirs(args.out_dir, exist_ok=True)
results = {}
for ex in todo:
print(f"\n===== building example: {ex.key} =====")
ranges_csv = os.path.join(args.out_dir, f"{ex.key}-ranges.csv")
ex_out = os.path.join(args.out_dir, ex.key)
os.makedirs(ex_out, exist_ok=True)
n_records, total_bytes, scan_files, _ = _run_index(ex, ranges_csv)
if n_records == 0:
print(f"WARNING: {ex.key} matched 0 records — skipping fetch/upload")
continue
warcs = _run_fetch(ex, ranges_csv, ex_out)
files = _upload(api, args.bucket, ex, warcs, ranges_csv)
results[ex.key] = dict(
n_records=n_records, total_bytes=total_bytes, scan_files=scan_files, files=files
)
print(f" -> {ex.key}: n_records={n_records} total_bytes={total_bytes} "
f"scan_files={scan_files} files={files}")
print("\n===== paste into src/examples.py (per example) =====")
for key, r in results.items():
print(f"# {key}: n_records={r['n_records']}, total_bytes={r['total_bytes']}, "
f"scan_files={r['scan_files']}, files={r['files']}")
return 0
if __name__ == "__main__":
sys.exit(main())