Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 53 additions & 0 deletions benchmarks/opsd/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
# 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, mean and p50 summaries, and p95 summaries
for latency metrics.

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.
236 changes: 236 additions & 0 deletions benchmarks/opsd/benchmark_hybrid_engine_rollout.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,236 @@
# 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


_LATENCY_FIELDS = ("prompt_expansion_ms", "generation_ms", "post_processing_ms", "total_ms")
_THROUGHPUT_FIELDS = ("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 _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


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()
Loading