|
Download RUNBOOK.md from fzzhang/svd-code: direct link, hf CLI and curl.
- Browser
- Download file 23.6 kB
-
https://huggingface.co/fzzhang/svd-code/resolve/main/RUNBOOK.md
- Command line
-
hf download hf://fzzhang/svd-code/RUNBOOK.md
-
curl -L -o RUNBOOK.md https://huggingface.co/fzzhang/svd-code/resolve/main/RUNBOOK.md
23.6 kB
| # Multi-Node SDG Pipeline Runbook | |
| End-to-end recipe for running the SDG pipeline distributed across **4 worker | |
| nodes Γ 8 H100 GPUs = 32 shards**, with cross-node MongoDB caching and HuggingFace | |
| output gathering. Tested on ByteDance Arnold (Kubernetes-backed Azure ND96isr_H100_v5). | |
| - Target config: `sdg/configs/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1.yaml` | |
| - Final dataset: published to `<your-hf-username>/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1` | |
| - Wall time: **~40β60 hours** end-to-end | |
| - HF username placeholder used below: `fzzhang` β replace with yours | |
| If the cluster dies, follow this runbook top-to-bottom on a fresh allocation. | |
| --- | |
| ## 0. Cluster prerequisites | |
| - 4 worker nodes, each with 8Γ H100 (80 GB), shared no-filesystem (per-pod overlay). | |
| - NVIDIA driver supporting at least CUDA 12.4 (driver 535+ is fine). | |
| - HF account with write access to the target dataset namespace. | |
| - Each worker has these Arnold-provided env vars: `ARNOLD_WORKER_0_HOST`, | |
| `ARNOLD_WORKER_0_PORT`, `MY_POD_NAME` (ends with `worker-N`). | |
| - Cross-worker IPv6 traffic on the `ARNOLD_WORKER_*_PORT` ports is open by default β | |
| verify in Β§3 below. | |
| --- | |
| ## 1. Per-worker bootstrap (run on EVERY worker, in parallel) | |
| ### 1.1 Clone the repo | |
| ```bash | |
| cd ~/lmc_muon # or your preferred dir | |
| # git clone or rsync this repo into ./self-verified-distillation | |
| cd self-verified-distillation | |
| ``` | |
| ### 1.2 Create the conda env and install pinned dependencies | |
| The driver caps at CUDA 12.4, so vLLM 0.8.5 + torch 2.6.0+cu124 is the highest | |
| stack that runs. Newer vLLM (β₯0.9) ships cu126/cu130 wheels which **will not | |
| run on driver 535**. | |
| ```bash | |
| conda create -n sdg python=3.12 -y | |
| conda activate sdg | |
| pip install -U uv | |
| uv pip install -r sdg/requirements.txt | |
| # Downgrade to cu12.4-compatible stack | |
| uv pip uninstall vllm torch torchvision torchaudio flashinfer-python triton | |
| uv pip install torch==2.6.0 torchvision==0.21.0 \ | |
| --index-url https://download.pytorch.org/whl/cu124 | |
| uv pip install "vllm==0.8.5" "transformers>=4.51.3,<5.0" | |
| # Sanity check β MUST print "OK 1.0" | |
| python -c "import torch; x=torch.zeros(1).cuda(); print('OK', (x+1).item())" | |
| ``` | |
| ### 1.3 HuggingFace login | |
| ```bash | |
| huggingface-cli login # paste token; or set HF_TOKEN env var | |
| python -c "from huggingface_hub import HfApi; print(HfApi().whoami()['name'])" | |
| # should print your HF username (e.g. fzzhang) | |
| ``` | |
| --- | |
| ## 2. MongoDB on worker_0 (ONLY worker_0) | |
| ### 2.1 Install mongod | |
| ```bash | |
| cd ~/lmc_muon/self-verified-distillation | |
| curl -O https://fastdl.mongodb.org/linux/mongodb-linux-x86_64-ubuntu2204-8.0.4.tgz | |
| tar xzf mongodb-linux-x86_64-ubuntu2204-8.0.4.tgz | |
| export PATH="$PWD/mongodb-linux-x86_64-ubuntu2204-8.0.4/bin:$PATH" | |
| which mongod # must print a path | |
| # Persist for new shells | |
| echo "export PATH=\"$PWD/mongodb-linux-x86_64-ubuntu2204-8.0.4/bin:\$PATH\"" >> ~/.bashrc | |
| ``` | |
| ### 2.2 Start mongod bound to all interfaces with IPv6 | |
| The `--ipv6` flag is mandatory; without it `--bind_ip_all` is IPv4-only and | |
| cross-worker traffic (IPv6 fabric) cannot reach mongod. | |
| ```bash | |
| mkdir -p /tmp/sdg_mongo | |
| nohup mongod --port $ARNOLD_WORKER_0_PORT --dbpath /tmp/sdg_mongo \ | |
| --bind_ip_all --ipv6 > mongo.log 2>&1 & | |
| sleep 5 | |
| grep -i "Listening on" mongo.log | |
| # Expected: BOTH 0.0.0.0:<port> AND [::]:<port> | |
| # Local sanity | |
| python3 -c "from pymongo import MongoClient; \ | |
| MongoClient(f'mongodb://[::1]:$ARNOLD_WORKER_0_PORT', serverSelectionTimeoutMS=5000).admin.command('ping'); \ | |
| print('mongo OK locally')" | |
| ``` | |
| --- | |
| ## 3. Verify cross-worker MongoDB reachability | |
| On workers 1, 2, 3: | |
| ```bash | |
| python3 -c " | |
| import os | |
| from pymongo import MongoClient | |
| host = os.environ['ARNOLD_WORKER_0_HOST'] | |
| port = os.environ['ARNOLD_WORKER_0_PORT'] | |
| uri = f'mongodb://[{host}]:{port}' | |
| MongoClient(uri, serverSelectionTimeoutMS=5000).admin.command('ping') | |
| print(f'mongo OK from {os.uname().nodename} β {uri}') | |
| " | |
| ``` | |
| You should see `mongo OK from ...` on each. If any worker times out, the | |
| platform is blocking inter-worker traffic on that port β do not proceed. | |
| --- | |
| ## 4. Patch the config on each worker | |
| Point `mongo_uri` at the IPv6 MongoDB and set `num_shards=32` to use all GPUs. | |
| ```bash | |
| # Run on EACH worker | |
| python3 <<'EOF' | |
| import os, yaml | |
| host = os.environ['ARNOLD_WORKER_0_HOST'] | |
| port = os.environ['ARNOLD_WORKER_0_PORT'] | |
| path = 'sdg/configs/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1.yaml' | |
| with open(path) as f: cfg = yaml.safe_load(f) | |
| cfg['mongo_uri'] = f'mongodb://[{host}]:{port}' | |
| cfg['num_shards'] = 32 | |
| with open(path, 'w') as f: yaml.dump(cfg, f, sort_keys=False) | |
| print('updated:', cfg['mongo_uri']) | |
| EOF | |
| ``` | |
| --- | |
| ## 5. Helper scripts | |
| All orchestration scripts live in `scripts/` (committed to the repo, so each | |
| worker has them after `git clone`). They are designed to be **runnable as-is**; | |
| override behavior with env vars rather than editing the files: | |
| | Script | Purpose | Run on | | |
| |---|---|---| | |
| | `scripts/launch_shards.sh` | Launch 8 shards on this worker, one per GPU. Honors `CONFIG`, `LIMIT`, `NUM_SHARDS`. | every worker | | |
| | `scripts/upload_my_shards.py` | Poll for this worker's 8 stats.json, then upload to staging HF repo. Honors `HF_USER`, `RUN_NAME`. | every worker | | |
| | `scripts/combine_and_publish.py` | Poll for 4 worker-done markers, gather, combine, publish to HF (full launcher.py stats schema). Honors `HF_USER`, `RUN_NAME`, `NUM_WORKERS`, `NUM_SHARDS`. | worker_0 only | | |
| | `scripts/monitor.sh` | Per-worker live status (GPUs, current shard progress, completion counts, mongo cache size via cross-worker IPv6). | any worker | | |
| Defaults target the math53K config + `HF_USER=fzzhang` + 4 workers Γ 32 shards. | |
| **Why scripts/, not inline heredocs:** during the original run, terminal | |
| autocomplete corrupted multi-line heredoc pastes (Python files came out | |
| unparseable). Committed files bypass paste entirely β you just `git pull` and | |
| run them. | |
| **Verify they're present and parse on each worker:** | |
| ```bash | |
| ls scripts/ | |
| python -c "import ast; ast.parse(open('scripts/upload_my_shards.py').read()); ast.parse(open('scripts/combine_and_publish.py').read()); print('OK')" | |
| bash -n scripts/launch_shards.sh && bash -n scripts/monitor.sh && echo "shell OK" | |
| ``` | |
| If you need a different HF account or a different config, export the env vars | |
| before running (no file edits needed): | |
| ```bash | |
| export HF_USER=youraccount | |
| export RUN_NAME=your_run_name_matching_your_config | |
| export CONFIG=sdg/configs/your_config.yaml | |
| ``` | |
| --- | |
| ## 6. Per-worker uploader | |
| `scripts/upload_my_shards.py` runs on **every worker** β sleep-polls every 2 | |
| minutes for that worker's 8 `stats.json` to appear, then uploads each | |
| `shards/shard_NNN/` dir to `<HF_USER>/<RUN_NAME>-staging` and writes a | |
| `workers/worker_N.done` marker so the combiner knows that worker is finished. | |
| Consumes near-zero resources while waiting β safe to launch any time (before, | |
| during, or after the shards run). No file edits needed; override `HF_USER` / | |
| `RUN_NAME` via env if you're not using defaults (see Β§5). | |
| --- | |
| ## 7. Combiner (worker_0 only) | |
| `scripts/combine_and_publish.py` runs on **worker_0 only** β sleep-polls for all | |
| 4 worker-done markers, downloads the staging repo, concatenates the 32 | |
| `output.jsonl` files, aggregates stats with the **full launcher.py schema** | |
| (`summary`, `checks`, `failure_breakdown`, `failure_breakdown_examples`, | |
| `num_shards`), enforces a row-count vs `total_passed` integrity check, and | |
| publishes to `<HF_USER>/<RUN_NAME>` with a HuggingFace dataset-card README. | |
| Set `HF_USER` / `RUN_NAME` / `NUM_WORKERS` / `NUM_SHARDS` via env if defaults | |
| don't match. | |
| --- | |
| ## 8. Optional: smoke test before committing 40+ hours | |
| Tiny end-to-end validation (~15 min on all 32 GPUs). Confirms multi-worker | |
| plumbing without burning real time. Run **on each worker**: | |
| ```bash | |
| LIMIT=128 ./scripts/launch_shards.sh | |
| ``` | |
| Wait ~15 min, then check completion (on each worker): | |
| ```bash | |
| grep -l "PIPELINE COMPLETE" logs/worker_*/shard_*.log | wc -l # should be 8 | |
| ``` | |
| If all four workers report 8, plumbing is verified. Clear the smoke outputs | |
| before the real run: | |
| ```bash | |
| # On each worker | |
| rm -rf output/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1 | |
| rm -rf logs/worker_* | |
| # Optionally wipe the cache too (on worker_0) | |
| pkill mongod && sleep 2 | |
| rm -rf /tmp/sdg_mongo/* | |
| # Restart mongod (see Β§2.2) | |
| ``` | |
| --- | |
| ## 9. Production launch | |
| ### 9.1 Launch shards on each worker | |
| ```bash | |
| ./scripts/launch_shards.sh | |
| ``` | |
| ### 9.2 Start the per-worker uploader on each worker | |
| It will sleep-poll until its 8 shards complete, then upload. | |
| ```bash | |
| nohup python scripts/upload_my_shards.py > upload.log 2>&1 & | |
| disown | |
| ``` | |
| ### 9.3 Start the combiner on worker_0 ONLY | |
| It will sleep-poll until all 4 worker-done markers exist on HF, then combine | |
| and publish. | |
| ```bash | |
| nohup python scripts/combine_and_publish.py > combine.log 2>&1 & | |
| disown | |
| ``` | |
| After this, you can disconnect all terminals. Work continues under `nohup`. | |
| ### 9.4 (Recommended) Verify integrity once shards complete | |
| Before the uploader kicks off uploads, you can sanity-check each shard locally. | |
| For `selection_policy: first_valid` (the default), `output.jsonl` line count | |
| must equal `total_passed` in `stats.json`. Run **on each worker**: | |
| ```bash | |
| RUN=output/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1/shards | |
| for d in $RUN/shard_*; do | |
| lines=$(wc -l < $d/output.jsonl 2>/dev/null) | |
| passed=$(python3 -c "import json; print(json.load(open('$d/stats.json'))['summary']['total_passed'])" 2>/dev/null) | |
| echo "$(basename $d): $lines rows | total_passed=$passed" | |
| done | |
| ``` | |
| Lines must match `total_passed` for every shard. The combiner re-enforces this | |
| check globally and will `raise RuntimeError` if anything diverges β catching it | |
| locally first is faster. | |
| --- | |
| ## 10. Monitoring | |
| ### 10.1 Per-worker live view: `scripts/monitor.sh` | |
| ```bash | |
| watch -n 30 ./scripts/monitor.sh | |
| # Or if `watch` isn't installed in the container: | |
| while true; do clear; ./scripts/monitor.sh; sleep 30; done | |
| ``` | |
| Each refresh shows: per-GPU util, the most-recent tqdm state for each shard | |
| (decoded from `\r`-overwritten progress bars), per-worker completion counts, | |
| and the cluster-wide MongoDB cache size (queried over IPv6 so it works from any | |
| worker, not just worker_0). | |
| β **Do not use `tail -1 | cut -cN`** to read shard progress β tqdm writes one | |
| physical line with `\r` overwrites, so `cut -cN` shows the line's *beginning* | |
| (which never changes) while progress advances at the *end*. `monitor.sh` | |
| correctly does `tail -c 500 | tr '\r' '\n' | tail -1`. | |
| Expected sequence in the per-shard line over the run: | |
| 1. `Loading model:` / `Capturing CUDA graphs` (~2 min vLLM warmup) | |
| 2. `Processed prompts: N/512 [...]` (Stage 3 generation, batch size 64 Γ n=8 = 512) | |
| 3. `[Generation] N/1656 examples (...%)` (between-batch progress lines) | |
| 4. `Round 0: M pending rows` (Stage 4 UQ starts; tqdm batches now sized by pending rows) | |
| 5. `Cycle: ... Factual: ... Correctness: ...` (inside a UQ round) | |
| 6. `PIPELINE COMPLETE` (shard finished) | |
| To definitively confirm a shard is in UQ rather than generation, grep the log: | |
| ```bash | |
| grep -E "Round [0-9]+:|Cycle:|Factual:|Correctness:|\[Generation\]|PIPELINE" \ | |
| logs/worker_0/shard_0.log | tail -20 | |
| ``` | |
| You'll see `[Generation] 1656/1656 examples (100.0%)` followed by `Round 0:` if | |
| generation is done. | |
| ### 10.2 Aggregate signal: MongoDB cache size | |
| `monitor.sh` already prints this. Standalone: | |
| ```bash | |
| python3 -c " | |
| from pymongo import MongoClient | |
| import os | |
| host = os.environ['ARNOLD_WORKER_0_HOST'] | |
| port = os.environ['ARNOLD_WORKER_0_PORT'] | |
| print(MongoClient(f'mongodb://[{host}]:{port}').sdg_cache.inference_cache.count_documents({})) | |
| " | |
| ``` | |
| Grows continuously. Expect hundreds of thousands by mid-run, well over a | |
| million during UQ. **Use the IPv6 host, not `127.0.0.1`** β only worker_0 can | |
| reach loopback mongod. | |
| ### 10.3 Uploader / combiner status | |
| ```bash | |
| tail -f upload.log # on any worker (file lives in working dir) | |
| tail -f combine.log # on worker_0 | |
| ``` | |
| The combiner prints `DONE: https://huggingface.co/datasets/<HF_USER>/<RUN_NAME>` | |
| when finished β that's the success signal for the entire pipeline. | |
| --- | |
| ## 11. Troubleshooting | |
| ### Shards exit immediately | |
| - `Exit 127`: `mongod` (or `python`) not on PATH. `which mongod`, re-source `~/.bashrc`. | |
| - `Qwen2Tokenizer has no attribute all_special_tokens_extended`: transformers | |
| too new for vLLM 0.8.5. `uv pip install "transformers>=4.51.3,<5.0"`. | |
| - `CUDA driver version is insufficient`: torch+CUDA wheel doesn't match driver. | |
| Reinstall via Β§1.2 exactly. | |
| ### MongoDB connection refused from other workers | |
| - mongod log shows only `0.0.0.0:PORT` (no `[::]:PORT`): missing `--ipv6` flag. | |
| `pkill mongod`, restart with `--bind_ip_all --ipv6`. | |
| ### `monitor.sh` shows shards stuck at `0%` for hours | |
| - Almost always a display lie, not real stuckness. `tail -1 | cut -cN` would | |
| show the beginning of a long `\r`-overwritten line forever; the actual | |
| progress is at the *end*. Use `scripts/monitor.sh` (uses `tail -c 500 | tr '\r' '\n' | tail -1`). | |
| - Confirm work is happening: GPU util β₯ 80%, log file `mtime` updates, | |
| MongoDB cache count growing. | |
| ### `monitor.sh` shows `mongo cache: <unreachable>` from a non-zero worker | |
| - Workers other than worker_0 cannot reach `127.0.0.1:PORT`. The script uses | |
| `ARNOLD_WORKER_0_HOST` over IPv6, but if that env var isn't set in the | |
| current shell, the query fails. Re-source whatever sets Arnold env vars or | |
| run from a fresh login shell. | |
| ### Inline Python heredoc files won't parse | |
| - Terminal autocomplete corrupts `cat > file <<'EOF' ... EOF`. Symptoms: | |
| `unterminated string literal`, `import os, re, time Path` (extra word | |
| appended), or scripts cut off mid-line. | |
| - **Fix**: don't paste β use the committed files in `scripts/`. They survive | |
| `git pull`. If you must hand-write, use `nano <file>` (editor paste bypasses | |
| shell autocomplete) or base64-encode locally and `echo '...' | base64 -d > file`. | |
| ### Trial gets reclaimed mid-run | |
| - Shard outputs on local disk are lost (no shared FS) and the MongoDB cache at | |
| `/tmp/sdg_mongo` dies with worker_0. | |
| - **Mitigation**: start the uploaders in Β§9.2 *at launch time*, not at the | |
| end. They sleep-poll for completion and push completed work to HF | |
| incrementally. If a worker dies after some shards finish but before its 8 | |
| are all done, you lose only the unfinished ones. | |
| - Once a worker has written its `worker_N.done` marker on HF, that worker's | |
| data is durable. The combiner can run from any machine afterward, including | |
| a fresh trial. | |
| ### Some shards fail mid-run | |
| - Each shard is independent; just relaunch the failed ones. | |
| - Find failures: `grep -L "PIPELINE COMPLETE" logs/worker_*/shard_*.log` | |
| - Relaunch one: `CUDA_VISIBLE_DEVICES=I nohup python -m sdg.generate --config <CONFIG> --shard-id N --num-shards 32 > logs/worker_W/shard_N.log 2>&1 &` | |
| - MongoDB cache means most work is replayed for free on the rerun. | |
| ### Published `stats.json` is missing `failure_breakdown` fields | |
| - Old `combine_and_publish.py` versions wrote a simplified stats schema. | |
| Current `scripts/combine_and_publish.py` includes the full schema. | |
| - If you have a published dataset with the old stats, you can re-aggregate | |
| from `gathered/shards/*/stats.json` (still present on worker_0 after the | |
| combine) and re-upload `stats.json` β same logic as the combiner's stats | |
| section, just runs standalone. | |
| --- | |
| ## 12. After publish: verify and clean up | |
| ### 12.1 Verify the published dataset loads | |
| ```bash | |
| python3 -c " | |
| from datasets import load_dataset | |
| ds = load_dataset('fzzhang/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1', split='train') | |
| print('rows:', len(ds)) | |
| print('keys:', list(ds[0].keys())) | |
| print('roles:', [m['role'] for m in ds[0]['messages']]) | |
| " | |
| ``` | |
| Should print the row count (~21k), `keys: ['messages']`, and the roles | |
| (`['user', 'assistant']`). Confirms the HF dataset-card README correctly | |
| auto-loads `output.jsonl`. | |
| ### 12.2 Clean up local + intermediate state | |
| ```bash | |
| # On each worker β free disk if you plan to reuse the trial | |
| rm -rf output/ logs/ gathered/ combined_output.jsonl combined_stats.json upload.log | |
| # On worker_0 | |
| pkill mongod | |
| rm -rf /tmp/sdg_mongo combine.log | |
| # Optionally delete the intermediate staging HF repo | |
| HF_USER=fzzhang RUN_NAME=qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1 python -c " | |
| import os | |
| from huggingface_hub import HfApi | |
| repo = f\"{os.environ['HF_USER']}/{os.environ['RUN_NAME']}-staging\" | |
| HfApi().delete_repo(repo_id=repo, repo_type='dataset') | |
| print('deleted', repo) | |
| " | |
| ``` | |
| The final HF dataset is independent of the staging repo and the MongoDB | |
| cache β deleting them after the run is purely cosmetic. | |
| --- | |
| ## Step-by-step playbook | |
| Every command in execution order. The "Where" column tells you which worker to | |
| run the command on. Commands marked "every worker" should be run on workers 0, | |
| 1, 2, AND 3 (in parallel terminals to save wall time). Steps numbered with a | |
| letter (e.g. 5a) are part of the same logical step but run on different | |
| workers. | |
| ### Step 1 β get the code (every worker) | |
| ```bash | |
| cd ~/lmc_muon # or wherever you keep the repo | |
| git clone <your-repo-url> self-verified-distillation | |
| cd self-verified-distillation | |
| ``` | |
| ### Step 2 β Python env + pinned deps (every worker) | |
| ```bash | |
| conda create -n sdg python=3.12 -y | |
| conda activate sdg | |
| pip install -U uv | |
| uv pip install -r sdg/requirements.txt | |
| uv pip uninstall vllm torch torchvision torchaudio flashinfer-python triton | |
| uv pip install torch==2.6.0 torchvision==0.21.0 \ | |
| --index-url https://download.pytorch.org/whl/cu124 | |
| uv pip install "vllm==0.8.5" "transformers>=4.51.3,<5.0" | |
| python -c "import torch; x=torch.zeros(1).cuda(); print('OK', (x+1).item())" | |
| # Must print: OK 1.0 | |
| ``` | |
| ### Step 3 β HuggingFace login (every worker) | |
| ```bash | |
| huggingface-cli login # paste your token | |
| python -c "from huggingface_hub import HfApi; print(HfApi().whoami()['name'])" | |
| # Must print your HF username | |
| ``` | |
| ### Step 4 β install mongod (worker_0 ONLY) | |
| ```bash | |
| cd ~/lmc_muon/self-verified-distillation | |
| curl -O https://fastdl.mongodb.org/linux/mongodb-linux-x86_64-ubuntu2204-8.0.4.tgz | |
| tar xzf mongodb-linux-x86_64-ubuntu2204-8.0.4.tgz | |
| export PATH="$PWD/mongodb-linux-x86_64-ubuntu2204-8.0.4/bin:$PATH" | |
| echo "export PATH=\"$PWD/mongodb-linux-x86_64-ubuntu2204-8.0.4/bin:\$PATH\"" >> ~/.bashrc | |
| which mongod # must print a path | |
| ``` | |
| ### Step 5 β start mongod (worker_0 ONLY) | |
| ```bash | |
| mkdir -p /tmp/sdg_mongo | |
| nohup mongod --port $ARNOLD_WORKER_0_PORT --dbpath /tmp/sdg_mongo \ | |
| --bind_ip_all --ipv6 > mongo.log 2>&1 & | |
| sleep 5 | |
| grep -i "Listening on" mongo.log | |
| # Must show BOTH 0.0.0.0:<port> AND [::]:<port> | |
| python3 -c "from pymongo import MongoClient; \ | |
| MongoClient(f'mongodb://[::1]:$ARNOLD_WORKER_0_PORT', serverSelectionTimeoutMS=5000).admin.command('ping'); \ | |
| print('mongo OK locally')" | |
| ``` | |
| ### Step 6 β verify cross-worker mongo reachability (workers 1, 2, 3) | |
| ```bash | |
| python3 -c " | |
| import os | |
| from pymongo import MongoClient | |
| host = os.environ['ARNOLD_WORKER_0_HOST'] | |
| port = os.environ['ARNOLD_WORKER_0_PORT'] | |
| uri = f'mongodb://[{host}]:{port}' | |
| MongoClient(uri, serverSelectionTimeoutMS=5000).admin.command('ping') | |
| print(f'mongo OK from {os.uname().nodename} β {uri}') | |
| " | |
| ``` | |
| All three must print `mongo OK from ...`. If any times out, the platform is | |
| blocking inter-worker traffic β stop here. | |
| ### Step 7 β patch the config (every worker) | |
| ```bash | |
| python3 <<'EOF' | |
| import os, yaml | |
| host = os.environ['ARNOLD_WORKER_0_HOST'] | |
| port = os.environ['ARNOLD_WORKER_0_PORT'] | |
| path = 'sdg/configs/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1.yaml' | |
| with open(path) as f: cfg = yaml.safe_load(f) | |
| cfg['mongo_uri'] = f'mongodb://[{host}]:{port}' | |
| cfg['num_shards'] = 32 | |
| with open(path, 'w') as f: yaml.dump(cfg, f, sort_keys=False) | |
| print('updated:', cfg['mongo_uri']) | |
| EOF | |
| ``` | |
| ### Step 8 β verify scripts/ are intact (every worker) | |
| ```bash | |
| python -c "import ast; ast.parse(open('scripts/upload_my_shards.py').read()); ast.parse(open('scripts/combine_and_publish.py').read()); print('OK')" | |
| bash -n scripts/launch_shards.sh && bash -n scripts/monitor.sh && echo "shell OK" | |
| ``` | |
| ### Step 9 β (optional) smoke test (every worker) | |
| ```bash | |
| LIMIT=128 ./scripts/launch_shards.sh | |
| # Wait ~15-20 min, then check: | |
| grep -l "PIPELINE COMPLETE" logs/worker_*/shard_*.log | wc -l # must be 8 | |
| ``` | |
| If smoke passes, clear outputs before the real run: | |
| ```bash | |
| # Every worker: | |
| rm -rf output/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1 | |
| rm -rf logs/worker_* | |
| # Worker_0 only β wipe cache and restart mongod: | |
| pkill mongod && sleep 2 && rm -rf /tmp/sdg_mongo/* | |
| # Then redo Step 5 on worker_0 | |
| ``` | |
| ### Step 10 β launch the real run (every worker) | |
| ```bash | |
| ./scripts/launch_shards.sh | |
| # Expected: 8 PIDs printed. ~5 min later, all 8 GPUs should be at ~95% util. | |
| ``` | |
| ### Step 11 β start the per-worker uploader (every worker) | |
| Safe to launch immediately β it sleep-polls until the worker's shards finish. | |
| ```bash | |
| nohup python scripts/upload_my_shards.py > upload.log 2>&1 & | |
| disown | |
| ``` | |
| ### Step 12 β start the combiner (worker_0 ONLY) | |
| ```bash | |
| nohup python scripts/combine_and_publish.py > combine.log 2>&1 & | |
| disown | |
| ``` | |
| ### Step 13 β monitor (any worker, anytime) | |
| ```bash | |
| watch -n 30 ./scripts/monitor.sh | |
| # Or, if `watch` isn't installed: | |
| while true; do clear; ./scripts/monitor.sh; sleep 30; done | |
| ``` | |
| After this you can disconnect all terminals. The shards, uploaders, and | |
| combiner all run under `nohup`. | |
| ### Step 14 β (recommended) per-shard integrity check after shards finish (every worker) | |
| When `monitor.sh` shows `stats.json: 8/8` and `PIPELINE COMPLETE: 8/8` on a | |
| worker, sanity-check its outputs: | |
| ```bash | |
| RUN=output/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1/shards | |
| for d in $RUN/shard_*; do | |
| lines=$(wc -l < $d/output.jsonl 2>/dev/null) | |
| passed=$(python3 -c "import json; print(json.load(open('$d/stats.json'))['summary']['total_passed'])" 2>/dev/null) | |
| echo "$(basename $d): $lines rows | total_passed=$passed" | |
| done | |
| ``` | |
| `lines` must equal `total_passed` for every shard. If they diverge, the | |
| combiner will refuse to publish β fix the broken shard before its uploader | |
| starts pushing. | |
| ### Step 15 β wait for `DONE` (worker_0) | |
| ```bash | |
| tail -f combine.log | |
| # Watch for: DONE: https://huggingface.co/datasets/<HF_USER>/<RUN_NAME> | |
| ``` | |
| ### Step 16 β verify the published dataset loads (any machine with HF access) | |
| ```bash | |
| python3 -c " | |
| from datasets import load_dataset | |
| ds = load_dataset('fzzhang/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1', split='train') | |
| print('rows:', len(ds)) | |
| print('keys:', list(ds[0].keys())) | |
| print('roles:', [m['role'] for m in ds[0]['messages']]) | |
| " | |
| ``` | |
| ### Step 17 β (optional) cleanup | |
| ```bash | |
| # Every worker | |
| rm -rf output/ logs/ gathered/ combined_output.jsonl combined_stats.json upload.log | |
| # Worker_0 only | |
| pkill mongod | |
| rm -rf /tmp/sdg_mongo combine.log | |
| # Delete the intermediate staging HF repo (worker_0) | |
| HF_USER=fzzhang \ | |
| RUN_NAME=qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1 \ | |
| python -c " | |
| import os | |
| from huggingface_hub import HfApi | |
| repo = f\"{os.environ['HF_USER']}/{os.environ['RUN_NAME']}-staging\" | |
| HfApi().delete_repo(repo_id=repo, repo_type='dataset') | |
| print('deleted', repo) | |
| " | |
| ``` | |
| --- | |
| ## Success signal | |
| The whole pipeline succeeds when `combine.log` on worker_0 prints: | |
| ``` | |
| DONE: https://huggingface.co/datasets/<HF_USER>/<RUN_NAME> | |
| ``` | |
| Everything else (per-shard logs, mongo cache count, GPU util, staging repo, | |
| upload markers) is plumbing that supports that one line. | |