diff --git a/benchmarks/pin_memory/README.md b/benchmarks/pin_memory/README.md index 8410e4277..7af4abc59 100644 --- a/benchmarks/pin_memory/README.md +++ b/benchmarks/pin_memory/README.md @@ -1,49 +1,36 @@ -# ZeRO-3 CPU-Offload Pinned-Memory Benchmark +# Pinned-memory experiments -This directory contains an end-to-end benchmark for ZeRO-3 CPU offload that -measures training step time with pinned vs unpinned host memory, plus an -opt-in ablation of registered vs unregistered pinned memory. +Harnesses for pin vs pageable host memory on CPU offload. Each subdirectory is +one experiment. Shared subprocess/JSON helpers live in `common.py`. -## Files in this Directory +Native backends `mlock` host memory: raise `RLIMIT_MEMLOCK` (`ulimit -l`) or +run as root for multi-GB models. -- **zero3_offload_bench.py**: Benchmarking script; the model can be a real - architecture fetched from the HuggingFace hub (random weights) or a - synthetic MLP stack that needs no network access. +## Layout -## What it Measures +| Folder | Blog experiment | Default command | +|--------|-----------------|-----------------| +| [`model_tensor_offload/`](model_tensor_offload/) | ZeRO CPU param/optimizer offload (stage 3 default; `--zero-stage 1\|2` optional) | `python model_tensor_offload/bench.py --hidden 2048 --layers 12 --batch 4 --seq 128` | +| [`activation_offload/`](activation_offload/) | Checkpoint hidden-state offload; `use_pin_memory` on/off, **async on** | `python activation_offload/bench.py --hidden 1024 --layers 8 --batch 1 --seq 2048` | +| [`h2d_d2h/`](h2d_d2h/) | Supporting H2D/D2H GB/s (pageable, torch, native-unregistered, native-registered) | `python h2d_d2h/bench.py` | +| [`grad_offload/`](grad_offload/) | Optional #8207-style grad offload (wraps model-tensor ZeRO-3) | `python grad_offload/bench.py --hidden 2048 --layers 12` | +| [`cpu_pin/`](cpu_pin/) | Optional CPU-only native vs Torch pin | `python cpu_pin/bench.py` | +| [`deepcompile_activation/`](deepcompile_activation/) | Optional `compile.offload_activation_pin_memory` | `python deepcompile_activation/bench.py` | -By default the script runs ZeRO-3 with `offload_optimizer` and `offload_param` -(both CPU) in two arms and reports the step-time comparison: +`zero3_offload_bench.py` at this directory root still runs **model-tensor ZeRO-3** (same flags as before). -| Arm | offload `pin_memory` | `DS_PIN_MEMORY_REGISTER_DEVICE` | -|-----|----------------------|---------------------------------| -| `unpinned` | `False` | (n/a) | -| `pinned` | `True` (`DS_PIN_MEMORY_BACKEND=native`) | `1` | +## Model-tensor arms -Works on any accelerator with native pin + `register_host_memory` support -(CUDA and XPU are tested). Each arm runs in its own subprocess with a fresh -rendezvous port so device state never leaks between arms. +| Arm | `offload_*.pin_memory` | Backend | +|-----|------------------------|---------| +| unpinned | `false` | n/a | +| pinned | `true` | `DS_PIN_MEMORY_BACKEND=native`, `DS_PIN_MEMORY_REGISTER_DEVICE=1` | -Power users can additionally ablate device registration of pinned buffers: +`--ablate-register` adds `pinned-unregistered`. CUDA-oriented; skip on XPU if `register_host_memory` is missing (`h2d_d2h/bench.py --skip-native-register`). ```bash -python zero3_offload_bench.py --ablate-register ... +python model_tensor_offload/bench.py --model Qwen/Qwen2.5-7B --batch 4 --seq 512 +python model_tensor_offload/bench.py --zero-stage 2 --hidden 2048 --layers 12 --batch 4 --seq 128 ``` -which adds a `pinned-unregistered` arm (`DS_PIN_MEMORY_REGISTER_DEVICE=0`). - -## Usage - -```bash -# real model architecture (config fetched from the HF hub, random weights) -python zero3_offload_bench.py --model Qwen/Qwen2.5-7B --batch 4 --seq 512 - -# synthetic MLP stack, no network needed -python zero3_offload_bench.py --hidden 2048 --layers 12 --batch 4 --seq 128 -``` - -Results are printed as a table (avg/min step time, GPU peak memory) and as a -JSON line (`DRIVERRESULT=...`) with per-arm details and the pinning speedup. - -> **Note**: the native backend mlocks host memory; raise `RLIMIT_MEMLOCK` -> (`ulimit -l`) or run as root for multi-GB models. +Each arm is a subprocess with a fresh rendezvous port. Results: table plus `DRIVERRESULT=` JSON. diff --git a/benchmarks/pin_memory/activation_offload/README.md b/benchmarks/pin_memory/activation_offload/README.md new file mode 100644 index 000000000..3e4c23025 --- /dev/null +++ b/benchmarks/pin_memory/activation_offload/README.md @@ -0,0 +1,7 @@ +# Activation / checkpoint hidden-state offload + +Compares `use_pin_memory` True vs False on `CheckpointHiddenStatesOffload` +with **async / side streams held on**. This is not the async-vs-blocking +table from DeepSpeed #8282. + +See the [parent README](../README.md). diff --git a/benchmarks/pin_memory/activation_offload/bench.py b/benchmarks/pin_memory/activation_offload/bench.py new file mode 100644 index 000000000..ca01dbb7a --- /dev/null +++ b/benchmarks/pin_memory/activation_offload/bench.py @@ -0,0 +1,162 @@ +# SPDX-License-Identifier: Apache-2.0 +# DeepSpeed Team +""" +Activation / checkpoint hidden-state CPU offload: pin vs pageable with async on. + +Holds use_streams=True. Compares use_pin_memory True vs False. Do not treat +this as DeepSpeed #8282's async-vs-blocking table. +""" + +from __future__ import annotations + +import argparse +import os +import sys +import time + +_PIN_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +if _PIN_ROOT not in sys.path: + sys.path.insert(0, _PIN_ROOT) + +from common import dist_env, patch_cpp_extension_drop_cxx17, print_arm_result, print_driver_result, run_arm_subprocess + + +def parse_args(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--hidden", type=int, default=1024) + parser.add_argument("--layers", type=int, default=8) + parser.add_argument("--batch", type=int, default=1) + parser.add_argument("--seq", type=int, default=2048) + parser.add_argument("--steps", type=int, default=4) + parser.add_argument("--warmup", type=int, default=2) + parser.add_argument("--pin", type=int, default=None, help="internal: use_pin_memory 0/1") + return parser.parse_args() + + +def run_arm(args): + dist_env() + patch_cpp_extension_drop_cxx17() + + import torch + from torch.utils.checkpoint import checkpoint + + from deepspeed.accelerator import get_accelerator + from deepspeed.runtime.activation_checkpointing.offload_activations import CheckpointHiddenStatesOffload + + accelerator = get_accelerator() + if not accelerator.is_available(): + raise RuntimeError(f"No {accelerator.device_name()} device is available") + accelerator.set_device(0) + device = accelerator.current_device_name() + + class Block(torch.nn.Module): + + def __init__(self, hidden): + super().__init__() + self.fc1 = torch.nn.Linear(hidden, 4 * hidden) + self.fc2 = torch.nn.Linear(4 * hidden, hidden) + + def forward(self, hidden_states): + return self.fc2(torch.nn.functional.gelu(self.fc1(hidden_states))) + + class Net(torch.nn.Module): + + def __init__(self, hidden, layers): + super().__init__() + self.blocks = torch.nn.ModuleList([Block(hidden) for _ in range(layers)]) + + def forward(self, hidden_states, offload): + x = hidden_states + for block in self.blocks: + offload.mark(x) + x = x + checkpoint(block, x, use_reentrant=False) + return x.sum() + + model = Net(args.hidden, args.layers).to(device) + opt = torch.optim.AdamW(model.parameters(), lr=1e-4) + x = torch.randn(args.batch, args.seq, args.hidden, device=device, requires_grad=True) + + # Async side stream stays on; pin vs pageable is the only axis. + offload = CheckpointHiddenStatesOffload(use_pin_memory=bool(args.pin), + use_streams=True, + min_offload_bytes=0, + keep_last_count=1) + + step_times = [] + for step in range(args.warmup + args.steps): + accelerator.synchronize() + t0 = time.perf_counter() + opt.zero_grad(set_to_none=True) + with offload: + loss = model(x, offload) + loss.backward() + opt.step() + accelerator.synchronize() + if step >= args.warmup: + step_times.append(time.perf_counter() - t0) + offload.reset() + + def _peak(): + try: + return round(torch.get_device_module(device).max_memory_allocated() / 1e9, 2) + except Exception: + return None + + print_arm_result({ + "experiment": "activation_offload", + "use_pin_memory": bool(args.pin), + "use_streams": True, + "device": accelerator.device_name(), + "hidden": args.hidden, + "layers": args.layers, + "batch": args.batch, + "seq": args.seq, + "steps": len(step_times), + "step_avg_s": sum(step_times) / len(step_times), + "step_min_s": min(step_times), + "gpu_peak_gb": _peak(), + }) + + +def run_driver(args): + script = os.path.abspath(__file__) + base = [ + "--hidden", + str(args.hidden), + "--layers", + str(args.layers), + "--batch", + str(args.batch), + "--seq", + str(args.seq), + "--steps", + str(args.steps), + "--warmup", + str(args.warmup), + ] + results = {} + for name, pin in (("pageable", 0), ("pinned", 1)): + results[name] = run_arm_subprocess(script, base + ["--pin", str(pin)]) + + pageable = results["pageable"] + pinned = results["pinned"] + print("\n================ Activation offload (async on) ================") + print(f"device: {pinned['device']} hidden: {pinned['hidden']} layers: {pinned['layers']} " + f"batch: {pinned['batch']} seq: {pinned['seq']}") + print(f"{'arm':<22}{'avg step (s)':>14}{'min step (s)':>14}{'GPU peak (GB)':>16}") + for name in ("pageable", "pinned"): + row = results[name] + print(f"{name:<22}{row['step_avg_s']:>14.3f}{row['step_min_s']:>14.3f}" + f"{(row['gpu_peak_gb'] or 0):>16.2f}") + speedup = pageable["step_avg_s"] / pinned["step_avg_s"] + print() + print(f"pin vs pageable step-time ratio: {speedup:.2f}x (async held on)") + print_driver_result({"pageable": pageable, "pinned": pinned, "speedup": speedup}) + + +if __name__ == "__main__": + parsed = parse_args() + if parsed.pin is None: + run_driver(parsed) + else: + run_arm(parsed) diff --git a/benchmarks/pin_memory/common.py b/benchmarks/pin_memory/common.py new file mode 100644 index 000000000..b33424ea1 --- /dev/null +++ b/benchmarks/pin_memory/common.py @@ -0,0 +1,74 @@ +# SPDX-License-Identifier: Apache-2.0 +# DeepSpeed Team +"""Shared helpers for pin_memory experiment drivers (subprocess arms, JSON lines).""" + +from __future__ import annotations + +import json +import os +import socket +import subprocess +import sys + + +def free_port(): + # A stale listener from an interrupted rank makes the next init hang in a + # collective, so always rendezvous on a fresh ephemeral port. + sock = socket.socket() + sock.bind(("127.0.0.1", 0)) + port = sock.getsockname()[1] + sock.close() + return port + + +def dist_env(): + os.environ.update( + MASTER_ADDR="127.0.0.1", + MASTER_PORT=str(free_port()), + RANK="0", + WORLD_SIZE="1", + LOCAL_RANK="0", + ) + + +def patch_cpp_extension_drop_cxx17(): + # torch-nightly requires C++20; SYCL toolchain flags may carry -std=c++17, + # which (appearing last) downgrades the dialect and breaks torch headers. + import torch.utils.cpp_extension as cpp_ext + + orig_load = cpp_ext.load + + def load_without_cxx17(*args, **kwargs): + for key in ("extra_cflags", "extra_cxxflags"): + if kwargs.get(key): + kwargs[key] = [flag for flag in kwargs[key] if flag != "-std=c++17"] + return orig_load(*args, **kwargs) + + cpp_ext.load = load_without_cxx17 + + +def print_arm_result(result): + print("ARMRESULT=" + json.dumps(result), flush=True) + + +def print_driver_result(summary): + print("DRIVERRESULT=" + json.dumps(summary), flush=True) + + +def parse_arm_result(stdout): + for line in stdout.splitlines(): + if line.startswith("ARMRESULT="): + return json.loads(line[len("ARMRESULT="):]) + return None + + +def run_arm_subprocess(script_path, extra_args, env=None): + command = [sys.executable, os.path.abspath(script_path), *extra_args] + print(f"[driver] {' '.join(command)}", flush=True) + proc = subprocess.run(command, env=env or os.environ.copy(), capture_output=True, text=True) + arm = parse_arm_result(proc.stdout) + if arm is None: + print(proc.stdout[-2000:]) + print(proc.stderr[-2000:], file=sys.stderr) + raise RuntimeError(f"arm produced no ARMRESULT (rc={proc.returncode})") + return arm diff --git a/benchmarks/pin_memory/cpu_pin/README.md b/benchmarks/pin_memory/cpu_pin/README.md new file mode 100644 index 000000000..c54674ae1 --- /dev/null +++ b/benchmarks/pin_memory/cpu_pin/README.md @@ -0,0 +1,6 @@ +# CPU-only pin + +Native `mlock` vs Torch pin on a CPU host. Torch typically cannot pin without +an accelerator-capable backend. + +See the [parent README](../README.md). diff --git a/benchmarks/pin_memory/cpu_pin/bench.py b/benchmarks/pin_memory/cpu_pin/bench.py new file mode 100644 index 000000000..e34b3d30a --- /dev/null +++ b/benchmarks/pin_memory/cpu_pin/bench.py @@ -0,0 +1,59 @@ +# SPDX-License-Identifier: Apache-2.0 +# DeepSpeed Team +"""CPU-only: native pin vs Torch (Torch cannot pin without an accelerator).""" + +from __future__ import annotations + +import argparse +import os +import sys + +_PIN_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +if _PIN_ROOT not in sys.path: + sys.path.insert(0, _PIN_ROOT) + +from common import print_arm_result, print_driver_result + + +def parse_args(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--numel", type=int, default=1024 * 1024) + return parser.parse_args() + + +def main(): + args = parse_args() + import torch + from deepspeed.accelerator import get_accelerator + + accel = get_accelerator() + os.environ["DS_PIN_MEMORY_BACKEND"] = "native" + host = torch.empty(args.numel, dtype=torch.float32) + native = accel.pin_memory(host.clone(), make_copy=False) + native_ok = bool(accel.is_pinned(native)) + accel.unpin_memory(native) + + os.environ["DS_PIN_MEMORY_BACKEND"] = "torch" + torch_ok = None + torch_error = None + try: + pinned = accel.pin_memory(torch.empty_like(host), make_copy=False) + torch_ok = bool(accel.is_pinned(pinned)) + except Exception as exc: + torch_error = type(exc).__name__ + ": " + str(exc) + + result = { + "experiment": "cpu_pin", + "accelerator": accel.device_name(), + "native_is_pinned": native_ok, + "torch_is_pinned": torch_ok, + "torch_error": torch_error, + } + print_arm_result(result) + print_driver_result(result) + if not native_ok: + raise SystemExit("native pin failed on this host") + + +if __name__ == "__main__": + main() diff --git a/benchmarks/pin_memory/deepcompile_activation/README.md b/benchmarks/pin_memory/deepcompile_activation/README.md new file mode 100644 index 000000000..d6efd33ee --- /dev/null +++ b/benchmarks/pin_memory/deepcompile_activation/README.md @@ -0,0 +1,5 @@ +# DeepCompile activation pin + +Optional `compile.offload_activation_pin_memory` on/off. Requires DeepCompile. + +See the [parent README](../README.md). diff --git a/benchmarks/pin_memory/deepcompile_activation/bench.py b/benchmarks/pin_memory/deepcompile_activation/bench.py new file mode 100644 index 000000000..1f435e532 --- /dev/null +++ b/benchmarks/pin_memory/deepcompile_activation/bench.py @@ -0,0 +1,116 @@ +# SPDX-License-Identifier: Apache-2.0 +# DeepSpeed Team +"""DeepCompile activation offload pin knob (compile.offload_activation_pin_memory).""" + +from __future__ import annotations + +import argparse +import os +import sys + +_PIN_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +if _PIN_ROOT not in sys.path: + sys.path.insert(0, _PIN_ROOT) + +from common import dist_env, patch_cpp_extension_drop_cxx17, print_arm_result, print_driver_result, run_arm_subprocess + + +def parse_args(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--hidden", type=int, default=1024) + parser.add_argument("--layers", type=int, default=4) + parser.add_argument("--batch", type=int, default=2) + parser.add_argument("--seq", type=int, default=128) + parser.add_argument("--steps", type=int, default=3) + parser.add_argument("--warmup", type=int, default=1) + parser.add_argument("--pin", type=int, default=None) + return parser.parse_args() + + +def run_arm(args): + dist_env() + patch_cpp_extension_drop_cxx17() + import time + import torch + import deepspeed + + class Net(torch.nn.Module): + + def __init__(self, hidden, layers): + super().__init__() + self.layers = torch.nn.ModuleList([torch.nn.Linear(hidden, hidden) for _ in range(layers)]) + + def forward(self, x): + for layer in self.layers: + x = torch.nn.functional.gelu(layer(x)) + return x.sum() + + model = Net(args.hidden, args.layers) + ds_config = { + "train_micro_batch_size_per_gpu": args.batch, + "gradient_accumulation_steps": 1, + "bf16": { + "enabled": True + }, + "optimizer": { + "type": "AdamW", + "params": { + "lr": 1e-4 + } + }, + "zero_optimization": { + "stage": 3 + }, + "compile": { + "enabled": True, + "offload_activation": True, + "offload_activation_pin_memory": bool(args.pin), + }, + } + engine, _, _, _ = deepspeed.initialize(model=model, config=ds_config) + dev = engine.device + x = torch.randn(args.batch, args.seq, args.hidden, device=dev) + times = [] + for step in range(args.warmup + args.steps): + t0 = time.perf_counter() + loss = engine(x) + engine.backward(loss) + engine.step() + if step >= args.warmup: + times.append(time.perf_counter() - t0) + print_arm_result({ + "experiment": "deepcompile_activation", + "offload_activation_pin_memory": bool(args.pin), + "device": str(dev), + "step_avg_s": sum(times) / len(times), + }) + + +def run_driver(args): + script = os.path.abspath(__file__) + base = [ + "--hidden", + str(args.hidden), + "--layers", + str(args.layers), + "--batch", + str(args.batch), + "--seq", + str(args.seq), + "--steps", + str(args.steps), + "--warmup", + str(args.warmup), + ] + results = {} + for name, pin in (("pageable", 0), ("pinned", 1)): + results[name] = run_arm_subprocess(script, base + ["--pin", str(pin)]) + print_driver_result(results) + + +if __name__ == "__main__": + parsed = parse_args() + if parsed.pin is None: + run_driver(parsed) + else: + run_arm(parsed) diff --git a/benchmarks/pin_memory/grad_offload/README.md b/benchmarks/pin_memory/grad_offload/README.md new file mode 100644 index 000000000..b05757441 --- /dev/null +++ b/benchmarks/pin_memory/grad_offload/README.md @@ -0,0 +1,4 @@ +# Gradient CPU offload + +Optional pin on/off for ZeRO-3 CPU offload (same as model-tensor stage 3). +See the [parent README](../README.md). diff --git a/benchmarks/pin_memory/grad_offload/bench.py b/benchmarks/pin_memory/grad_offload/bench.py new file mode 100644 index 000000000..20e421b88 --- /dev/null +++ b/benchmarks/pin_memory/grad_offload/bench.py @@ -0,0 +1,12 @@ +# SPDX-License-Identifier: Apache-2.0 +# DeepSpeed Team +"""Async gradient CPU offload pin on/off (#8207). Reuses model-tensor ZeRO-3 offload.""" + +from __future__ import annotations + +import os +import sys + +_HERE = os.path.dirname(os.path.abspath(__file__)) +_MODEL = os.path.join(os.path.dirname(_HERE), "model_tensor_offload", "bench.py") +os.execv(sys.executable, [sys.executable, _MODEL, "--zero-stage", "3", *sys.argv[1:]]) diff --git a/benchmarks/pin_memory/h2d_d2h/README.md b/benchmarks/pin_memory/h2d_d2h/README.md new file mode 100644 index 000000000..5e65e20bf --- /dev/null +++ b/benchmarks/pin_memory/h2d_d2h/README.md @@ -0,0 +1,6 @@ +# H2D/D2H microbench + +Pageable vs Torch pin vs native-unregistered vs native-registered. Use +`--skip-native-register` on accelerators without `register_host_memory`. + +See the [parent README](../README.md). diff --git a/benchmarks/pin_memory/h2d_d2h/bench.py b/benchmarks/pin_memory/h2d_d2h/bench.py new file mode 100644 index 000000000..f016c513c --- /dev/null +++ b/benchmarks/pin_memory/h2d_d2h/bench.py @@ -0,0 +1,149 @@ +# SPDX-License-Identifier: Apache-2.0 +# DeepSpeed Team +"""H2D/D2H bandwidth: pageable vs Torch pin vs native unregistered vs native registered.""" + +from __future__ import annotations + +import argparse +import json +import os +import subprocess +import sys + +import torch + +from deepspeed.accelerator import get_accelerator + +ARMS = { + "pageable": { + "DS_PIN_MEMORY_BACKEND": "torch", + "DS_PIN_MEMORY_REGISTER_DEVICE": "0" + }, + "torch": { + "DS_PIN_MEMORY_BACKEND": "torch", + "DS_PIN_MEMORY_REGISTER_DEVICE": "1" + }, + "native-unregistered": { + "DS_PIN_MEMORY_BACKEND": "native", + "DS_PIN_MEMORY_REGISTER_DEVICE": "0" + }, + "native-registered": { + "DS_PIN_MEMORY_BACKEND": "native", + "DS_PIN_MEMORY_REGISTER_DEVICE": "1" + }, +} + + +def _parse_args(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--arm", choices=list(ARMS)) + parser.add_argument("--sizes-mib", type=int, nargs="+", default=[4, 64, 256]) + parser.add_argument("--warmup", type=int, default=10) + parser.add_argument("--iters", type=int, default=50) + parser.add_argument("--skip-native-register", + action="store_true", + help="Skip native-registered (XPU / no register_host_memory).") + return parser.parse_args() + + +def _time_copy(accelerator, copy_fn, stream, warmup, iters): + with accelerator.stream(stream): + for _ in range(warmup): + copy_fn() + stream.synchronize() + start = accelerator.Event(enable_timing=True) + end = accelerator.Event(enable_timing=True) + start.record(stream) + for _ in range(iters): + copy_fn() + end.record(stream) + stream.synchronize() + return start.elapsed_time(end) / 1000.0 / iters + + +def _allocate_host(accelerator, numel, arm): + raw = torch.empty(numel, dtype=torch.float32) + if arm == "pageable": + return raw + if arm == "torch": + return accelerator._torch_pin_memory(raw) + return accelerator.pin_memory(raw, make_copy=False) + + +def _run_arm(args): + for key, value in ARMS[args.arm].items(): + os.environ[key] = value + + accelerator = get_accelerator() + if not accelerator.is_available(): + raise RuntimeError(f"No {accelerator.device_name()} device is available") + accelerator.set_device(0) + stream = accelerator.Stream() + + for size_mib in args.sizes_mib: + num_bytes = size_mib * 1024 * 1024 + numel = num_bytes // torch.tensor([], dtype=torch.float32).element_size() + host = _allocate_host(accelerator, numel, args.arm) + device = torch.empty_like(host, device=accelerator.current_device_name()) + + h2d_seconds = _time_copy(accelerator, lambda: device.copy_(host, non_blocking=True), stream, args.warmup, + args.iters) + d2h_seconds = _time_copy(accelerator, lambda: host.copy_(device, non_blocking=True), stream, args.warmup, + args.iters) + + result = { + "experiment": "h2d_d2h", + "arm": args.arm, + "size_mib": size_mib, + "h2d_gbps": num_bytes / h2d_seconds / 1e9, + "d2h_gbps": num_bytes / d2h_seconds / 1e9, + "torch_is_pinned": bool(accelerator._torch_is_pinned(host)) if args.arm != "pageable" else False, + "accelerator_is_pinned": bool(accelerator.is_pinned(host)) if args.arm != "pageable" else False, + } + print(f"RESULT={json.dumps(result, sort_keys=True)}", flush=True) + if args.arm != "pageable": + accelerator.unpin_memory(host) + + +def _arms_to_run(args): + arms = list(ARMS) + if args.skip_native_register: + arms = [arm for arm in arms if arm != "native-registered"] + return arms + + +def _run_all(args): + results = [] + for arm in _arms_to_run(args): + command = [ + sys.executable, + os.path.abspath(__file__), + "--arm", + arm, + "--sizes-mib", + *(str(size) for size in args.sizes_mib), + "--warmup", + str(args.warmup), + "--iters", + str(args.iters), + ] + process = subprocess.run(command, check=True, text=True, capture_output=True) + if process.stderr: + print(process.stderr, file=sys.stderr, end="") + for line in process.stdout.splitlines(): + print(line) + if line.startswith("RESULT="): + results.append(json.loads(line.removeprefix("RESULT="))) + + print("\narm,size_mib,h2d_gbps,d2h_gbps,torch_is_pinned,accelerator_is_pinned") + for result in results: + print(f"{result['arm']},{result['size_mib']},{result['h2d_gbps']:.2f},{result['d2h_gbps']:.2f}," + f"{result['torch_is_pinned']},{result['accelerator_is_pinned']}") + + +if __name__ == "__main__": + arguments = _parse_args() + if arguments.arm: + _run_arm(arguments) + else: + _run_all(arguments) diff --git a/benchmarks/pin_memory/model_tensor_offload/README.md b/benchmarks/pin_memory/model_tensor_offload/README.md new file mode 100644 index 000000000..563236ede --- /dev/null +++ b/benchmarks/pin_memory/model_tensor_offload/README.md @@ -0,0 +1,7 @@ +# Model-tensor CPU offload + +Pin vs pageable step time for ZeRO CPU offload of parameters and/or optimizer +states. Default `--zero-stage 3` (param + optimizer). `--zero-stage 1` or `2` +offloads the optimizer only. + +See the [parent README](../README.md). diff --git a/benchmarks/pin_memory/model_tensor_offload/bench.py b/benchmarks/pin_memory/model_tensor_offload/bench.py new file mode 100644 index 000000000..b2383dc36 --- /dev/null +++ b/benchmarks/pin_memory/model_tensor_offload/bench.py @@ -0,0 +1,225 @@ +# SPDX-License-Identifier: Apache-2.0 +# DeepSpeed Team +""" +Model-tensor CPU offload: pinned vs unpinned host memory (ZeRO stage 1/2/3). + +Default is ZeRO-3 with offload_optimizer + offload_param. Stages 1 and 2 +offload the optimizer only. Pinned arms use DS_PIN_MEMORY_BACKEND=native. +Pass --ablate-register for registered vs unregistered (CUDA). +""" + +from __future__ import annotations + +import argparse +import json +import os +import sys +import time + +_PIN_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +if _PIN_ROOT not in sys.path: + sys.path.insert(0, _PIN_ROOT) + +from common import dist_env, patch_cpp_extension_drop_cxx17, print_arm_result, print_driver_result, run_arm_subprocess + + +def parse_args(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--model", + type=str, + default=None, + help="HF model id; architecture from config, random weights. Omit for synthetic.") + parser.add_argument("--hidden", type=int, default=2048) + parser.add_argument("--layers", type=int, default=12) + parser.add_argument("--batch", type=int, default=4) + parser.add_argument("--seq", type=int, default=128) + parser.add_argument("--steps", type=int, default=4) + parser.add_argument("--warmup", type=int, default=2) + parser.add_argument("--zero-stage", type=int, choices=(1, 2, 3), default=3) + parser.add_argument("--pin", type=int, default=None, help="internal: offload pin_memory 0/1") + parser.add_argument("--register", type=int, default=None, help="internal: DS_PIN_MEMORY_REGISTER_DEVICE") + parser.add_argument("--ablate-register", action="store_true") + return parser.parse_args() + + +def _build_synthetic(hidden, layers): + import torch + + class Block(torch.nn.Module): + + def __init__(self, h): + super().__init__() + self.fc1 = torch.nn.Linear(h, 4 * h) + self.fc2 = torch.nn.Linear(4 * h, h) + + def forward(self, x): + return self.fc2(torch.nn.functional.gelu(self.fc1(x))) + + class Net(torch.nn.Module): + + def __init__(self, h, n): + super().__init__() + self.emb = torch.nn.Embedding(32000, h) + self.blocks = torch.nn.ModuleList([Block(h) for _ in range(n)]) + self.head = torch.nn.Linear(h, 32000, bias=False) + + def forward(self, idx): + x = self.emb(idx) + for block in self.blocks: + x = x + block(x) + return self.head(x).sum() + + return Net(hidden, layers), 32000 + + +def _zero_config(stage, pin): + offload_optimizer = {"device": "cpu", "pin_memory": bool(pin)} + cfg = {"stage": stage, "offload_optimizer": offload_optimizer} + if stage == 3: + cfg["offload_param"] = {"device": "cpu", "pin_memory": bool(pin)} + return cfg + + +def run_arm(args): + os.environ["DS_PIN_MEMORY_BACKEND"] = "native" + os.environ["DS_PIN_MEMORY_REGISTER_DEVICE"] = str(args.register if args.pin else 0) + dist_env() + patch_cpp_extension_drop_cxx17() + + import torch + import deepspeed + + if args.model: + from transformers import AutoConfig, AutoModelForCausalLM + config = AutoConfig.from_pretrained(args.model) + model = AutoModelForCausalLM.from_config(config) + vocab = model.get_input_embeddings().weight.shape[0] + + def forward_loss(engine, ids): + return engine(ids).logits.sum() + else: + model, vocab = _build_synthetic(args.hidden, args.layers) + + def forward_loss(engine, ids): + return engine(ids) + + n_params = sum(p.numel() for p in model.parameters()) + ds_config = { + "train_micro_batch_size_per_gpu": args.batch, + "gradient_accumulation_steps": 1, + "bf16": { + "enabled": True + }, + "optimizer": { + "type": "AdamW", + "params": { + "lr": 1e-4 + } + }, + "zero_optimization": _zero_config(args.zero_stage, args.pin), + } + + engine, _, _, _ = deepspeed.initialize(model=model, config=ds_config) + dev = engine.device + dev_api = getattr(torch, dev.type) + ids = torch.randint(0, vocab, (args.batch, args.seq), dtype=torch.long, device=dev) + + step_times = [] + for step in range(args.warmup + args.steps): + dev_api.synchronize() + t0 = time.perf_counter() + loss = forward_loss(engine, ids) + engine.backward(loss) + engine.step() + dev_api.synchronize() + if step >= args.warmup: + step_times.append(time.perf_counter() - t0) + + def _mem(fn): + try: + return round(fn() / 1e9, 2) + except Exception: + return None + + print_arm_result({ + "experiment": "model_tensor_offload", + "zero_stage": args.zero_stage, + "pin_memory": bool(args.pin), + "register": bool(args.register) if args.pin else None, + "device": dev.type, + "model": args.model or f"synthetic-h{args.hidden}-l{args.layers}", + "params_b": round(n_params / 1e9, 3), + "batch": args.batch, + "seq": args.seq, + "tokens_per_step": args.batch * args.seq, + "steps": len(step_times), + "step_avg_s": sum(step_times) / len(step_times), + "step_min_s": min(step_times), + "gpu_peak_gb": _mem(dev_api.max_memory_allocated), + }) + + +def _passthrough(args): + extra = [ + "--zero-stage", + str(args.zero_stage), + "--hidden", + str(args.hidden), + "--layers", + str(args.layers), + "--batch", + str(args.batch), + "--seq", + str(args.seq), + "--steps", + str(args.steps), + "--warmup", + str(args.warmup), + ] + if args.model: + extra.extend(["--model", args.model]) + return extra + + +def run_driver(args): + if args.ablate_register: + arms = [("unpinned", 0, 0), ("pinned-unregistered", 1, 0), ("pinned-registered", 1, 1)] + else: + arms = [("unpinned", 0, 0), ("pinned", 1, 1)] + script = os.path.abspath(__file__) + results = {} + for name, pin, reg in arms: + extra = _passthrough(args) + ["--pin", str(pin), "--register", str(reg)] + results[name] = run_arm_subprocess(script, extra) + + unpinned = results["unpinned"] + pinned = results["pinned-registered"] if args.ablate_register else results["pinned"] + print("\n================ Model-tensor CPU-offload step time ================") + print(f"stage: {pinned['zero_stage']} model: {pinned['model']} params: {pinned['params_b']}B " + f"device: {pinned['device']}") + print(f"batch: {pinned['batch']} x seq: {pinned['seq']} ({pinned['tokens_per_step']} tokens/step)") + print() + print(f"{'arm':<22}{'avg step (s)':>14}{'min step (s)':>14}{'GPU peak (GB)':>16}") + for name, _, _ in arms: + row = results[name] + print(f"{name:<22}{row['step_avg_s']:>14.3f}{row['step_min_s']:>14.3f}" + f"{(row['gpu_peak_gb'] or 0):>16.2f}") + speedup = unpinned["step_avg_s"] / pinned["step_avg_s"] + saved = unpinned["step_avg_s"] - pinned["step_avg_s"] + tok_pin = pinned["tokens_per_step"] / pinned["step_avg_s"] + tok_unpin = unpinned["tokens_per_step"] / unpinned["step_avg_s"] + print() + print(f"pinning speedup: {speedup:.2f}x saved: {saved:.3f} s/step " + f"throughput: {tok_pin:.0f} vs {tok_unpin:.0f} tok/s") + summary = {"unpinned": unpinned, "pinned": pinned, "speedup": speedup} + if args.ablate_register: + summary["pinned-unregistered"] = results["pinned-unregistered"] + print_driver_result(summary) + + +if __name__ == "__main__": + parsed = parse_args() + if parsed.pin is None: + run_driver(parsed) + else: + run_arm(parsed) diff --git a/benchmarks/pin_memory/zero3_offload_bench.py b/benchmarks/pin_memory/zero3_offload_bench.py index ab35f9a04..12857859a 100644 --- a/benchmarks/pin_memory/zero3_offload_bench.py +++ b/benchmarks/pin_memory/zero3_offload_bench.py @@ -1,268 +1,11 @@ # SPDX-License-Identifier: Apache-2.0 # DeepSpeed Team -""" -ZeRO-3 CPU-offload end-to-end benchmark: pinned vs unpinned host memory. +"""Backward-compatible entry: model-tensor ZeRO-3 offload (see model_tensor_offload/).""" -Measures training step time with ZeRO-3 and offload_optimizer + offload_param -(both cpu), comparing pin_memory=True (native backend, registered with the -device) against an unpinned baseline. Pass --ablate-register to additionally -compare registered vs unregistered pinned memory. Works on CUDA and XPU (any -accelerator with native pin + register_host_memory support). +from __future__ import annotations -Each arm runs in its own subprocess with a fresh rendezvous port so device -state never leaks between arms. - -Examples: - # real model architecture (config fetched from the HF hub, random weights) - python zero3_offload_bench.py --model Qwen/Qwen2.5-7B --batch 4 --seq 512 - - # synthetic MLP stack, no network needed - python zero3_offload_bench.py --hidden 2048 --layers 12 --batch 4 --seq 128 - -Note: the native backend mlocks host memory; raise RLIMIT_MEMLOCK -(ulimit -l) or run as root for multi-GB models. -""" - -import argparse -import json import os -import socket -import subprocess import sys -import time - - -def str2bool(v): - return str(v).lower() in ("1", "true", "yes", "on") - - -def parse_args(): - p = argparse.ArgumentParser(description=__doc__) - p.add_argument("--model", - type=str, - default=None, - help="HF model id (e.g. Qwen/Qwen2.5-7B); architecture is loaded from " - "the config with random weights. Omit to use the synthetic model.") - p.add_argument("--hidden", type=int, default=2048, help="synthetic model hidden size") - p.add_argument("--layers", type=int, default=12, help="synthetic model layer count") - p.add_argument("--batch", type=int, default=4, help="micro batch size per gpu") - p.add_argument("--seq", type=int, default=128, help="sequence length") - p.add_argument("--steps", type=int, default=4, help="timed steps per arm") - p.add_argument("--warmup", type=int, default=2, help="warmup steps per arm") - p.add_argument("--pin", type=int, default=None, help="internal: run a single arm with offload pin_memory=0/1") - p.add_argument("--register", - type=int, - default=None, - help="internal: run a single pinned arm with DS_PIN_MEMORY_REGISTER_DEVICE=0/1") - p.add_argument("--ablate-register", - action="store_true", - help="also report registered vs unregistered pinned memory (for power users)") - return p.parse_args() - - -def _free_port(): - # A stale listener from an interrupted rank makes the next init hang in a - # collective, so always rendezvous on a fresh ephemeral port. - s = socket.socket() - s.bind(("127.0.0.1", 0)) - port = s.getsockname()[1] - s.close() - return port - - -def _build_synthetic(hidden, layers): - import torch - - class Block(torch.nn.Module): - - def __init__(self, h): - super().__init__() - self.fc1 = torch.nn.Linear(h, 4 * h) - self.fc2 = torch.nn.Linear(4 * h, h) - - def forward(self, x): - return self.fc2(torch.nn.functional.gelu(self.fc1(x))) - - class Net(torch.nn.Module): - - def __init__(self, h, n): - super().__init__() - self.emb = torch.nn.Embedding(32000, h) - self.blocks = torch.nn.ModuleList([Block(h) for _ in range(n)]) - self.head = torch.nn.Linear(h, 32000, bias=False) - - def forward(self, idx): - x = self.emb(idx) - for b in self.blocks: - x = x + b(x) - return self.head(x).sum() - - return Net(hidden, layers), 32000 - - -def run_arm(args): - """Single arm: one process, one pinning/register setting.""" - os.environ["DS_PIN_MEMORY_BACKEND"] = "native" - # Registering only matters once memory is pinned; keep it off otherwise. - os.environ["DS_PIN_MEMORY_REGISTER_DEVICE"] = str(args.register if args.pin else 0) - os.environ.update(MASTER_ADDR="127.0.0.1", MASTER_PORT=str(_free_port()), RANK="0", WORLD_SIZE="1", LOCAL_RANK="0") - - import torch - import torch.utils.cpp_extension as _ce - - # torch-nightly requires C++20; SYCL toolchain flags may carry -std=c++17, - # which (appearing last) downgrades the dialect and breaks torch headers. - _orig_ce_load = _ce.load - - def _ce_load_without_cxx17(*a, **kw): - for key in ("extra_cflags", "extra_cxxflags"): - if kw.get(key): - kw[key] = [f for f in kw[key] if f != "-std=c++17"] - return _orig_ce_load(*a, **kw) - - _ce.load = _ce_load_without_cxx17 - - import deepspeed - - if args.model: - from transformers import AutoConfig, AutoModelForCausalLM - config = AutoConfig.from_pretrained(args.model) - model = AutoModelForCausalLM.from_config(config) - # Newer config classes may nest vocab_size, so read it off the model. - vocab = model.get_input_embeddings().weight.shape[0] - - def forward_loss(engine, ids): - return engine(ids).logits.sum() - else: - model, vocab = _build_synthetic(args.hidden, args.layers) - - def forward_loss(engine, ids): - return engine(ids) - - n_params = sum(p.numel() for p in model.parameters()) - - ds_config = { - "train_micro_batch_size_per_gpu": args.batch, - "gradient_accumulation_steps": 1, - "bf16": { - "enabled": True - }, - "optimizer": { - "type": "AdamW", - "params": { - "lr": 1e-4 - } - }, - "zero_optimization": { - "stage": 3, - "offload_optimizer": { - "device": "cpu", - "pin_memory": bool(args.pin) - }, - "offload_param": { - "device": "cpu", - "pin_memory": bool(args.pin) - }, - }, - } - - engine, _, _, _ = deepspeed.initialize(model=model, config=ds_config) - dev = engine.device - # Device API via getattr so the accelerator-agnostic rule (no hardcoded - # backend namespaces) holds while working on cuda and xpu alike. - dev_api = getattr(torch, dev.type) - sync = dev_api.synchronize - ids = torch.randint(0, vocab, (args.batch, args.seq), dtype=torch.long, device=dev) - - step_times = [] - for step in range(args.warmup + args.steps): - sync() - t0 = time.perf_counter() - loss = forward_loss(engine, ids) - engine.backward(loss) - engine.step() - sync() - dt = time.perf_counter() - t0 - if step >= args.warmup: - step_times.append(dt) - - def _mem(fn): - try: - return round(fn() / 1e9, 2) - except Exception: - return None - - result = { - "pin_memory": bool(args.pin), - "register": bool(args.register) if args.pin else None, - "device": dev.type, - "model": args.model or f"synthetic-h{args.hidden}-l{args.layers}", - "params_b": round(n_params / 1e9, 3), - "batch": args.batch, - "seq": args.seq, - "tokens_per_step": args.batch * args.seq, - "steps": len(step_times), - "step_avg_s": sum(step_times) / len(step_times), - "step_min_s": min(step_times), - "gpu_peak_gb": _mem(dev_api.max_memory_allocated), - } - print("ARMRESULT=" + json.dumps(result), flush=True) - - -def run_driver(args): - """Run the arms in subprocesses and print the comparison.""" - if args.ablate_register: - arms = [("unpinned", 0, 0), ("pinned-unregistered", 1, 0), ("pinned-registered", 1, 1)] - else: - arms = [("unpinned", 0, 0), ("pinned", 1, 1)] - results = {} - for name, pin, reg in arms: - cmd = [sys.executable, os.path.abspath(__file__), "--pin", str(pin), "--register", str(reg), "--model"] + \ - ([args.model] if args.model else ["None"]) + \ - ["--hidden", str(args.hidden), "--layers", str(args.layers), "--batch", str(args.batch), - "--seq", str(args.seq), "--steps", str(args.steps), "--warmup", str(args.warmup)] - # argparse cannot take a literal None for --model; drop it instead. - if not args.model: - cmd = cmd[:cmd.index("--model")] + cmd[cmd.index("--model") + 2:] - print(f"[driver] running arm {name} ...", flush=True) - proc = subprocess.run(cmd, env=os.environ.copy(), capture_output=True, text=True) - arm = None - for line in proc.stdout.splitlines(): - if line.startswith("ARMRESULT="): - arm = json.loads(line[len("ARMRESULT="):]) - if arm is None: - print(proc.stdout[-2000:]) - print(proc.stderr[-2000:], file=sys.stderr) - raise RuntimeError(f"arm {name} produced no result (rc={proc.returncode})") - results[name] = arm - - unpinned = results["unpinned"] - pinned = results["pinned-registered"] if args.ablate_register else results["pinned"] - print("\n================ ZeRO-3 CPU-offload step time ================") - print(f"model: {pinned['model']} params: {pinned['params_b']}B device: {pinned['device']}") - print(f"batch: {pinned['batch']} x seq: {pinned['seq']} ({pinned['tokens_per_step']} tokens/step)") - print() - print(f"{'arm':<22}{'avg step (s)':>14}{'min step (s)':>14}{'GPU peak (GB)':>16}") - for name, _, _ in arms: - r = results[name] - print(f"{name:<22}{r['step_avg_s']:>14.3f}{r['step_min_s']:>14.3f}" - f"{(r['gpu_peak_gb'] or 0):>16.2f}") - speedup = unpinned['step_avg_s'] / pinned['step_avg_s'] - saved = unpinned['step_avg_s'] - pinned['step_avg_s'] - tok_pin = pinned['tokens_per_step'] / pinned['step_avg_s'] - tok_unpin = unpinned['tokens_per_step'] / unpinned['step_avg_s'] - print() - print(f"pinning speedup: {speedup:.2f}x saved: {saved:.3f} s/step " - f"throughput: {tok_pin:.0f} vs {tok_unpin:.0f} tok/s") - summary = {'unpinned': unpinned, 'pinned': pinned, 'speedup': speedup} - if args.ablate_register: - summary['pinned-unregistered'] = results['pinned-unregistered'] - print(f"DRIVERRESULT={json.dumps(summary)}", flush=True) - -if __name__ == "__main__": - _args = parse_args() - if _args.pin is None: - run_driver(_args) - else: - run_arm(_args) +_NEW = os.path.join(os.path.dirname(os.path.abspath(__file__)), "model_tensor_offload", "bench.py") +os.execv(sys.executable, [sys.executable, _NEW, "--zero-stage", "3", *sys.argv[1:]])