Buckets:
| """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.