From 45c74e6af4b7a31e79a6ee250af0f045da224e36 Mon Sep 17 00:00:00 2001 From: nathon-lee Date: Sun, 23 Aug 2026 20:06:28 +0800 Subject: [PATCH 1/2] perf(opsd): add HybridEngine rollout profiling benchmark Signed-off-by: nathon-lee --- benchmarks/opsd/README.md | 52 ++++ .../opsd/benchmark_hybrid_engine_rollout.py | 229 ++++++++++++++++++ .../test_benchmark_hybrid_engine_rollout.py | 110 +++++++++ 3 files changed, 391 insertions(+) create mode 100644 benchmarks/opsd/README.md create mode 100644 benchmarks/opsd/benchmark_hybrid_engine_rollout.py create mode 100644 benchmarks/opsd/tests/test_benchmark_hybrid_engine_rollout.py diff --git a/benchmarks/opsd/README.md b/benchmarks/opsd/README.md new file mode 100644 index 000000000..9e07fa752 --- /dev/null +++ b/benchmarks/opsd/README.md @@ -0,0 +1,52 @@ +# OPSD HybridEngine rollout benchmark + +`benchmark_hybrid_engine_rollout.py` measures rollout-level performance for an +OPSD workload backed by DeepSpeed HybridEngine. It runs a matrix of synthetic, +exact-length prompts and reports prompt expansion, generation, post-processing, +total latency, generated-token throughput, and peak accelerator memory. Each +case includes raw iteration profiles plus mean, p50, and p95 summaries. + +This benchmark depends on the rollout profiling API introduced by +DeepSpeed PR #8295: + +https://github.com/deepspeedai/DeepSpeed/pull/8295 + +Use a DeepSpeed checkout that contains that API and place it first on +`PYTHONPATH`. The current validation scope is one process, one GPU, and ZeRO-0. +This is a HybridEngine rollout benchmark, not a complete OPSD training-step +benchmark; it does not measure teacher inference, loss computation, backward, +or optimizer work. + +## Usage + +The workload matrix is controlled by `--batch-sizes`, +`--samples-per-prompt`, `--prompt-lengths`, and `--response-lengths`. Both +`--dtype fp16` and `--dtype bf16` are supported. `--warmup` and `--iterations` +control unreported warmup calls and recorded calls. Pass +`--release-inference-cache` to release the inference cache after generation; +otherwise the benchmark retains it. `--temperature`, `--top-p`, `--seed`, and +`--output` control sampling, reproducibility, and the JSON output path. + +The largest effective batch (`batch_size * samples_per_prompt`) executes first +so HybridEngine initializes a sufficiently large inference workspace. Results +in the output JSON retain the matrix order requested on the command line. + +From the DeepSpeedExamples repository root, run a single-GPU benchmark with: + +```bash +PYTHONPATH=/workspace/DeepSpeed_woo:/workspace/DeepSpeedExamples \ +torchrun --nproc_per_node=1 \ + benchmarks/opsd/benchmark_hybrid_engine_rollout.py \ + --model facebook/opt-6.7b \ + --batch-sizes 1 \ + --samples-per-prompt 1 4 \ + --prompt-lengths 128 512 \ + --response-lengths 32 128 \ + --warmup 1 \ + --iterations 2 \ + --output /tmp/opsd_rollout_profile_examples.json +``` + +Use `--help` for the complete argument list. Model download, GPU memory, and +the fused inference kernels supported by the selected model can limit which +matrix shapes run successfully. diff --git a/benchmarks/opsd/benchmark_hybrid_engine_rollout.py b/benchmarks/opsd/benchmark_hybrid_engine_rollout.py new file mode 100644 index 000000000..d8f031862 --- /dev/null +++ b/benchmarks/opsd/benchmark_hybrid_engine_rollout.py @@ -0,0 +1,229 @@ +# SPDX-License-Identifier: Apache-2.0 +# DeepSpeed Team +"""Benchmark stage-level profiling for HybridEngine-backed OPSD rollout. + +The benchmark uses synthetic token IDs so prompt lengths are exact and model +tokenization does not become part of the measured rollout. Run it with one +accelerator process; ZeRO-3 and multi-rank measurements are intentionally left +for a later benchmark stage. +""" + +import argparse +import itertools +import json +import math +import os +import statistics +from pathlib import Path + + +_TIMING_FIELDS = ("prompt_expansion_ms", "generation_ms", "post_processing_ms", "total_ms", "tokens_per_second") + + +def _percentile(values, percentile): + ordered = sorted(values) + index = max(0, math.ceil(len(ordered) * percentile) - 1) + return ordered[index] + + +def _summarize(profiles): + summary = {} + for field in _TIMING_FIELDS: + values = [profile[field] for profile in profiles] + summary[field] = { + "mean": statistics.mean(values), + "p50": statistics.median(values), + "p95": _percentile(values, 0.95), + } + return summary + + +def _build_parser(): + parser = argparse.ArgumentParser(description="Benchmark HybridEngine rollout stage profiling") + parser.add_argument("--model", default="facebook/opt-6.7b", help="HuggingFace model ID or local model path") + parser.add_argument("--dtype", choices=["fp16", "bf16"], default="fp16") + parser.add_argument("--batch-sizes", type=int, nargs="+", default=[1]) + parser.add_argument("--samples-per-prompt", type=int, nargs="+", default=[1, 4]) + parser.add_argument("--prompt-lengths", type=int, nargs="+", default=[128, 512]) + parser.add_argument("--response-lengths", type=int, nargs="+", default=[32, 128]) + parser.add_argument("--temperature", type=float, default=0.0) + parser.add_argument("--top-p", type=float, default=1.0) + parser.add_argument("--warmup", type=int, default=5) + parser.add_argument("--iterations", type=int, default=20) + parser.add_argument("--release-inference-cache", action="store_true") + parser.add_argument("--seed", type=int, default=1234) + parser.add_argument("--output", default="opsd_rollout_profile.json") + return parser + + +def _validate_args(args): + positive_values = [*args.batch_sizes, *args.samples_per_prompt, *args.prompt_lengths, *args.response_lengths] + positive_values.extend([args.warmup, args.iterations]) + if any(value <= 0 for value in positive_values): + raise ValueError("Batch sizes, sequence lengths, warmup, and iterations must all be positive") + if args.temperature < 0.0: + raise ValueError("temperature must be non-negative") + if not 0.0 < args.top_p <= 1.0: + raise ValueError("top-p must be in the interval (0, 1]") + + +def _load_model_and_tokenizer(model_name, dtype, device): + from transformers import AutoModelForCausalLM, AutoTokenizer + + tokenizer = AutoTokenizer.from_pretrained(model_name) + if tokenizer.pad_token_id is None: + if tokenizer.eos_token_id is None: + raise ValueError("Tokenizer must define pad_token_id or eos_token_id") + tokenizer.pad_token = tokenizer.eos_token + + model = AutoModelForCausalLM.from_pretrained(model_name, torch_dtype=dtype, low_cpu_mem_usage=True) + return model.to(device), tokenizer + + +def _build_engine(model, args): + import deepspeed + + use_bf16 = args.dtype == "bf16" + max_effective_batch_size = max(args.batch_sizes) * max(args.samples_per_prompt) + ds_config = { + "train_batch_size": max_effective_batch_size, + "train_micro_batch_size_per_gpu": max_effective_batch_size, + "fp16": { + "enabled": not use_bf16, + }, + "bf16": { + "enabled": use_bf16, + }, + "zero_optimization": { + "stage": 0, + }, + "hybrid_engine": { + "enabled": True, + "max_out_tokens": max(args.prompt_lengths) + max(args.response_lengths), + "release_inference_cache": args.release_inference_cache, + }, + } + engine, _, _, _ = deepspeed.initialize(model=model, config=ds_config) + if not hasattr(engine, "_generate"): + raise RuntimeError("The model architecture did not create HybridEngine inference containers") + engine.eval() + return engine + + +def _make_request(model, batch_size, prompt_length, device): + import torch + from deepspeed.runtime.rollout.base import RolloutRequest + + vocab_size = model.config.vocab_size + first_token_id = min(3, vocab_size - 1) + prompt_ids = torch.randint(first_token_id, vocab_size, (batch_size, prompt_length), device=device) + prompt_attention_mask = torch.ones_like(prompt_ids) + return RolloutRequest(prompt_ids=prompt_ids, prompt_attention_mask=prompt_attention_mask) + + +def _run_case(rollout, model, args, batch_size, samples_per_prompt, prompt_length, response_length, device): + from deepspeed.accelerator import get_accelerator + from deepspeed.runtime.rollout.base import SamplingConfig + + request = _make_request(model, batch_size, prompt_length, device) + sampling = SamplingConfig( + max_new_tokens=response_length, + temperature=args.temperature, + top_p=args.top_p, + n_samples_per_prompt=samples_per_prompt, + ) + + for _ in range(args.warmup): + rollout.generate(request, sampling) + + accelerator = get_accelerator() + accelerator.reset_peak_memory_stats() + profiles = [] + for _ in range(args.iterations): + rollout.generate(request, sampling) + profiles.append(dict(rollout.get_last_profile())) + + return { + "batch_size": batch_size, + "samples_per_prompt": samples_per_prompt, + "prompt_length": prompt_length, + "requested_response_length": response_length, + "returned_response_length": profiles[-1]["response_length"], + "peak_memory_mb": accelerator.max_memory_allocated() / (1024**2), + "summary": _summarize(profiles), + "profiles": profiles, + } + + +def _ordered_case_specs(args): + requested = list( + itertools.product(args.batch_sizes, args.samples_per_prompt, args.prompt_lengths, args.response_lengths)) + execution_order = sorted(range(len(requested)), + key=lambda index: requested[index][0] * requested[index][1], + reverse=True) + return requested, execution_order + + +def _build_result(args, device, cases): + return { + "model": args.model, + "dtype": args.dtype, + "device": device, + "warmup": args.warmup, + "iterations": args.iterations, + "temperature": args.temperature, + "top_p": args.top_p, + "release_inference_cache": args.release_inference_cache, + "cases": cases, + } + + +def _run(args): + import torch + import deepspeed.comm as dist + from deepspeed.accelerator import get_accelerator + from deepspeed.runtime.rollout.hybrid_engine_rollout import HybridEngineRollout, HybridEngineRolloutConfig + + _validate_args(args) + world_size = int(os.getenv("WORLD_SIZE", "1")) + if world_size != 1: + raise RuntimeError("This initial benchmark supports exactly one accelerator process") + + torch.manual_seed(args.seed) + local_rank = int(os.getenv("LOCAL_RANK", "0")) + accelerator = get_accelerator() + accelerator.set_device(local_rank) + device = accelerator.device_name(local_rank) + dtype = torch.bfloat16 if args.dtype == "bf16" else torch.float16 + + try: + model, tokenizer = _load_model_and_tokenizer(args.model, dtype, device) + engine = _build_engine(model, args) + rollout = HybridEngineRollout(engine, tokenizer, HybridEngineRolloutConfig(enable_profiling=True)) + + case_specs, execution_order = _ordered_case_specs(args) + cases = [None] * len(case_specs) + # HybridEngine sizes its inference workspace on the first forward. Run + # the largest effective batch first while preserving requested output order. + for case_index in execution_order: + batch_size, samples_per_prompt, prompt_length, response_length = case_specs[case_index] + cases[case_index] = _run_case(rollout, engine.module, args, batch_size, samples_per_prompt, prompt_length, + response_length, device) + + result = _build_result(args, device, cases) + if not dist.is_initialized() or dist.get_rank() == 0: + output_path = Path(args.output) + output_path.parent.mkdir(parents=True, exist_ok=True) + output_path.write_text(json.dumps(result, indent=2) + "\n", encoding="utf-8") + print(f"Wrote benchmark results to {output_path}") + finally: + if dist.is_initialized(): + dist.destroy_process_group() + + +def main(): + _run(_build_parser().parse_args()) + + +if __name__ == "__main__": + main() diff --git a/benchmarks/opsd/tests/test_benchmark_hybrid_engine_rollout.py b/benchmarks/opsd/tests/test_benchmark_hybrid_engine_rollout.py new file mode 100644 index 000000000..c39b16096 --- /dev/null +++ b/benchmarks/opsd/tests/test_benchmark_hybrid_engine_rollout.py @@ -0,0 +1,110 @@ +# SPDX-License-Identifier: Apache-2.0 +# DeepSpeed Team + +import importlib.util +import json +from pathlib import Path +import unittest + + +_BENCHMARK_PATH = Path(__file__).resolve().parents[1] / "benchmark_hybrid_engine_rollout.py" +_SPEC = importlib.util.spec_from_file_location("benchmark_hybrid_engine_rollout", _BENCHMARK_PATH) +benchmark = importlib.util.module_from_spec(_SPEC) +_SPEC.loader.exec_module(benchmark) + + +class BenchmarkHybridEngineRolloutTest(unittest.TestCase): + + def test_argument_parsing(self): + args = benchmark._build_parser().parse_args([ + "--model", + "local-model", + "--dtype", + "bf16", + "--batch-sizes", + "2", + "4", + "--samples-per-prompt", + "1", + "3", + "--prompt-lengths", + "8", + "--response-lengths", + "16", + "--release-inference-cache", + ]) + + self.assertEqual(args.model, "local-model") + self.assertEqual(args.dtype, "bf16") + self.assertEqual(args.batch_sizes, [2, 4]) + self.assertEqual(args.samples_per_prompt, [1, 3]) + self.assertEqual(args.prompt_lengths, [8]) + self.assertEqual(args.response_lengths, [16]) + self.assertTrue(args.release_inference_cache) + + def test_summary_percentiles(self): + profiles = [] + for value in range(1, 21): + profiles.append({field: float(value) for field in benchmark._TIMING_FIELDS}) + + summary = benchmark._summarize(profiles) + + self.assertEqual(summary["total_ms"], {"mean": 10.5, "p50": 10.5, "p95": 19.0}) + self.assertEqual(summary["tokens_per_second"]["p95"], 19.0) + + def test_largest_effective_batch_executes_first_without_reordering_results(self): + args = benchmark._build_parser().parse_args([ + "--batch-sizes", + "1", + "2", + "--samples-per-prompt", + "1", + "4", + "--prompt-lengths", + "8", + "--response-lengths", + "8", + ]) + + requested, execution_order = benchmark._ordered_case_specs(args) + + self.assertEqual(requested, [(1, 1, 8, 8), (1, 4, 8, 8), (2, 1, 8, 8), (2, 4, 8, 8)]) + self.assertEqual([requested[index] for index in execution_order], [ + (2, 4, 8, 8), + (1, 4, 8, 8), + (2, 1, 8, 8), + (1, 1, 8, 8), + ]) + + def test_json_output_has_basic_fields(self): + args = benchmark._build_parser().parse_args(["--model", "local-model", "--iterations", "1"]) + case = { + "batch_size": 1, + "samples_per_prompt": 1, + "prompt_length": 8, + "requested_response_length": 8, + "returned_response_length": 8, + "peak_memory_mb": 12.5, + "summary": {}, + "profiles": [{"total_ms": 1.0}], + } + + result = json.loads(json.dumps(benchmark._build_result(args, "cuda:0", [case]))) + + self.assertEqual(set(result), { + "model", + "dtype", + "device", + "warmup", + "iterations", + "temperature", + "top_p", + "release_inference_cache", + "cases", + }) + self.assertEqual(result["cases"][0]["peak_memory_mb"], 12.5) + self.assertEqual(result["cases"][0]["profiles"], [{"total_ms": 1.0}]) + + +if __name__ == "__main__": + unittest.main() From efdb5abb21e9f63bd7f5971bb2c9cb673e544f29 Mon Sep 17 00:00:00 2001 From: nathon-lee Date: Mon, 24 Aug 2026 19:45:03 +0800 Subject: [PATCH 2/2] perf(opsd): address benchmark review feedback Signed-off-by: nathon-lee --- benchmarks/opsd/README.md | 3 +- .../opsd/benchmark_hybrid_engine_rollout.py | 11 +- .../test_benchmark_hybrid_engine_rollout.py | 110 ------------------ 3 files changed, 11 insertions(+), 113 deletions(-) delete mode 100644 benchmarks/opsd/tests/test_benchmark_hybrid_engine_rollout.py diff --git a/benchmarks/opsd/README.md b/benchmarks/opsd/README.md index 9e07fa752..e58be3800 100644 --- a/benchmarks/opsd/README.md +++ b/benchmarks/opsd/README.md @@ -4,7 +4,8 @@ OPSD workload backed by DeepSpeed HybridEngine. It runs a matrix of synthetic, exact-length prompts and reports prompt expansion, generation, post-processing, total latency, generated-token throughput, and peak accelerator memory. Each -case includes raw iteration profiles plus mean, p50, and p95 summaries. +case includes raw iteration profiles, mean and p50 summaries, and p95 summaries +for latency metrics. This benchmark depends on the rollout profiling API introduced by DeepSpeed PR #8295: diff --git a/benchmarks/opsd/benchmark_hybrid_engine_rollout.py b/benchmarks/opsd/benchmark_hybrid_engine_rollout.py index d8f031862..3d4b610ff 100644 --- a/benchmarks/opsd/benchmark_hybrid_engine_rollout.py +++ b/benchmarks/opsd/benchmark_hybrid_engine_rollout.py @@ -17,7 +17,8 @@ from pathlib import Path -_TIMING_FIELDS = ("prompt_expansion_ms", "generation_ms", "post_processing_ms", "total_ms", "tokens_per_second") +_LATENCY_FIELDS = ("prompt_expansion_ms", "generation_ms", "post_processing_ms", "total_ms") +_THROUGHPUT_FIELDS = ("tokens_per_second", ) def _percentile(values, percentile): @@ -28,13 +29,19 @@ def _percentile(values, percentile): def _summarize(profiles): summary = {} - for field in _TIMING_FIELDS: + for field in _LATENCY_FIELDS: values = [profile[field] for profile in profiles] summary[field] = { "mean": statistics.mean(values), "p50": statistics.median(values), "p95": _percentile(values, 0.95), } + for field in _THROUGHPUT_FIELDS: + values = [profile[field] for profile in profiles] + summary[field] = { + "mean": statistics.mean(values), + "p50": statistics.median(values), + } return summary diff --git a/benchmarks/opsd/tests/test_benchmark_hybrid_engine_rollout.py b/benchmarks/opsd/tests/test_benchmark_hybrid_engine_rollout.py deleted file mode 100644 index c39b16096..000000000 --- a/benchmarks/opsd/tests/test_benchmark_hybrid_engine_rollout.py +++ /dev/null @@ -1,110 +0,0 @@ -# SPDX-License-Identifier: Apache-2.0 -# DeepSpeed Team - -import importlib.util -import json -from pathlib import Path -import unittest - - -_BENCHMARK_PATH = Path(__file__).resolve().parents[1] / "benchmark_hybrid_engine_rollout.py" -_SPEC = importlib.util.spec_from_file_location("benchmark_hybrid_engine_rollout", _BENCHMARK_PATH) -benchmark = importlib.util.module_from_spec(_SPEC) -_SPEC.loader.exec_module(benchmark) - - -class BenchmarkHybridEngineRolloutTest(unittest.TestCase): - - def test_argument_parsing(self): - args = benchmark._build_parser().parse_args([ - "--model", - "local-model", - "--dtype", - "bf16", - "--batch-sizes", - "2", - "4", - "--samples-per-prompt", - "1", - "3", - "--prompt-lengths", - "8", - "--response-lengths", - "16", - "--release-inference-cache", - ]) - - self.assertEqual(args.model, "local-model") - self.assertEqual(args.dtype, "bf16") - self.assertEqual(args.batch_sizes, [2, 4]) - self.assertEqual(args.samples_per_prompt, [1, 3]) - self.assertEqual(args.prompt_lengths, [8]) - self.assertEqual(args.response_lengths, [16]) - self.assertTrue(args.release_inference_cache) - - def test_summary_percentiles(self): - profiles = [] - for value in range(1, 21): - profiles.append({field: float(value) for field in benchmark._TIMING_FIELDS}) - - summary = benchmark._summarize(profiles) - - self.assertEqual(summary["total_ms"], {"mean": 10.5, "p50": 10.5, "p95": 19.0}) - self.assertEqual(summary["tokens_per_second"]["p95"], 19.0) - - def test_largest_effective_batch_executes_first_without_reordering_results(self): - args = benchmark._build_parser().parse_args([ - "--batch-sizes", - "1", - "2", - "--samples-per-prompt", - "1", - "4", - "--prompt-lengths", - "8", - "--response-lengths", - "8", - ]) - - requested, execution_order = benchmark._ordered_case_specs(args) - - self.assertEqual(requested, [(1, 1, 8, 8), (1, 4, 8, 8), (2, 1, 8, 8), (2, 4, 8, 8)]) - self.assertEqual([requested[index] for index in execution_order], [ - (2, 4, 8, 8), - (1, 4, 8, 8), - (2, 1, 8, 8), - (1, 1, 8, 8), - ]) - - def test_json_output_has_basic_fields(self): - args = benchmark._build_parser().parse_args(["--model", "local-model", "--iterations", "1"]) - case = { - "batch_size": 1, - "samples_per_prompt": 1, - "prompt_length": 8, - "requested_response_length": 8, - "returned_response_length": 8, - "peak_memory_mb": 12.5, - "summary": {}, - "profiles": [{"total_ms": 1.0}], - } - - result = json.loads(json.dumps(benchmark._build_result(args, "cuda:0", [case]))) - - self.assertEqual(set(result), { - "model", - "dtype", - "device", - "warmup", - "iterations", - "temperature", - "top_p", - "release_inference_cache", - "cases", - }) - self.assertEqual(result["cases"][0]["peak_memory_mb"], 12.5) - self.assertEqual(result["cases"][0]["profiles"], [{"total_ms": 1.0}]) - - -if __name__ == "__main__": - unittest.main()