Buckets:

Mercity/SkillsStorage / code /ship_delta.py
Pranav2748's picture
download
raw
3.56 kB
"""Ship finished crawler shards: convert -> upload raw + processed to the bucket -> delete local copies.
Keeps local disk flat during long crawls. Runs as a loop; idempotent (shipped.txt) and safe to restart.
The long-running parent only orchestrates. Every shard is converted AND uploaded by a short-lived child process
(`--child SHARD`), so all memory (pyarrow buffers, the hf-xet upload client's caches) is returned to the OS when
the child exits — an in-process uploader was seen to grow ~200 MB/hour. SHIP_WORKERS (default 3) children run concurrently.
Usage: python code/ship_delta.py RAW_DIR RUN_NAME [--once]
raw shard -> hf://buckets/Mercity/SkillsStorage/raw/github_delta/<RUN>/<shard>
processed -> hf://buckets/Mercity/SkillsStorage/processed/github_delta/<RUN>/{skills,files}/<shard>
repos.jsonl / done.txt / shipped.txt are re-uploaded each round (crawl log + resume state).
"""
import os, sys, time, subprocess, threading
from concurrent.futures import ThreadPoolExecutor
from pathlib import Path
ROOT = Path(__file__).resolve().parent.parent
RAW, RUN = Path(sys.argv[1]), sys.argv[2]
ONCE = "--once" in sys.argv
BUCKET = "Mercity/SkillsStorage"
TMP = ROOT / "work/ship" / RUN
shipped_f = RAW / "shipped.txt"
def upload(pairs):
"""pairs: [(local_path, remote_path_in_bucket)] uploaded in one batch call, with retries (child processes only)."""
from dotenv import load_dotenv; load_dotenv(ROOT / ".env")
from huggingface_hub import HfApi
api, err = HfApi(), None
for i in range(6):
try:
api.batch_bucket_files(BUCKET, add=[(str(l), r) for l, r in pairs]); return
except Exception as e:
err = e; time.sleep(15 * (i + 1))
raise RuntimeError(f"upload failed: {[str(l) for l, _ in pairs]}: {err}")
def child_ship(shard):
sys.path.insert(0, str(ROOT / "code"))
from build_github_delta import convert_shard
sk, fi, ns, nf = convert_shard(shard, TMP)
pairs = [(sk, f"processed/github_delta/{RUN}/skills/{shard.name}"), (shard, f"raw/github_delta/{RUN}/{shard.name}")]
if fi: pairs.append((fi, f"processed/github_delta/{RUN}/files/{shard.name}"))
upload(pairs)
sk.unlink(); fi and fi.unlink()
print(ns, nf)
if "--child" in sys.argv:
child_ship(Path(sys.argv[sys.argv.index("--child") + 1])); sys.exit(0)
if "--meta" in sys.argv:
upload([(RAW / m, f"raw/github_delta/{RUN}/{m}") for m in ("repos.jsonl", "done.txt", "shipped.txt") if (RAW / m).exists()])
sys.exit(0)
lock = threading.Lock()
def ship(shard):
r = subprocess.run([sys.executable, __file__, str(RAW), RUN, "--child", str(shard)], capture_output=True, text=True)
if r.returncode: raise RuntimeError(f"ship failed {shard.name}: {r.stderr[-400:]}")
ns, nf = map(int, r.stdout.split()[-2:])
with lock:
with open(shipped_f, "a") as f: f.write(shard.name + "\n")
shard.unlink()
print(f"shipped {shard.name}: {ns} skills, {nf} files", flush=True)
while True:
shipped = set(shipped_f.read_text().split()) if shipped_f.exists() else set()
ready = [p for p in sorted(RAW.glob("files-*.parquet")) if p.name not in shipped and (ONCE or time.time() - p.stat().st_mtime > 30)] # shards are written atomically
with ThreadPoolExecutor(int(os.environ.get("SHIP_WORKERS", "3"))) as ex: # uploads are latency-bound: more in flight
for f in [ex.submit(ship, s) for s in ready]: f.result()
subprocess.run([sys.executable, __file__, str(RAW), RUN, "--meta"], capture_output=True)
if ONCE: break
time.sleep(60)

Xet Storage Details

Size:
3.56 kB
·
Xet hash:
c1f99fe9cfbcafd45c3fb2f733f227372054eabf52a6e893ef788e915e548a6b

Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.