diff --git a/experiments/dsv4_h20_sglang_tp_dp_matrix/adaptive_config.env b/experiments/dsv4_h20_sglang_tp_dp_matrix/adaptive_config.env new file mode 100644 index 0000000..1eda1f0 --- /dev/null +++ b/experiments/dsv4_h20_sglang_tp_dp_matrix/adaptive_config.env @@ -0,0 +1,49 @@ +# Adaptive concurrency search settings. +# +# For each fixed (TP, DP, ISL, OSL), probe: +# C = start, start * multiplier, ... up to max +# and stop after Total TPS has less than TPS_MIN_GAIN_PCT meaningful growth for +# PLATEAU_PATIENCE consecutive points. + +SEARCH_START_CONCURRENCY="${SEARCH_START_CONCURRENCY:-1}" +SEARCH_MAX_CONCURRENCY="${SEARCH_MAX_CONCURRENCY:-512}" +SEARCH_MULTIPLIER="${SEARCH_MULTIPLIER:-2}" +NUM_PROMPTS_MULTIPLIER="${NUM_PROMPTS_MULTIPLIER:-5}" + +# A gain below 2% is treated as throughput saturation. Two consecutive +# low-gain points prevent one noisy measurement from stopping the search. +TPS_MIN_GAIN_PCT="${TPS_MIN_GAIN_PCT:-2.0}" +PLATEAU_PATIENCE="${PLATEAU_PATIENCE:-2}" + +# Keep the same random workload semantics as the fixed matrix baseline. +# DATASET_PATH must contain at least SEARCH_MAX_CONCURRENCY times +# NUM_PROMPTS_MULTIPLIER valid two-turn conversations. Set this explicitly to +# random-ids to use generated token IDs without a ShareGPT seed dataset. +BENCH_DATASET_NAME="${BENCH_DATASET_NAME:-random}" +# SGLang interprets 0.0 as Uniform[1, requested_len]. Use 1.0 for fixed +# ISL/OSL points; lower values intentionally benchmark a length distribution. +RANDOM_RANGE_RATIO="${RANDOM_RANGE_RATIO:-1.0}" +# Before each measured point, warm up with the same concurrency so lazy kernel +# compilation and CUDA graph capture are excluded from TTFT/TPS. 0 means no +# cap; set a positive cap only when very high-concurrency warmup is impractical. +BENCH_WARMUP_MAX_REQUESTS="${BENCH_WARMUP_MAX_REQUESTS:-0}" + +# Reject a point if the completed request count or actual token lengths do not +# match the requested workload. +INPUT_LENGTH_TOLERANCE_PCT="${INPUT_LENGTH_TOLERANCE_PCT:-5.0}" +OUTPUT_LENGTH_TOLERANCE_PCT="${OUTPUT_LENGTH_TOLERANCE_PCT:-10.0}" + +MAX_POINT_RETRIES="${MAX_POINT_RETRIES:-1}" +SERVER_RESTART_COOLDOWN_S="${SERVER_RESTART_COOLDOWN_S:-10}" +SCENARIO_TIMEOUT_S="${SCENARIO_TIMEOUT_S:-1800}" +GPU_MEM_SAMPLE_INTERVAL_S="${GPU_MEM_SAMPLE_INTERVAL_S:-1}" + +# Optional space-separated filters, useful for smoke tests: +# TP_LIST="8" ISL_LIST="1024" OSL_LIST="128" +TP_LIST="${TP_LIST:-}" +ISL_LIST="${ISL_LIST:-}" +OSL_LIST="${OSL_LIST:-}" + +DRY_RUN="${DRY_RUN:-0}" +# Counts ISL/OSL shapes per TP/DP config, not individual concurrency probes. +GRID_LIMIT="${GRID_LIMIT:-0}" diff --git a/experiments/dsv4_h20_sglang_tp_dp_matrix/compare.py b/experiments/dsv4_h20_sglang_tp_dp_matrix/compare.py new file mode 100644 index 0000000..cfea1dd --- /dev/null +++ b/experiments/dsv4_h20_sglang_tp_dp_matrix/compare.py @@ -0,0 +1,162 @@ +#!/usr/bin/env python3 +"""Cross TP×DP configuration comparison for dsv4_h20_sglang_tp_dp_matrix. + +Usage: + python3 compare.py --run-root results/ [--output comparison.md] +""" +import argparse +import json +import re +from collections import defaultdict +from pathlib import Path + + +def load_result(result_root: Path) -> dict: + path = result_root / "results.json" + with open(path, "r", encoding="utf-8") as f: + return json.load(f) + + +def slo_status(ttft_p95_ms: float, tpot_mean_ms: float, + ttft_limit_ms: float = 3000.0, tpot_limit_ms: float = 50.0) -> str: + ttft_ok = ttft_p95_ms < ttft_limit_ms + tpot_ok = tpot_mean_ms < tpot_limit_ms + if ttft_ok and tpot_ok: + return "PASS" + if ttft_ok or tpot_ok: + return "PARTIAL" + return "FAIL" + + +def gpu_memory_str(gpu: dict | None) -> str: + if not gpu: + return "-" + peak = gpu.get("peak_used_mb", 0) + total = gpu.get("memory_total_mb", 0) + if total: + return f"{peak:.0f}/{total:.0f} ({100*peak/total:.1f}%)" + return f"{peak:.0f}" + + +def main(): + parser = argparse.ArgumentParser() + parser.add_argument("--run-root", type=Path, required=True) + parser.add_argument("-o", "--output", type=Path, default=Path("comparison.md")) + parser.add_argument("--ttft-limit", type=float, default=3000.0) + parser.add_argument("--tpot-limit", type=float, default=50.0) + args = parser.parse_args() + + # Discover configurations: tp*_dp* directories. + configs = [] + for subdir in sorted(args.run_root.iterdir()): + if not subdir.is_dir(): + continue + name = subdir.name + if not (name.startswith("tp") and "_dp" in name): + continue + results_json = subdir / "results.json" + if not results_json.exists(): + continue + configs.append((name, load_result(subdir))) + + if not configs: + print(f"No tp*_dp* results found under {args.run_root}") + return + + model = configs[0][1].get("metadata", {}).get("model", "unknown") + hardware = configs[0][1].get("metadata", {}).get("hardware", "unknown") + + # Group by scenario name. + by_scenario: dict[str, dict[str, dict]] = defaultdict(dict) + skipped: dict[str, dict[str, str]] = defaultdict(dict) + for label, data in configs: + for s in data.get("scenarios", []): + key = s["name"] + if s.get("status") == "skipped_oom": + skipped[key][label] = s.get("note", "skipped") + else: + by_scenario[key][label] = s + + with open(args.output, "w", encoding="utf-8") as f: + f.write(f"# SGLang TP×DP matrix comparison ({hardware})\n\n") + f.write("## Summary\n\n") + f.write(f"- Model: `{model}`\n") + f.write(f"- Hardware: {hardware}\n") + f.write("- Backend: SGLang (Docker)\n") + f.write("- Benchmark client: `sglang.bench_serving`\n") + f.write(f"- SLO reference: TTFT P95 < {args.ttft_limit}ms, TPOT mean < {args.tpot_limit}ms\n\n") + + # Configuration overview. + f.write("### Configurations\n\n") + f.write("| Config | TP | DP | GPUs/replica | Notes |\n") + f.write("|---|---:|---:|---:|---|\n") + for label, data in configs: + cfg = data.get("config", {}) + tp = cfg.get("tp", "?") + dp = cfg.get("dp", "?") + f.write(f"| {label} | {tp} | {dp} | {tp} | server args recorded per ISL in results.json |\n") + f.write("\n") + + # Side-by-side table. + f.write("## Side-by-side results\n\n") + headers = [ + "Scenario", "ISL", "OSL", "Config", "Conc", "Req/s", "OutTok/s", + "TTFT P95(ms)", "TTFT P99(ms)", "TPOT Mean(ms)", "TPOT P95(ms)", + "TPOT P99(ms)", "E2E P99(ms)", "Peak GPU mem", "SLO" + ] + f.write("| " + " | ".join(headers) + " |\n") + f.write("|" + "|".join(["---"] * len(headers)) + "|\n") + + for scenario_name in sorted(by_scenario.keys(), key=lambda x: tuple(map(int, re.findall(r"\d+", x)))): + _, isl, osl = re.findall(r"\d+", scenario_name) + for label, data in configs: + s = by_scenario[scenario_name].get(label) + if s is None: + if scenario_name in skipped and label in skipped[scenario_name]: + note = skipped[scenario_name][label] + f.write(f"| {scenario_name} | {isl} | {osl} | {label} | - | - | - | - | - | - | - | - | - | - | {note} |\n") + continue + cfg = s["config"] + m = s["metrics"] + status = slo_status(m["ttft_ms"]["p95"], m["tpot_ms"]["mean"], args.ttft_limit, args.tpot_limit) + gpu = m.get("gpu_memory") + f.write( + f"| {scenario_name} | {isl} | {osl} | {label} | {cfg['concurrency']} | " + f"{m['request_throughput']:.2f} | {m['output_token_throughput']:.2f} | " + f"{m['ttft_ms']['p95']:.2f} | {m['ttft_ms']['p99']:.2f} | " + f"{m['tpot_ms']['mean']:.2f} | {m['tpot_ms']['p95']:.2f} | {m['tpot_ms']['p99']:.2f} | " + f"{m['e2e_ms']['p99']:.2f} | {gpu_memory_str(gpu)} | {status} |\n" + ) + + # Best throughput per ISL/OSL. + f.write("\n## Best throughput per (ISL, OSL)\n\n") + f.write("| ISL | OSL | Best Config | Concurrency | OutTok/s | TTFT P95(ms) | TPOT Mean(ms) | SLO |\n") + f.write("|---:|---:|---|---:|---:|---:|---:|---:|\n") + best_by_shape: dict[tuple[int, int], tuple[float, str, dict]] = {} + for scenario_name, backends in by_scenario.items(): + _, isl, osl = re.findall(r"\d+", scenario_name) + isl_i, osl_i = int(isl), int(osl) + for label, s in backends.items(): + m = s["metrics"] + out_tok = m["output_token_throughput"] + if (isl_i, osl_i) not in best_by_shape or out_tok > best_by_shape[(isl_i, osl_i)][0]: + best_by_shape[(isl_i, osl_i)] = (out_tok, label, s) + for (isl_i, osl_i), (out_tok, label, s) in sorted(best_by_shape.items()): + m = s["metrics"] + status = slo_status(m["ttft_ms"]["p95"], m["tpot_ms"]["mean"], args.ttft_limit, args.tpot_limit) + f.write( + f"| {isl_i} | {osl_i} | {label} | {s['config']['concurrency']} | " + f"{out_tok:.2f} | {m['ttft_ms']['p95']:.2f} | {m['tpot_ms']['mean']:.2f} | {status} |\n" + ) + + f.write("\n## Notes\n\n") + f.write("- SLO check uses TTFT P95 and TPOT mean.\n") + f.write("- A PARTIAL indicates one of the two metrics is out of target; FAIL indicates both are out.\n") + f.write("- `Peak GPU mem` shows peak used / total MB and utilization percentage.\n") + f.write("- Optional (P) combinations that failed are marked as skipped/OOM and do not break the run.\n") + + print(f"Wrote comparison to {args.output}") + + +if __name__ == "__main__": + main() diff --git a/experiments/dsv4_h20_sglang_tp_dp_matrix/config.env b/experiments/dsv4_h20_sglang_tp_dp_matrix/config.env new file mode 100644 index 0000000..764ffb6 --- /dev/null +++ b/experiments/dsv4_h20_sglang_tp_dp_matrix/config.env @@ -0,0 +1,70 @@ +# TP×DP matrix experiment for DeepSeek-V4-Flash on H20 (8 GPUs) using SGLang. +# Tests SGLang with three parallel configurations: +# TP=2, DP=4 -> 2 GPUs per replica, 4 replicas +# TP=4, DP=2 -> 4 GPUs per replica, 2 replicas +# TP=8, DP=1 -> 8 GPUs, no data parallelism + +EXPERIMENT="dsv4_h20_sglang_tp_dp_matrix" +MODEL_NAME="DeepSeek-V4-Flash" +MODEL_PATH="/data1/hf_models/DeepSeek-V4-Flash" +SERVED_MODEL_NAME="deepseek-v4-flash" + +SGLANG_PORT="${SGLANG_PORT:-30031}" + +# Python interpreter for orchestration scripts (parse_backend.py, compare.py, etc.) +# and the benchmark client. Defaults to the system python3 if the sglang venv +# does not exist on the host. +VENV_CLIENT="${VENV_CLIENT:-/data1/yy/sskj/envs/sglang}" + +# Run the benchmark client natively (0) or inside Docker (1). +USE_DOCKER_CLIENT="${USE_DOCKER_CLIENT:-1}" + +export CUDA_VISIBLE_DEVICES="${CUDA_VISIBLE_DEVICES:-0,1,2,3,4,5,6,7}" + +# Runtime working directory for logs, pid files, and tmp. Defaults to a local +# directory under this experiment so the benchmark is self-contained. +RUNTIME_BASE="${RUNTIME_BASE:-${SCRIPT_DIR}/runtime}" + +# Parallel configurations to test. Format: "TP DP" +declare -a PARALLEL_CONFIGS=( + "2 4" + "4 2" + "8 1" +) + +# SGLang server settings +CONTEXT_LENGTH="${CONTEXT_LENGTH:-1048576}" +MAX_RUNNING_REQUESTS="${MAX_RUNNING_REQUESTS:-128}" +MEM_FRACTION_STATIC="${MEM_FRACTION_STATIC:-0.88}" +MOE_RUNNER_BACKEND="${MOE_RUNNER_BACKEND:-marlin}" + +# Deployment switch. 0 = native sglang venv, 1 = Docker. +USE_DOCKER="${USE_DOCKER:-1}" +DOCKER_IMAGE="${DOCKER_IMAGE:-lmsysorg/sglang:latest}" + +# Dataset used by sglang.bench_serving --dataset-name random. +# The random sampler needs a ShareGPT-style JSON file locally; it falls back to +# downloading from HuggingFace, which usually fails on offline H20 nodes. +DATASET_PATH="${DATASET_PATH:-/data1/yy/sskj/datasets/ShareGPT_V3_unfiltered_cleaned_split.json}" + +# Matrix and concurrency rules are defined in matrix.json by default. +MATRIX_FILE="${MATRIX_FILE:-${SCRIPT_DIR:-.}/matrix.json}" +MATRIX_MODE="${MATRIX_MODE:-Y}" + +# Sampling density for concurrency. +# 0 = use the default heuristic in generate_scenarios.py (6-8 points). +# 2 = only test the low and high endpoints. +export CONCURRENCY_SAMPLES="${CONCURRENCY_SAMPLES:-2}" + +# Per-scenario timeout to avoid hangs (seconds). +SCENARIO_TIMEOUT_S="${SCENARIO_TIMEOUT_S:-1800}" + +# GPU memory sampling interval (seconds). +GPU_MEM_SAMPLE_INTERVAL_S="${GPU_MEM_SAMPLE_INTERVAL_S:-1}" + +# Dry-run mode: if 1, only log the server args and scenario plan without starting +# any server or sending requests. +DRY_RUN="${DRY_RUN:-0}" + +# Per-config scenario limit for quick smoke tests. 0 = run all generated scenarios. +GRID_LIMIT="${GRID_LIMIT:-0}" diff --git a/experiments/dsv4_h20_sglang_tp_dp_matrix/generate_scenarios.py b/experiments/dsv4_h20_sglang_tp_dp_matrix/generate_scenarios.py new file mode 100644 index 0000000..2bfb152 --- /dev/null +++ b/experiments/dsv4_h20_sglang_tp_dp_matrix/generate_scenarios.py @@ -0,0 +1,98 @@ +#!/usr/bin/env python3 +"""Generate the scenario list for the TP×DP matrix experiment. + +Reads matrix.json and prints TSV lines: + + mark input_len output_len concurrency num_prompts + +mark is one of Y/P/N. The caller (run_bench.sh) decides how to treat each. +""" +import argparse +import json +import math +import os +from pathlib import Path + + +def sample_concurrency(low: int, high: int, target: int) -> list[int]: + """Return only the low and high concurrency values in [low, high]. + + For the TP×DP matrix we only need the two endpoints of the concurrency + range (e.g. 1 and 128 for ISL=1024). The `target` argument is kept for + API compatibility but is ignored. + """ + assert 1 <= low <= high, f"invalid concurrency range: {low}-{high}" + if low == high: + return [low] + return [low, high] + + +def generate_scenarios(matrix_path: Path, mode: str, target_samples: int) -> list[dict]: + with open(matrix_path, "r", encoding="utf-8") as f: + data = json.load(f) + + matrix = data["matrix"] + concurrency_cfg = data["concurrency"] + + scenarios = [] + for isl_str in sorted(matrix.keys(), key=int): + osl_map = matrix[isl_str] + low = concurrency_cfg[isl_str]["low"] + high = concurrency_cfg[isl_str]["high"] + concurrencies = sample_concurrency(low, high, target_samples) + + for osl_str in sorted(osl_map.keys(), key=int): + mark = osl_map[osl_str] + if mode == "Y" and mark != "Y": + continue + if mode == "Y+P" and mark not in ("Y", "P"): + continue + # mode == "all" keeps everything, including N. + + for conc in concurrencies: + scenarios.append( + { + "mark": mark, + "input_len": int(isl_str), + "output_len": int(osl_str), + "concurrency": conc, + "num_prompts": conc * 5, + } + ) + + return scenarios + + +def main() -> None: + parser = argparse.ArgumentParser() + parser.add_argument("--matrix", type=Path, default=Path("matrix.json")) + parser.add_argument("--mode", choices=["Y", "Y+P", "all"], default=None, + help="Scenario selection mode. Defaults to matrix.mode.") + parser.add_argument("--target-samples", type=int, default=0, + help="Target number of concurrency samples. 0 = heuristic (6-8).") + args = parser.parse_args() + + with open(args.matrix, "r", encoding="utf-8") as f: + data = json.load(f) + + mode = args.mode if args.mode else data.get("mode", "Y+P") + + target_samples = args.target_samples + if target_samples <= 0: + env_samples = os.getenv("CONCURRENCY_SAMPLES", "0") + try: + target_samples = int(env_samples) + except ValueError: + target_samples = 0 + if target_samples <= 0: + target_samples = 7 + + scenarios = generate_scenarios(args.matrix, mode, target_samples) + + print("mark\tinput_len\toutput_len\tconcurrency\tnum_prompts") + for s in scenarios: + print(f"{s['mark']}\t{s['input_len']}\t{s['output_len']}\t{s['concurrency']}\t{s['num_prompts']}") + + +if __name__ == "__main__": + main() diff --git a/experiments/dsv4_h20_sglang_tp_dp_matrix/matrix.json b/experiments/dsv4_h20_sglang_tp_dp_matrix/matrix.json new file mode 100644 index 0000000..b7e45e3 --- /dev/null +++ b/experiments/dsv4_h20_sglang_tp_dp_matrix/matrix.json @@ -0,0 +1,81 @@ +{ + "comment": "ISL/OSL matrix for dsv4_h20_sglang_tp_dp_matrix. Y=must test, P=optional (record N on failure), N=skip.", + "mode": "Y", + "comment": "Only mandatory (Y) combinations are tested; 1M ISL is excluded per user request.", + "matrix": { + "1024": { + "128": "Y", + "256": "Y", + "512": "Y", + "1024": "Y", + "2048": "Y", + "4096": "Y" + }, + "4096": { + "128": "Y", + "256": "Y", + "512": "Y", + "1024": "Y", + "2048": "Y", + "4096": "Y" + }, + "16384": { + "128": "Y", + "256": "Y", + "512": "Y", + "1024": "Y", + "2048": "Y", + "4096": "P" + }, + "65536": { + "128": "Y", + "256": "Y", + "512": "Y", + "1024": "Y", + "2048": "P", + "4096": "N" + }, + "131072": { + "128": "Y", + "256": "Y", + "512": "Y", + "1024": "P", + "2048": "N", + "4096": "N" + }, + "262144": { + "128": "Y", + "256": "Y", + "512": "P", + "1024": "N", + "2048": "N", + "4096": "N" + }, + "524288": { + "128": "Y", + "256": "P", + "512": "N", + "1024": "N", + "2048": "N", + "4096": "N" + }, + "1048576": { + "128": "Y", + "256": "P", + "512": "N", + "1024": "N", + "2048": "N", + "4096": "N" + } + }, + "concurrency": { + "1024": { "low": 1, "high": 128 }, + "4096": { "low": 1, "high": 64 }, + "16384": { "low": 1, "high": 32 }, + "65536": { "low": 1, "high": 8 }, + "131072": { "low": 1, "high": 4 }, + "262144": { "low": 1, "high": 2 }, + "524288": { "low": 1, "high": 2 }, + "1048576": { "low": 1, "high": 2 } + } +} diff --git a/experiments/dsv4_h20_sglang_tp_dp_matrix/run_adaptive_concurrency.sh b/experiments/dsv4_h20_sglang_tp_dp_matrix/run_adaptive_concurrency.sh new file mode 100755 index 0000000..0979a46 --- /dev/null +++ b/experiments/dsv4_h20_sglang_tp_dp_matrix/run_adaptive_concurrency.sh @@ -0,0 +1,181 @@ +#!/usr/bin/env bash +# Find the Total-TPS saturation concurrency for each SGLang TP/DP/ISL/OSL shape. +set -Eeuo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +EXPERIMENT_NAME="$(basename "$SCRIPT_DIR")" + +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/lib.sh" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/platform.sh" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/config.env" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/adaptive_config.env" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/adaptive_bench_lib.sh" + +ENGINE="sglang" +ENGINE_PORT="$SGLANG_PORT" +RESULT_BASE="${RESULT_BASE:-${SCRIPT_DIR}/adaptive_results}" +ACTIVE_ENGINE_SERVER_LOG="" + +if [[ -x "${VENV_CLIENT}/bin/python" ]]; then + PYTHON="${VENV_CLIENT}/bin/python" +else + PYTHON="$(command -v python3)" +fi + +DOCKER_IMAGE="${DOCKER_IMAGE:-lmsysorg/sglang:latest}" + +engine_is_healthy() { + curl --fail --silent --show-error --max-time 5 \ + "http://127.0.0.1:${ENGINE_PORT}/health" >/dev/null 2>&1 +} + +engine_stop_server() { + local tp="$1" + local dp="$2" + local pid_file="${RUNTIME_BASE}/${EXPERIMENT}_sglang_tp${tp}_dp${dp}.pid" + + if [[ -f "$pid_file" ]]; then + local pid + pid="$(cat "$pid_file")" + if [[ -n "${CONTAINER_NAME:-}" ]]; then + if docker exec "$CONTAINER_NAME" kill -0 "$pid" 2>/dev/null; then + log "stopping sglang in persistent container pid=${pid} tp=${tp} dp=${dp}" + docker exec "$CONTAINER_NAME" kill "$pid" 2>/dev/null || true + sleep 5 + docker exec "$CONTAINER_NAME" kill -9 "$pid" 2>/dev/null || true + fi + elif kill -0 "$pid" 2>/dev/null; then + log "stopping sglang server pid=${pid} tp=${tp} dp=${dp}" + kill "$pid" 2>/dev/null || true + sleep 5 + kill -9 "$pid" 2>/dev/null || true + fi + rm -f "$pid_file" + fi + if [[ -z "${CONTAINER_NAME:-}" ]]; then + docker rm -f "${EXPERIMENT}_sglang_tp${tp}_dp${dp}" >/dev/null 2>&1 || true + fi + ACTIVE_ENGINE_SERVER_LOG="" + sleep 2 +} + +engine_build_server_args() { + local tp="$1" + local dp="$2" + local -a args=( + sglang serve --model-path "$MODEL_PATH" + --trust-remote-code + --tp-size "$tp" + --moe-runner-backend "$MOE_RUNNER_BACKEND" + --context-length "$CONTEXT_LENGTH" + --max-running-requests "$MAX_RUNNING_REQUESTS" + --mem-fraction-static "$MEM_FRACTION_STATIC" + --host 0.0.0.0 + --port "$ENGINE_PORT" + ) + if (( dp > 1 )); then + args+=(--dp-size "$dp") + fi + printf '%q ' "${args[@]}" +} + +engine_start_server() { + local tp="$1" + local dp="$2" + local outer_log="${ADAPTIVE_LOG_DIR}/sglang_tp${tp}_dp${dp}.server.outer.log" + log "starting sglang server tp=${tp} dp=${dp}" + if [[ -n "${CONTAINER_NAME:-}" ]]; then + bash "${SCRIPT_DIR}/run_sglang_in_container.sh" "$tp" "$dp" >> "$outer_log" 2>&1 + else + bash "${SCRIPT_DIR}/start_sglang_dp.sh" "$tp" "$dp" >> "$outer_log" 2>&1 + fi + if ! engine_is_healthy; then + log "ERROR: sglang health check failed tp=${tp} dp=${dp}" + return 1 + fi + if [[ -z "${CONTAINER_NAME:-}" ]]; then + ACTIVE_ENGINE_SERVER_LOG="$( + find "${RUNTIME_BASE}/logs" -maxdepth 1 -type f \ + -name "${EXPERIMENT}_sglang*tp${tp}_dp${dp}_*.log" \ + -printf '%T@ %p\n' 2>/dev/null | sort -nr | head -n 1 | cut -d' ' -f2- + )" + fi + log "sglang server healthy tp=${tp} dp=${dp} log=${ACTIVE_ENGINE_SERVER_LOG:-container:/tmp/sglang_server.log}" +} + +engine_detect_oom() { + local detail_log="$1" + local tp="$2" + local dp="$3" + local pattern='CUDA out of memory|torch\.OutOfMemoryError|OutOfMemory|out of memory|OOM|RESOURCE_EXHAUSTED|Failed to allocate memory' + local outer_log="${ADAPTIVE_LOG_DIR}/sglang_tp${tp}_dp${dp}.server.outer.log" + local -a logs=("$detail_log" "$outer_log") + if [[ -n "$ACTIVE_ENGINE_SERVER_LOG" ]]; then + logs+=("$ACTIVE_ENGINE_SERVER_LOG") + fi + if grep -Eiq "$pattern" "${logs[@]}" 2>/dev/null; then + return 0 + fi + if [[ -n "${CONTAINER_NAME:-}" ]]; then + docker exec "$CONTAINER_NAME" grep -Eiq "$pattern" /tmp/sglang_server.log 2>/dev/null + return $? + fi + return 1 +} + +engine_run_bench() { + local isl="$1" + local osl="$2" + local concurrency="$3" + local num_prompts="$4" + local output_file="$5" + local warmup_requests + warmup_requests="$(adaptive_warmup_request_count "$concurrency")" + local -a bench_args=( + --backend sglang + --host 127.0.0.1 + --port "$ENGINE_PORT" + --dataset-name "$BENCH_DATASET_NAME" + --random-input-len "$isl" + --random-output-len "$osl" + --random-range-ratio "$RANDOM_RANGE_RATIO" + --num-prompts "$num_prompts" + --max-concurrency "$concurrency" + --request-rate 10000 + --warmup-requests "$warmup_requests" + --output-file "$output_file" + --output-details + --disable-tqdm + ) + if [[ "$BENCH_DATASET_NAME" == "random" ]]; then + bench_args+=(--dataset-path "$DATASET_PATH") + else + bench_args+=(--tokenize-prompt) + fi + + if [[ "$USE_DOCKER_CLIENT" == "1" ]]; then + local -a volume_args=(-v "${MODEL_PATH}:${MODEL_PATH}:ro" -v "${RESULT_BASE}:${RESULT_BASE}") + if [[ "$BENCH_DATASET_NAME" == "random" ]]; then + volume_args+=(-v "${DATASET_PATH}:${DATASET_PATH}:ro") + fi + docker run --rm \ + --network host \ + "${volume_args[@]}" \ + -e PYTHONUNBUFFERED=1 \ + "$DOCKER_IMAGE" \ + python -m sglang.bench_serving "${bench_args[@]}" + else + "$PYTHON" -m sglang.bench_serving "${bench_args[@]}" + fi +} + +export -f engine_run_bench +export ENGINE_PORT MODEL_PATH RESULT_BASE DOCKER_IMAGE USE_DOCKER_CLIENT +export BENCH_DATASET_NAME DATASET_PATH RANDOM_RANGE_RATIO BENCH_WARMUP_MAX_REQUESTS PYTHON + +adaptive_main "$@" diff --git a/experiments/dsv4_h20_sglang_tp_dp_matrix/run_adaptive_concurrency_add16.sh b/experiments/dsv4_h20_sglang_tp_dp_matrix/run_adaptive_concurrency_add16.sh new file mode 100755 index 0000000..b9b28f3 --- /dev/null +++ b/experiments/dsv4_h20_sglang_tp_dp_matrix/run_adaptive_concurrency_add16.sh @@ -0,0 +1,184 @@ +#!/usr/bin/env bash +# Find the Total-TPS saturation concurrency for each SGLang TP/DP/ISL/OSL shape. +set -Eeuo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +EXPERIMENT_NAME="$(basename "$SCRIPT_DIR")" + +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/lib.sh" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/platform.sh" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/config.env" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/adaptive_config.env" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/adaptive_bench_lib.sh" + +ENGINE="sglang" +ENGINE_PORT="$SGLANG_PORT" +RESULT_BASE="${RESULT_BASE:-${SCRIPT_DIR}/adaptive_results}" +ACTIVE_ENGINE_SERVER_LOG="" + +if [[ -x "${VENV_CLIENT}/bin/python" ]]; then + PYTHON="${VENV_CLIENT}/bin/python" +else + PYTHON="$(command -v python3)" +fi + +DOCKER_IMAGE="${DOCKER_IMAGE:-lmsysorg/sglang:latest}" + +engine_is_healthy() { + curl --fail --silent --show-error --max-time 5 \ + "http://127.0.0.1:${ENGINE_PORT}/health" >/dev/null 2>&1 +} + +engine_stop_server() { + local tp="$1" + local dp="$2" + local pid_file="${RUNTIME_BASE}/${EXPERIMENT}_sglang_tp${tp}_dp${dp}.pid" + + if [[ -f "$pid_file" ]]; then + local pid + pid="$(cat "$pid_file")" + if [[ -n "${CONTAINER_NAME:-}" ]]; then + if docker exec "$CONTAINER_NAME" kill -0 "$pid" 2>/dev/null; then + log "stopping sglang in persistent container pid=${pid} tp=${tp} dp=${dp}" + docker exec "$CONTAINER_NAME" kill "$pid" 2>/dev/null || true + sleep 5 + docker exec "$CONTAINER_NAME" kill -9 "$pid" 2>/dev/null || true + fi + elif kill -0 "$pid" 2>/dev/null; then + log "stopping sglang server pid=${pid} tp=${tp} dp=${dp}" + kill "$pid" 2>/dev/null || true + sleep 5 + kill -9 "$pid" 2>/dev/null || true + fi + rm -f "$pid_file" + fi + if [[ -z "${CONTAINER_NAME:-}" ]]; then + docker rm -f "${EXPERIMENT}_sglang_tp${tp}_dp${dp}" >/dev/null 2>&1 || true + fi + ACTIVE_ENGINE_SERVER_LOG="" + sleep 2 +} + +engine_build_server_args() { + local tp="$1" + local dp="$2" + local -a args=( + sglang serve --model-path "$MODEL_PATH" + --trust-remote-code + --tp-size "$tp" + --moe-runner-backend "$MOE_RUNNER_BACKEND" + --context-length "$CONTEXT_LENGTH" + --max-running-requests "$MAX_RUNNING_REQUESTS" + --mem-fraction-static "$MEM_FRACTION_STATIC" + --host 0.0.0.0 + --port "$ENGINE_PORT" + ) + if (( dp > 1 )); then + args+=(--dp-size "$dp") + fi + printf '%q ' "${args[@]}" +} + +engine_start_server() { + local tp="$1" + local dp="$2" + local outer_log="${ADAPTIVE_LOG_DIR}/sglang_tp${tp}_dp${dp}.server.outer.log" + log "starting sglang server tp=${tp} dp=${dp}" + if [[ -n "${CONTAINER_NAME:-}" ]]; then + bash "${SCRIPT_DIR}/run_sglang_in_container.sh" "$tp" "$dp" >> "$outer_log" 2>&1 + else + bash "${SCRIPT_DIR}/start_sglang_dp.sh" "$tp" "$dp" >> "$outer_log" 2>&1 + fi + if ! engine_is_healthy; then + log "ERROR: sglang health check failed tp=${tp} dp=${dp}" + return 1 + fi + if [[ -z "${CONTAINER_NAME:-}" ]]; then + ACTIVE_ENGINE_SERVER_LOG="$( + find "${RUNTIME_BASE}/logs" -maxdepth 1 -type f \ + -name "${EXPERIMENT}_sglang*tp${tp}_dp${dp}_*.log" \ + -printf '%T@ %p\n' 2>/dev/null | sort -nr | head -n 1 | cut -d' ' -f2- + )" + fi + log "sglang server healthy tp=${tp} dp=${dp} log=${ACTIVE_ENGINE_SERVER_LOG:-container:/tmp/sglang_server.log}" +} + +engine_detect_oom() { + local detail_log="$1" + local tp="$2" + local dp="$3" + local pattern='CUDA out of memory|torch\.OutOfMemoryError|OutOfMemory|out of memory|OOM|RESOURCE_EXHAUSTED|Failed to allocate memory' + local outer_log="${ADAPTIVE_LOG_DIR}/sglang_tp${tp}_dp${dp}.server.outer.log" + local -a logs=("$detail_log" "$outer_log") + if [[ -n "$ACTIVE_ENGINE_SERVER_LOG" ]]; then + logs+=("$ACTIVE_ENGINE_SERVER_LOG") + fi + if grep -Eiq "$pattern" "${logs[@]}" 2>/dev/null; then + return 0 + fi + if [[ -n "${CONTAINER_NAME:-}" ]]; then + docker exec "$CONTAINER_NAME" grep -Eiq "$pattern" /tmp/sglang_server.log 2>/dev/null + return $? + fi + return 1 +} + +engine_run_bench() { + local isl="$1" + local osl="$2" + local concurrency="$3" + local num_prompts="$4" + local output_file="$5" + local warmup_requests + warmup_requests="$(adaptive_warmup_request_count "$concurrency")" + local -a bench_args=( + --backend sglang + --host 127.0.0.1 + --port "$ENGINE_PORT" + --dataset-name "$BENCH_DATASET_NAME" + --random-input-len "$isl" + --random-output-len "$osl" + --random-range-ratio "$RANDOM_RANGE_RATIO" + --num-prompts "$num_prompts" + --max-concurrency "$concurrency" + --request-rate 10000 + --warmup-requests "$warmup_requests" + --output-file "$output_file" + --output-details + --disable-tqdm + ) + if [[ "$BENCH_DATASET_NAME" == "random" ]]; then + bench_args+=(--dataset-path "$DATASET_PATH") + else + bench_args+=(--tokenize-prompt) + fi + + if [[ "$USE_DOCKER_CLIENT" == "1" ]]; then + local -a volume_args=(-v "${MODEL_PATH}:${MODEL_PATH}:ro" -v "${RESULT_BASE}:${RESULT_BASE}") + if [[ "$BENCH_DATASET_NAME" == "random" ]]; then + volume_args+=(-v "${DATASET_PATH}:${DATASET_PATH}:ro") + fi + docker run --rm \ + --network host \ + "${volume_args[@]}" \ + -e PYTHONUNBUFFERED=1 \ + "$DOCKER_IMAGE" \ + python -m sglang.bench_serving "${bench_args[@]}" + else + "$PYTHON" -m sglang.bench_serving "${bench_args[@]}" + fi +} + +export -f engine_run_bench +export ENGINE_PORT MODEL_PATH RESULT_BASE DOCKER_IMAGE USE_DOCKER_CLIENT +export BENCH_DATASET_NAME DATASET_PATH RANDOM_RANGE_RATIO BENCH_WARMUP_MAX_REQUESTS PYTHON + +export SEARCH_START_CONCURRENCY=16 +export SEARCH_ADDEND=16 + +adaptive_main "$@" diff --git a/experiments/dsv4_h20_sglang_tp_dp_matrix/run_bench.sh b/experiments/dsv4_h20_sglang_tp_dp_matrix/run_bench.sh new file mode 100755 index 0000000..d27b870 --- /dev/null +++ b/experiments/dsv4_h20_sglang_tp_dp_matrix/run_bench.sh @@ -0,0 +1,558 @@ +#!/usr/bin/env bash +# TP×DP matrix benchmark for DeepSeek-V4-Flash on SGLang (Docker). +set -Eeuo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +EXPERIMENT_NAME="$(basename "$SCRIPT_DIR")" + +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/lib.sh" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/platform.sh" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/config.env" + +RUN_ID="${RUN_ID:-$(date '+%Y%m%d-%H%M%S')}" +RESULT_BASE="${SCRIPT_DIR}/results" +MATRIX_FILE="${MATRIX_FILE:-${SCRIPT_DIR}/matrix.json}" +MATRIX_MODE="${MATRIX_MODE:-Y}" +SCENARIO_TIMEOUT_S="${SCENARIO_TIMEOUT_S:-1800}" +GPU_MEM_SAMPLE_INTERVAL_S="${GPU_MEM_SAMPLE_INTERVAL_S:-1}" +DRY_RUN="${DRY_RUN:-0}" +GRID_LIMIT="${GRID_LIMIT:-0}" + +# Export variables used inside functions that are called via bash -c subshells. +export DATASET_PATH MODEL_PATH RESULT_BASE DOCKER_IMAGE USE_DOCKER_CLIENT + +if [[ -x "${VENV_CLIENT}/bin/python" ]]; then + PYTHON="${VENV_CLIENT}/bin/python" +else + PYTHON="$(command -v python3)" +fi + +DOCKER_IMAGE="${DOCKER_IMAGE:-lmsysorg/sglang:latest}" + +log_dir_global="${RESULT_BASE}/${RUN_ID}/logs" +mkdir -p "$log_dir_global" +log_init "${log_dir_global}/orchestrator.log" + +log "experiment=${EXPERIMENT_NAME} run_id=${RUN_ID} platform=${PLATFORM} hardware=${HARDWARE}" +log "matrix_mode=${MATRIX_MODE} matrix_file=${MATRIX_FILE} dry_run=${DRY_RUN} grid_limit=${GRID_LIMIT}" + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + +is_server_healthy() { + curl --fail --silent --show-error --max-time 5 "http://127.0.0.1:${SGLANG_PORT}/health" >/dev/null 2>&1 +} + +stop_server() { + local tp="$1" + local dp="$2" + + local pid_file="${RUNTIME_BASE}/${EXPERIMENT}_sglang_tp${tp}_dp${dp}.pid" + if [[ -f "$pid_file" ]]; then + local pid + pid="$(cat "$pid_file")" + # In container-reuse mode, pid is the server PID inside the container. + if [[ -n "${CONTAINER_NAME:-}" ]]; then + if docker exec "$CONTAINER_NAME" kill -0 "$pid" 2>/dev/null; then + log "stopping sglang server inside container (tp=${tp}, dp=${dp}, pid=${pid})" + docker exec "$CONTAINER_NAME" kill "$pid" 2>/dev/null || true + sleep 5 + docker exec "$CONTAINER_NAME" kill -9 "$pid" 2>/dev/null || true + fi + else + if kill -0 "$pid" 2>/dev/null; then + log "stopping sglang server pid=${pid} (tp=${tp}, dp=${dp})" + kill "$pid" 2>/dev/null || true + sleep 5 + kill -9 "$pid" 2>/dev/null || true + fi + fi + rm -f "$pid_file" + fi + + # Fallback: remove any Docker container started by this experiment (legacy mode). + if [[ -z "${CONTAINER_NAME:-}" ]]; then + docker rm -f "${EXPERIMENT}_sglang_tp${tp}_dp${dp}" >/dev/null 2>&1 || true + fi + + # Fallback: kill any sglang serve processes for this model. + pkill -9 -f "sglang serve.*${MODEL_NAME}" 2>/dev/null || true + pkill -9 -f "sglang serve.*${MODEL_PATH}" 2>/dev/null || true + sleep 2 +} + +build_server_args() { + local tp="$1" + local dp="$2" + + local args=( + "sglang serve" "$MODEL_PATH" + --trust-remote-code + --tp-size "$tp" + --moe-runner-backend "$MOE_RUNNER_BACKEND" + --context-length "$CONTEXT_LENGTH" + --max-running-requests "$MAX_RUNNING_REQUESTS" + --mem-fraction-static "$MEM_FRACTION_STATIC" + --host 0.0.0.0 + --port "$SGLANG_PORT" + ) + if [[ "$dp" -gt 1 ]]; then + args+=( + --dp-size "$dp" + ) + fi + printf '%s ' "${args[@]}" +} + +start_server() { + local tp="$1" + local dp="$2" + + log "starting sglang server tp=${tp} dp=${dp}" + if [[ -n "${CONTAINER_NAME:-}" ]]; then + bash "${SCRIPT_DIR}/run_sglang_in_container.sh" "$tp" "$dp" \ + >> "${log_dir_global}/sglang_tp${tp}_dp${dp}.server.outer.log" 2>&1 + else + bash "${SCRIPT_DIR}/start_sglang_dp.sh" "$tp" "$dp" \ + >> "${log_dir_global}/sglang_tp${tp}_dp${dp}.server.outer.log" 2>&1 + fi + + if ! is_server_healthy; then + log "error: sglang server tp=${tp} dp=${dp} failed health check on port ${SGLANG_PORT}" + return 1 + fi + log "sglang server tp=${tp} dp=${dp} is healthy on port ${SGLANG_PORT}" +} + +restart_server() { + local tp="$1" + local dp="$2" + log "restarting sglang server tp=${tp} dp=${dp} after non-OOM failure" + stop_server "$tp" "$dp" + sleep 10 + start_server "$tp" "$dp" +} + +run_bench_serving() { + # Run sglang.bench_serving either natively or inside the SGLang Docker image. + if [[ "${USE_DOCKER_CLIENT:-1}" == "1" ]]; then + local vol_args=() + vol_args+=("-v" "${MODEL_PATH}:${MODEL_PATH}:ro") + if [[ -n "${DATASET_PATH:-}" ]]; then + vol_args+=("-v" "${DATASET_PATH}:${DATASET_PATH}:ro") + fi + vol_args+=("-v" "${RESULT_BASE}:${RESULT_BASE}") + docker run --rm \ + --network host \ + "${vol_args[@]}" \ + -e PYTHONUNBUFFERED=1 \ + "${DOCKER_IMAGE}" \ + python -m sglang.bench_serving "$@" + else + "$PYTHON" -m sglang.bench_serving "$@" + fi +} +export -f run_bench_serving + +run_warmup() { + local input_len="$1" + local output_len="$2" + + log "warming up (input=${input_len}, output=${output_len}, num=1)" + bash -c ' + run_bench_serving \ + --backend sglang \ + --host 127.0.0.1 \ + --port "'"$SGLANG_PORT"'" \ + --dataset-name random \ + --dataset-path "'"$DATASET_PATH"'" \ + --random-input-len "'"$input_len"'" \ + --random-output-len "'"$output_len"'" \ + --num-prompts 1 \ + --max-concurrency 1 \ + --request-rate 10000 \ + --output-file /dev/null \ + --output-details \ + >> "'"${log_dir_global}/warmup.log"'" 2>&1 + ' + log "warmup completed" +} + +scenario_already_completed() { + local output_file="$1" + local expected="$2" + [[ -s "$output_file" ]] || return 1 + local completed + completed="$("$PYTHON" -c " +import json, sys +path = sys.argv[1] +try: + with open(path, 'r', encoding='utf-8') as f: + for line in f: + line = line.strip() + if line: + data = json.loads(line) + print(data.get('completed', 0)) + break +except Exception: + print(0) +" "$output_file")" + [[ "${completed:-0}" -ge "$expected" ]] +} + +scenario_already_processed() { + local result_root="$1" + local scenario_name="$2" + local json_path="${result_root}/results.json" + [[ -f "$json_path" ]] || return 1 + "$PYTHON" -c " +import json, sys +path, name = sys.argv[1], sys.argv[2] +try: + with open(path, 'r', encoding='utf-8') as f: + data = json.load(f) + for s in data.get('scenarios', []): + if s.get('name') == name: + if s.get('status') or s.get('metrics', {}).get('success', 0) > 0: + sys.exit(0) +except Exception: + pass +sys.exit(1) +" "$json_path" "$scenario_name" +} + +detect_oom() { + local detail_log="$1" + local server_outer_log="$2" + local pattern='CUDA out of memory|torch\.OutOfMemoryError|OutOfMemory|out of memory|OOM|RESOURCE_EXHAUSTED|Failed to allocate memory' + if grep -Eiq "$pattern" "$detail_log" "$server_outer_log" 2>/dev/null; then + return 0 + fi + return 1 +} + +start_gpu_monitor() { + local csv_path="$1" + mkdir -p "$(dirname "$csv_path")" + nvidia-smi \ + --query-gpu=timestamp,index,memory.used,memory.total,utilization.gpu \ + --format=csv \ + -l "$GPU_MEM_SAMPLE_INTERVAL_S" \ + > "$csv_path" 2>/dev/null & + echo $! +} + +stop_gpu_monitor() { + local pid="$1" + if kill -0 "$pid" 2>/dev/null; then + kill "$pid" 2>/dev/null || true + sleep 1 + kill -9 "$pid" 2>/dev/null || true + fi +} + +append_scenario_record() { + local result_root="$1" + local json_path="$result_root/results.json" + shift + local scenario_json + scenario_json="$("$PYTHON" -c " +import json, sys +pairs = [a.split('=', 1) for a in sys.argv[1:]] +d = {} +for k, v in pairs: + try: + d[k] = json.loads(v) + except json.JSONDecodeError: + d[k] = v +print(json.dumps(d, ensure_ascii=False)) +" "$@")" + PYTHON="$PYTHON" append_scenario_to_json "$json_path" "$scenario_json" +} + +record_skipped_csv() { + local csv_path="$1" + shift + # Args: key=value + local row + row="$("$PYTHON" -c " +import csv, json, sys, io +pairs = [a.split('=', 1) for a in sys.argv[1:]] +d = {} +for k, v in pairs: + try: + d[k] = json.loads(v) + except json.JSONDecodeError: + d[k] = v +buf = io.StringIO() +writer = csv.DictWriter(buf, fieldnames=['engine','tp','dp','mark','isl','osl','concurrency','status','reason','detail_log'], extrasaction='ignore') +writer.writerow(d) +print(buf.getvalue().strip()) +" "$@")" + echo "$row" >> "$csv_path" +} + +skip_remaining_scenarios() { + local result_root="$1" + local scenario_tsv="$2" + local start_index="$3" + local status="$4" + local reason="$5" + local tp="$6" + local dp="$7" + local skipped_csv="${RESULT_BASE}/${RUN_ID}/skipped_after_oom.csv" + + local i=0 + tail -n +2 "$scenario_tsv" | while IFS=$'\t' read -r mark isl osl conc num; do + if (( i < start_index )); then + i=$((i + 1)) + continue + fi + i=$((i + 1)) + local sname="c${conc}_i${isl}_o${osl}" + if scenario_already_processed "$result_root" "$sname"; then + continue + fi + append_scenario_record "$result_root" \ + "name=${sname}" \ + "config=$(jq -n --arg phase main --argjson c "$conc" --argjson i "$isl" --argjson o "$osl" --arg dataset random --argjson n "$num" '{phase: $phase, concurrency: $c, input_len: $i, output_len: $o, dataset: $dataset, num_prompts: $n}')" \ + "status=\"${status}\"" \ + "note=\"${reason}\"" + record_skipped_csv "$skipped_csv" \ + "engine=sglang" "tp=${tp}" "dp=${dp}" "mark=${mark}" "isl=${isl}" "osl=${osl}" "concurrency=${conc}" "status=${status}" "reason=${reason}" + done +} + +# --------------------------------------------------------------------------- +# Per-configuration runner +# --------------------------------------------------------------------------- + +run_parallel_config() { + local tp="$1" + local dp="$2" + local config_label="tp${tp}_dp${dp}" + local result_root="${RESULT_BASE}/${RUN_ID}/${config_label}" + local raw_dir="${result_root}/raw_outputs" + local gpu_log_dir="${result_root}/gpu_logs" + local phase_log_dir="${result_root}/logs" + mkdir -p "$raw_dir" "$gpu_log_dir" "$phase_log_dir" + + log "===== ${config_label} START =====" + + # Generate scenario list for this config. + local scenario_tsv="${result_root}/scenarios.tsv" + "$PYTHON" "${SCRIPT_DIR}/generate_scenarios.py" \ + --matrix "$MATRIX_FILE" \ + --mode "$MATRIX_MODE" \ + > "$scenario_tsv" + local total_scenarios + total_scenarios="$(tail -n +2 "$scenario_tsv" | wc -l)" + log "generated ${total_scenarios} scenarios for ${config_label}" + + # Write metadata. + ensure_result_root "$result_root" + write_metadata_json \ + "${result_root}/results.json" \ + "${EXPERIMENT_NAME}_${config_label}" \ + "$RUN_ID" \ + "$MODEL_PATH" \ + "sglang" \ + "sglang" \ + "$HARDWARE" \ + "$ACCELERATOR" \ + "$CHIP" \ + "experiments/${EXPERIMENT_NAME}/run_bench.sh" \ + "$DOCKER_IMAGE" \ + "H20 SGLang TP×DP matrix for DeepSeek-V4-Flash" + + local server_args_str + server_args_str="$(build_server_args "$tp" "$dp")" + jq --arg tp "$tp" --arg dp "$dp" --arg cuda "$CUDA_VISIBLE_DEVICES" --arg args "$server_args_str" \ + '.config = { + "tp": ($tp | tonumber), + "dp": ($dp | tonumber), + "cuda_visible_devices": $cuda, + "backend": "sglang", + "server_start_script": "experiments/'${EXPERIMENT_NAME}'/start_sglang_dp.sh", + "server_args": $args + }' "${result_root}/results.json" > "${result_root}/results.json.tmp" && \ + mv "${result_root}/results.json.tmp" "${result_root}/results.json" + + if [[ "$DRY_RUN" == "1" ]]; then + log "DRY_RUN: would start server with args: ${server_args_str}" + local line + tail -n +2 "$scenario_tsv" | while IFS=$'\t' read -r mark isl osl conc num; do + log "DRY_RUN: ${config_label} scenario mark=${mark} c=${conc} i=${isl} o=${osl} n=${num}" + done + log "===== ${config_label} DONE (dry run) =====" + return 0 + fi + + # Initialize skipped_after_oom.csv for this run. + local skipped_csv="${RESULT_BASE}/${RUN_ID}/skipped_after_oom.csv" + if [[ ! -f "$skipped_csv" ]]; then + echo "engine,tp,dp,mark,isl,osl,concurrency,status,reason,detail_log" > "$skipped_csv" + fi + + # Start server once for this TP×DP config. + if ! start_server "$tp" "$dp"; then + log "ERROR: ${config_label} failed to start; skipping all scenarios" + skip_remaining_scenarios "$result_root" "$scenario_tsv" 0 "SKIPPED_SERVICE_START_FAILED" "service failed to start" "$tp" "$dp" + log "===== ${config_label} DONE =====" + return 0 + fi + + # Warmup with a small prompt before the first scenario. + run_warmup 1024 128 || true + + # Read scenarios into an array so we can skip remaining entries on failure. + local -a scenarios=() + while IFS= read -r line; do + scenarios+=("$line") + done < <(tail -n +2 "$scenario_tsv") + + local i mark isl osl conc num + local output_file detail_log gpu_csv sname bench_rc + for (( i = 0; i < ${#scenarios[@]}; i++ )); do + IFS=$'\t' read -r mark isl osl conc num <<< "${scenarios[$i]}" + + if [[ "$GRID_LIMIT" -gt 0 && "$i" -ge "$GRID_LIMIT" ]]; then + log "GRID_LIMIT=${GRID_LIMIT} reached; skipping remaining scenarios" + skip_remaining_scenarios "$result_root" "$scenario_tsv" "$i" "SKIPPED_GRID_LIMIT" "GRID_LIMIT reached" "$tp" "$dp" + break + fi + + sname="c${conc}_i${isl}_o${osl}" + output_file="${raw_dir}/sglang_main_${conc}_${isl}_${osl}.jsonl" + detail_log="${phase_log_dir}/sglang_${config_label}_${sname}.log" + gpu_csv="${gpu_log_dir}/gpu_mem_${conc}_${isl}_${osl}.csv" + + if scenario_already_completed "$output_file" "$num" || scenario_already_processed "$result_root" "$sname"; then + log "skipping already-processed ${config_label} scenario: ${sname}" + continue + fi + + log "running ${config_label} scenario: mark=${mark} c=${conc} i=${isl} o=${osl} n=${num}" + + local gpu_pid + gpu_pid="$(start_gpu_monitor "$gpu_csv")" + + bench_rc=0 + timeout "$SCENARIO_TIMEOUT_S" bash -c ' + run_bench_serving \ + --backend sglang \ + --host 127.0.0.1 \ + --port "'"$SGLANG_PORT"'" \ + --dataset-name random \ + --dataset-path "'"$DATASET_PATH"'" \ + --random-input-len "'"$isl"'" \ + --random-output-len "'"$osl"'" \ + --num-prompts "'"$num"'" \ + --max-concurrency "'"$conc"'" \ + --request-rate 10000 \ + --output-file "'"$output_file"'" \ + --output-details \ + > "'"$detail_log"'" 2>&1 + ' || bench_rc=$? + + stop_gpu_monitor "$gpu_pid" + + if [[ "$bench_rc" -eq 0 ]]; then + log "finished ${config_label} scenario: output=${output_file}" + append_scenario_record "$result_root" \ + "name=${sname}" \ + "config=$(jq -n --arg phase main --argjson c "$conc" --argjson i "$isl" --argjson o "$osl" --arg dataset random --argjson n "$num" '{phase: $phase, concurrency: $c, input_len: $i, output_len: $o, dataset: $dataset, num_prompts: $n}')" \ + "status=\"completed\"" \ + "note=\"benchmark finished successfully\"" + continue + fi + + # Failure handling. + if detect_oom "$detail_log" "${log_dir_global}/sglang_tp${tp}_dp${dp}.server.outer.log"; then + log "ERROR: ${config_label} scenario ${sname} triggered OOM; stopping config" + append_scenario_record "$result_root" \ + "name=${sname}" \ + "config=$(jq -n --arg phase main --argjson c "$conc" --argjson i "$isl" --argjson o "$osl" --arg dataset random --argjson n "$num" '{phase: $phase, concurrency: $c, input_len: $i, output_len: $o, dataset: $dataset, num_prompts: $n}')" \ + "status=\"OOM\"" \ + "note=\"detected CUDA out-of-memory\"" + record_skipped_csv "$skipped_csv" \ + "engine=sglang" "tp=${tp}" "dp=${dp}" "mark=${mark}" "isl=${isl}" "osl=${osl}" "concurrency=${conc}" "status=OOM" "reason=detected CUDA out-of-memory" "detail_log=${detail_log}" + stop_server "$tp" "$dp" + skip_remaining_scenarios "$result_root" "$scenario_tsv" "$((i + 1))" "SKIPPED_AFTER_OOM" "previous case OOM" "$tp" "$dp" + break + fi + + log "ERROR: ${config_label} scenario ${sname} failed (rc=${bench_rc}); see ${detail_log}" + if [[ "$mark" == "P" ]]; then + log "optional (P) scenario failed; recording as skipped and continuing" + append_scenario_record "$result_root" \ + "name=${sname}" \ + "config=$(jq -n --arg phase main --argjson c "$conc" --argjson i "$isl" --argjson o "$osl" --arg dataset random --argjson n "$num" '{phase: $phase, concurrency: $c, input_len: $i, output_len: $o, dataset: $dataset, num_prompts: $n}')" \ + "status=\"skipped_optional\"" \ + "note=\"optional scenario failed (rc=${bench_rc})\"" + record_skipped_csv "$skipped_csv" \ + "engine=sglang" "tp=${tp}" "dp=${dp}" "mark=${mark}" "isl=${isl}" "osl=${osl}" "concurrency=${conc}" "status=skipped_optional" "reason=optional scenario failed (rc=${bench_rc})" "detail_log=${detail_log}" + continue + fi + + # Mandatory scenario failed but not OOM: try to restart the server. + if restart_server "$tp" "$dp"; then + run_warmup 1024 128 || true + log "resuming ${config_label} after server restart" + continue + fi + + log "ERROR: ${config_label} server restart failed; skipping remaining scenarios" + append_scenario_record "$result_root" \ + "name=${sname}" \ + "config=$(jq -n --arg phase main --argjson c "$conc" --argjson i "$isl" --argjson o "$osl" --arg dataset random --argjson n "$num" '{phase: $phase, concurrency: $c, input_len: $i, output_len: $o, dataset: $dataset, num_prompts: $n}')" \ + "status=\"FAILED\"" \ + "note=\"scenario failed and server restart failed (rc=${bench_rc})\"" + record_skipped_csv "$skipped_csv" \ + "engine=sglang" "tp=${tp}" "dp=${dp}" "mark=${mark}" "isl=${isl}" "osl=${osl}" "concurrency=${conc}" "status=FAILED" "reason=scenario failed and server restart failed" "detail_log=${detail_log}" + skip_remaining_scenarios "$result_root" "$scenario_tsv" "$((i + 1))" "SKIPPED_RESTART_FAILED" "server restart failed" "$tp" "$dp" + break + done + + stop_server "$tp" "$dp" + + # Parse results. + log "parsing ${config_label} results" + "$PYTHON" "${SCRIPT_DIR}/../../scripts/common/parse_backend.py" "$result_root" --backend sglang \ + >> "${phase_log_dir}/parse.log" 2>&1 || { + log "WARNING: parser failed for ${config_label}; see ${phase_log_dir}/parse.log" + } + + log "===== ${config_label} DONE =====" +} + +# --------------------------------------------------------------------------- +# Main +# --------------------------------------------------------------------------- + +# Cleanup any leftovers. +for cfg in "${PARALLEL_CONFIGS[@]}"; do + read -r tp dp <<< "$cfg" + stop_server "$tp" "$dp" +done + +# Run each parallel configuration. +for cfg in "${PARALLEL_CONFIGS[@]}"; do + read -r tp dp <<< "$cfg" + run_parallel_config "$tp" "$dp" +done + +# Generate cross-configuration comparison. +log "generating comparison report" +"$PYTHON" "${SCRIPT_DIR}/compare.py" \ + --run-root "${RESULT_BASE}/${RUN_ID}" \ + --output "${RESULT_BASE}/${RUN_ID}/comparison.md" \ + >> "${log_dir_global}/compare.log" 2>&1 || { + log "WARNING: comparison script failed; see ${log_dir_global}/compare.log" + } + +log "all results saved to ${RESULT_BASE}/${RUN_ID}" diff --git a/experiments/dsv4_h20_sglang_tp_dp_matrix/run_sglang_in_container.sh b/experiments/dsv4_h20_sglang_tp_dp_matrix/run_sglang_in_container.sh new file mode 100755 index 0000000..69a08b0 --- /dev/null +++ b/experiments/dsv4_h20_sglang_tp_dp_matrix/run_sglang_in_container.sh @@ -0,0 +1,93 @@ +#!/usr/bin/env bash +# Run SGLang server inside the already-running benchmark container. +# Usage: run_sglang_in_container.sh +set -e + +TP="${1}" +DP="${2}" + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/config.env" + +RUNTIME_BASE="${RUNTIME_BASE:-${SCRIPT_DIR}/runtime}" +mkdir -p "${RUNTIME_BASE}/logs" + +PORT="${SGLANG_PORT:-30031}" +NAME="${EXPERIMENT}_sglang_container" +PID_FILE="${RUNTIME_BASE}/${EXPERIMENT}_sglang_tp${TP}_dp${DP}.pid" + +LOG="${RUNTIME_BASE}/logs/${EXPERIMENT}_sglang_in_container_tp${TP}_dp${DP}_$(date +%Y%m%d_%H%M%S).log" +rm -f "$PID_FILE" + +# Build server args. +SERVER_ARGS=( + serve + --model-path "$MODEL_PATH" + --trust-remote-code + --tp-size "$TP" + --moe-runner-backend "$MOE_RUNNER_BACKEND" + --context-length "$CONTEXT_LENGTH" + --max-running-requests "$MAX_RUNNING_REQUESTS" + --mem-fraction-static "$MEM_FRACTION_STATIC" + --host 0.0.0.0 + --port "$PORT" +) + +if [[ "$DP" -gt 1 ]]; then + SERVER_ARGS+=( + --dp-size "$DP" + ) +fi + +SERVER_ARGS_STR="${SERVER_ARGS[*]}" + +echo "=== Starting SGLang server inside container (TP=${TP}, DP=${DP}) ===" +echo "Container: $NAME" +echo "Command: sglang ${SERVER_ARGS_STR}" +echo "Log: $LOG" + +# Kill any existing sglang process inside container first. +docker exec "$NAME" pkill -9 -f "sglang serve" 2>/dev/null || true +sleep 3 + +# Start sglang serve inside the container in background. +# We use nohup so it survives after docker exec returns. +docker exec -d \ + -e CUDA_VISIBLE_DEVICES="${CUDA_VISIBLE_DEVICES}" \ + -e PYTHONUNBUFFERED=1 \ + "$NAME" \ + bash -c "nohup sglang ${SERVER_ARGS_STR} > /tmp/sglang_server.log 2>&1 &" + +# Wait for the server process to appear inside container. +sleep 2 +SERVER_PID=$(docker exec "$NAME" pgrep -f "sglang serve" | head -n 1 || true) +if [[ -z "$SERVER_PID" ]]; then + echo "ERROR: sglang server process not found inside container" + docker exec "$NAME" cat /tmp/sglang_server.log 2>/dev/null | tail -50 || true + exit 1 +fi + +echo "Server PID inside container: $SERVER_PID" +echo "$SERVER_PID" > "$PID_FILE" + +echo "Waiting for health on port ${PORT}..." +for i in $(seq 1 360); do + if curl --fail --silent --show-error --max-time 5 "http://127.0.0.1:${PORT}/health" >/dev/null 2>&1; then + echo "SGLang server is ready at http://127.0.0.1:${PORT}" + echo "Log: $LOG" + exit 0 + fi + # Check if server process is still alive. + if ! docker exec "$NAME" kill -0 "$SERVER_PID" 2>/dev/null; then + echo "ERROR: SGLang server exited early" + docker exec "$NAME" cat /tmp/sglang_server.log 2>/dev/null | tail -200 || true + exit 1 + fi + echo "Waiting... ($i/360)" + sleep 5 +done + +echo "ERROR: SGLang server not healthy after 360 retries (30 mins)" +docker exec "$NAME" cat /tmp/sglang_server.log 2>/dev/null | tail -200 || true +exit 1 diff --git a/experiments/dsv4_h20_sglang_tp_dp_matrix/start_sglang_container.sh b/experiments/dsv4_h20_sglang_tp_dp_matrix/start_sglang_container.sh new file mode 100755 index 0000000..2cf469c --- /dev/null +++ b/experiments/dsv4_h20_sglang_tp_dp_matrix/start_sglang_container.sh @@ -0,0 +1,70 @@ +#!/usr/bin/env bash +# Start a long-running Docker container for SGLang benchmarking. +# The container stays alive (sleep infinity) so we can docker exec into it +# to run different TP×DP configurations without losing DeepGEMM JIT cache. +# +# Usage: start_sglang_container.sh +set -e + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/config.env" + +RUNTIME_BASE="${RUNTIME_BASE:-${SCRIPT_DIR}/runtime}" +mkdir -p "${RUNTIME_BASE}/logs" "${RUNTIME_BASE}/tmp" + +IMAGE="${DOCKER_IMAGE:-lmsysorg/sglang:latest}" +PORT="${SGLANG_PORT:-30031}" +NAME="${EXPERIMENT}_sglang_container" +PID_FILE="${RUNTIME_BASE}/${EXPERIMENT}_container.pid" + +# Persistent cache directory on host for DeepGEMM JIT kernels. +CACHE_DIR="${RUNTIME_BASE}/cache" +mkdir -p "${CACHE_DIR}/deep_gemm" "${CACHE_DIR}/tvm-ffi" + +# Clean up any stale container. +docker rm -f "$NAME" >/dev/null 2>&1 || true + +LOG="${RUNTIME_BASE}/logs/${EXPERIMENT}_container_$(date +%Y%m%d_%H%M%S).log" +rm -f "$PID_FILE" + +echo "=== Starting SGLang benchmark container ===" +echo "Image: $IMAGE" +echo "Container name: $NAME" +echo "Host port: $PORT" +echo "Cache dir: $CACHE_DIR" +echo "Log: $LOG" + +# Start a long-running container. +# We override the entrypoint to sleep infinity so the container stays alive. +docker run -d \ + --name "$NAME" \ + --gpus all \ + --ipc host \ + --shm-size 16g \ + --entrypoint /bin/bash \ + -p "${PORT}:${PORT}" \ + -v "${MODEL_PATH}:${MODEL_PATH}:ro" \ + -v "${CACHE_DIR}/deep_gemm:/root/.cache/deep_gemm" \ + -v "${CACHE_DIR}/tvm-ffi:/root/.cache/tvm-ffi" \ + -v "${RUNTIME_BASE}/tmp:/tmp" \ + -e CUDA_VISIBLE_DEVICES="${CUDA_VISIBLE_DEVICES}" \ + -e PYTHONUNBUFFERED=1 \ + "$IMAGE" \ + -c "sleep infinity" \ + > "$LOG" 2>&1 + +# Wait a moment for container to be ready. +sleep 2 + +if ! docker ps --filter "name=$NAME" --format '{{.Names}}' | grep -q "$NAME"; then + echo "ERROR: Container failed to start" + tail -50 "$LOG" + exit 1 +fi + +# Record the container ID for later use. +docker inspect -f '{{.Id}}' "$NAME" > "$PID_FILE" + +echo "Container $NAME is running" +echo "Log: $LOG" diff --git a/experiments/dsv4_h20_sglang_tp_dp_matrix/start_sglang_docker.sh b/experiments/dsv4_h20_sglang_tp_dp_matrix/start_sglang_docker.sh new file mode 100755 index 0000000..5753e46 --- /dev/null +++ b/experiments/dsv4_h20_sglang_tp_dp_matrix/start_sglang_docker.sh @@ -0,0 +1,99 @@ +#!/usr/bin/env bash +# Start SGLang server in Docker for a given TPxDP configuration. +# Usage: start_sglang_docker.sh +# +# Uses the lmsysorg/sglang image and the same argument set as the +# bare-metal start script. The container is removed automatically on stop. +set -e + +TP="${1}" +DP="${2}" + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/config.env" + +RUNTIME_BASE="${RUNTIME_BASE:-${SCRIPT_DIR}/runtime}" +mkdir -p "${RUNTIME_BASE}/logs" "${RUNTIME_BASE}/tmp" + +IMAGE="${DOCKER_IMAGE:-lmsysorg/sglang:latest}" +PORT="${SGLANG_PORT:-30031}" +NAME="${EXPERIMENT}_sglang_tp${TP}_dp${DP}" +PID_FILE="${RUNTIME_BASE}/${EXPERIMENT}_sglang_tp${TP}_dp${DP}.pid" + +LOG="${RUNTIME_BASE}/logs/${EXPERIMENT}_sglang_docker_tp${TP}_dp${DP}_$(date +%Y%m%d_%H%M%S).log" +rm -f "$PID_FILE" + +# Clean up any stale container with the same name. +docker rm -f "$NAME" >/dev/null 2>&1 || true + +SERVER_ARGS=( + serve + --model-path "$MODEL_PATH" + --trust-remote-code + --tp-size "$TP" + --moe-runner-backend "$MOE_RUNNER_BACKEND" + --context-length "$CONTEXT_LENGTH" + --max-running-requests "$MAX_RUNNING_REQUESTS" + --mem-fraction-static "$MEM_FRACTION_STATIC" + --host 0.0.0.0 + --port "$PORT" +) + +if [[ "$DP" -gt 1 ]]; then + SERVER_ARGS+=( + --dp-size "$DP" + ) +fi + +SERVER_ARGS_STR="${SERVER_ARGS[*]}" + +echo "=== Starting SGLang server in Docker (TP=${TP}, DP=${DP}) ===" +echo "Image: $IMAGE" +echo "Model: $MODEL_PATH" +echo "Container name: $NAME" +echo "Host port: $PORT" +echo "Command: sglang ${SERVER_ARGS_STR}" +echo "Log: $LOG" + +# Run docker in the foreground so that killing the host process stops the +# container (the --rm flag ensures cleanup). nohup lets us background it and +# capture the host PID in the same way as the bare-metal start script. +nohup docker run --rm \ + --name "$NAME" \ + --gpus all \ + --ipc host \ + --shm-size 16g \ + --entrypoint sglang \ + -p "${PORT}:${PORT}" \ + -v "${MODEL_PATH}:${MODEL_PATH}:ro" \ + -v "${RUNTIME_BASE}/tmp:/tmp" \ + -e CUDA_VISIBLE_DEVICES="${CUDA_VISIBLE_DEVICES}" \ + -e PYTHONUNBUFFERED=1 \ + "$IMAGE" \ + "${SERVER_ARGS[@]}" \ + > "$LOG" 2>&1 & + +PID=$! +echo $PID > "$PID_FILE" +echo "PID: $PID" +echo "Waiting for health on port ${PORT}..." + +for i in $(seq 1 240); do + if curl --fail --silent --show-error --max-time 5 "http://127.0.0.1:${PORT}/health" >/dev/null 2>&1; then + echo "SGLang server is ready at http://127.0.0.1:${PORT}" + echo "Log: $LOG" + exit 0 + fi + if ! kill -0 $PID 2>/dev/null; then + echo "ERROR: Docker SGLang server exited early" + tail -200 "$LOG" + exit 1 + fi + echo "Waiting... ($i/240)" + sleep 5 +done + +echo "ERROR: Docker SGLang server not healthy after 240 retries" +tail -200 "$LOG" +exit 1 diff --git a/experiments/dsv4_h20_sglang_tp_dp_matrix/start_sglang_dp.sh b/experiments/dsv4_h20_sglang_tp_dp_matrix/start_sglang_dp.sh new file mode 100755 index 0000000..230bf88 --- /dev/null +++ b/experiments/dsv4_h20_sglang_tp_dp_matrix/start_sglang_dp.sh @@ -0,0 +1,86 @@ +#!/usr/bin/env bash +# Start SGLang server for a given TP×DP configuration. +# Usage: start_sglang_dp.sh +# +# By default this delegates to the Docker start script because the experiment +# is intended to run SGLang inside a container. Set USE_DOCKER=0 to use the +# local VENV_CLIENT environment instead. +set -e + +TP="${1}" +DP="${2}" + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/config.env" + +if [[ "${USE_DOCKER:-1}" == "1" ]]; then + exec "${SCRIPT_DIR}/start_sglang_docker.sh" "$@" +fi + +RUNTIME_BASE="${RUNTIME_BASE:-${SCRIPT_DIR}/runtime}" +mkdir -p "${RUNTIME_BASE}/logs" "${RUNTIME_BASE}/tmp" + +VENV="${VENV_CLIENT}" +export PATH="$VENV/bin:$PATH" +export PYTHONUNBUFFERED=1 +export TMPDIR="${RUNTIME_BASE}/tmp" +export CUDA_VISIBLE_DEVICES="${CUDA_VISIBLE_DEVICES}" + +LOG="${RUNTIME_BASE}/logs/${EXPERIMENT}_sglang_tp${TP}_dp${DP}_$(date +%Y%m%d_%H%M%S).log" +PID_FILE="${RUNTIME_BASE}/${EXPERIMENT}_sglang_tp${TP}_dp${DP}.pid" + +rm -f "$PID_FILE" + +SERVER_ARGS=( + sglang serve + --model-path "$MODEL_PATH" + --trust-remote-code + --tp-size "$TP" + --moe-runner-backend "$MOE_RUNNER_BACKEND" + --context-length "$CONTEXT_LENGTH" + --max-running-requests "$MAX_RUNNING_REQUESTS" + --mem-fraction-static "$MEM_FRACTION_STATIC" + --host 0.0.0.0 + --port "$SGLANG_PORT" +) + +if [[ "$DP" -gt 1 ]]; then + SERVER_ARGS+=( + --dp-size "$DP" + ) +fi + +SERVER_ARGS_STR="${SERVER_ARGS[*]}" + +echo "=== Starting SGLang server (TP=${TP}, DP=${DP}) ===" +echo "Model: $MODEL_PATH" +echo "Port: $SGLANG_PORT" +echo "Command: $SERVER_ARGS_STR" +echo "Log: $LOG" + +nohup "${SERVER_ARGS[@]}" > "$LOG" 2>&1 & + +PID=$! +echo $PID > "$PID_FILE" +echo "PID: $PID" +echo "Waiting for health on port ${SGLANG_PORT}..." + +for i in $(seq 1 240); do + if curl --fail --silent --show-error --max-time 5 "http://127.0.0.1:${SGLANG_PORT}/health" >/dev/null 2>&1; then + echo "SGLang server is ready at http://127.0.0.1:${SGLANG_PORT}" + echo "Log: $LOG" + exit 0 + fi + if ! kill -0 $PID 2>/dev/null; then + echo "ERROR: SGLang server exited early" + tail -200 "$LOG" + exit 1 + fi + echo "Waiting... ($i/240)" + sleep 5 +done + +echo "ERROR: SGLang server not healthy after 240 retries" +tail -200 "$LOG" +exit 1 diff --git a/experiments/dsv4_h20_vllm_tp_dp_matrix/adaptive_config.env b/experiments/dsv4_h20_vllm_tp_dp_matrix/adaptive_config.env new file mode 100644 index 0000000..1eda1f0 --- /dev/null +++ b/experiments/dsv4_h20_vllm_tp_dp_matrix/adaptive_config.env @@ -0,0 +1,49 @@ +# Adaptive concurrency search settings. +# +# For each fixed (TP, DP, ISL, OSL), probe: +# C = start, start * multiplier, ... up to max +# and stop after Total TPS has less than TPS_MIN_GAIN_PCT meaningful growth for +# PLATEAU_PATIENCE consecutive points. + +SEARCH_START_CONCURRENCY="${SEARCH_START_CONCURRENCY:-1}" +SEARCH_MAX_CONCURRENCY="${SEARCH_MAX_CONCURRENCY:-512}" +SEARCH_MULTIPLIER="${SEARCH_MULTIPLIER:-2}" +NUM_PROMPTS_MULTIPLIER="${NUM_PROMPTS_MULTIPLIER:-5}" + +# A gain below 2% is treated as throughput saturation. Two consecutive +# low-gain points prevent one noisy measurement from stopping the search. +TPS_MIN_GAIN_PCT="${TPS_MIN_GAIN_PCT:-2.0}" +PLATEAU_PATIENCE="${PLATEAU_PATIENCE:-2}" + +# Keep the same random workload semantics as the fixed matrix baseline. +# DATASET_PATH must contain at least SEARCH_MAX_CONCURRENCY times +# NUM_PROMPTS_MULTIPLIER valid two-turn conversations. Set this explicitly to +# random-ids to use generated token IDs without a ShareGPT seed dataset. +BENCH_DATASET_NAME="${BENCH_DATASET_NAME:-random}" +# SGLang interprets 0.0 as Uniform[1, requested_len]. Use 1.0 for fixed +# ISL/OSL points; lower values intentionally benchmark a length distribution. +RANDOM_RANGE_RATIO="${RANDOM_RANGE_RATIO:-1.0}" +# Before each measured point, warm up with the same concurrency so lazy kernel +# compilation and CUDA graph capture are excluded from TTFT/TPS. 0 means no +# cap; set a positive cap only when very high-concurrency warmup is impractical. +BENCH_WARMUP_MAX_REQUESTS="${BENCH_WARMUP_MAX_REQUESTS:-0}" + +# Reject a point if the completed request count or actual token lengths do not +# match the requested workload. +INPUT_LENGTH_TOLERANCE_PCT="${INPUT_LENGTH_TOLERANCE_PCT:-5.0}" +OUTPUT_LENGTH_TOLERANCE_PCT="${OUTPUT_LENGTH_TOLERANCE_PCT:-10.0}" + +MAX_POINT_RETRIES="${MAX_POINT_RETRIES:-1}" +SERVER_RESTART_COOLDOWN_S="${SERVER_RESTART_COOLDOWN_S:-10}" +SCENARIO_TIMEOUT_S="${SCENARIO_TIMEOUT_S:-1800}" +GPU_MEM_SAMPLE_INTERVAL_S="${GPU_MEM_SAMPLE_INTERVAL_S:-1}" + +# Optional space-separated filters, useful for smoke tests: +# TP_LIST="8" ISL_LIST="1024" OSL_LIST="128" +TP_LIST="${TP_LIST:-}" +ISL_LIST="${ISL_LIST:-}" +OSL_LIST="${OSL_LIST:-}" + +DRY_RUN="${DRY_RUN:-0}" +# Counts ISL/OSL shapes per TP/DP config, not individual concurrency probes. +GRID_LIMIT="${GRID_LIMIT:-0}" diff --git a/experiments/dsv4_h20_vllm_tp_dp_matrix/compare.py b/experiments/dsv4_h20_vllm_tp_dp_matrix/compare.py new file mode 100644 index 0000000..002fa0c --- /dev/null +++ b/experiments/dsv4_h20_vllm_tp_dp_matrix/compare.py @@ -0,0 +1,163 @@ +#!/usr/bin/env python3 +"""Cross TP×DP configuration comparison for dsv4_h20_vllm_tp_dp_matrix. + +Usage: + python3 compare.py --run-root results/ [--output comparison.md] +""" +import argparse +import json +import re +from collections import defaultdict +from pathlib import Path + + +def load_result(result_root: Path) -> dict: + path = result_root / "results.json" + with open(path, "r", encoding="utf-8") as f: + return json.load(f) + + +def slo_status(ttft_p95_ms: float, tpot_mean_ms: float, + ttft_limit_ms: float = 3000.0, tpot_limit_ms: float = 50.0) -> str: + ttft_ok = ttft_p95_ms < ttft_limit_ms + tpot_ok = tpot_mean_ms < tpot_limit_ms + if ttft_ok and tpot_ok: + return "PASS" + if ttft_ok or tpot_ok: + return "PARTIAL" + return "FAIL" + + +def gpu_memory_str(gpu: dict | None) -> str: + if not gpu: + return "-" + peak = gpu.get("peak_used_mb", 0) + total = gpu.get("memory_total_mb", 0) + if total: + return f"{peak:.0f}/{total:.0f} ({100*peak/total:.1f}%)" + return f"{peak:.0f}" + + +def main(): + parser = argparse.ArgumentParser() + parser.add_argument("--run-root", type=Path, required=True) + parser.add_argument("-o", "--output", type=Path, default=Path("comparison.md")) + parser.add_argument("--ttft-limit", type=float, default=3000.0) + parser.add_argument("--tpot-limit", type=float, default=50.0) + args = parser.parse_args() + + # Discover configurations: tp*_dp* directories. + configs = [] + for subdir in sorted(args.run_root.iterdir()): + if not subdir.is_dir(): + continue + name = subdir.name + if not (name.startswith("tp") and "_dp" in name): + continue + results_json = subdir / "results.json" + if not results_json.exists(): + continue + configs.append((name, load_result(subdir))) + + if not configs: + print(f"No tp*_dp* results found under {args.run_root}") + return + + model = configs[0][1].get("metadata", {}).get("model", "unknown") + hardware = configs[0][1].get("metadata", {}).get("hardware", "unknown") + + # Group by scenario name. + by_scenario: dict[str, dict[str, dict]] = defaultdict(dict) + skipped: dict[str, dict[str, str]] = defaultdict(dict) + for label, data in configs: + for s in data.get("scenarios", []): + key = s["name"] + if s.get("status") == "skipped_oom": + skipped[key][label] = s.get("note", "skipped") + else: + by_scenario[key][label] = s + + with open(args.output, "w", encoding="utf-8") as f: + f.write(f"# vLLM TP×DP matrix comparison ({hardware})\n\n") + f.write("## Summary\n\n") + f.write(f"- Model: `{model}`\n") + f.write(f"- Hardware: {hardware}\n") + f.write("- Backend: vLLM (Docker)\n") + f.write("- Benchmark client: `sglang.bench_serving`\n") + f.write(f"- SLO reference: TTFT P95 < {args.ttft_limit}ms, TPOT mean < {args.tpot_limit}ms\n\n") + + # Configuration overview. + f.write("### Configurations\n\n") + f.write("| Config | TP | DP | GPUs/replica | Notes |\n") + f.write("|---|---:|---:|---:|---|\n") + for label, data in configs: + cfg = data.get("config", {}) + tp = cfg.get("tp", "?") + dp = cfg.get("dp", "?") + f.write(f"| {label} | {tp} | {dp} | {tp} | server args recorded per ISL in results.json |\n") + f.write("\n") + + # Side-by-side table. + f.write("## Side-by-side results\n\n") + headers = [ + "Scenario", "ISL", "OSL", "Config", "Conc", "Req/s", "OutTok/s", + "TTFT P95(ms)", "TTFT P99(ms)", "TPOT Mean(ms)", "TPOT P95(ms)", + "TPOT P99(ms)", "E2E P99(ms)", "Peak GPU mem", "SLO" + ] + f.write("| " + " | ".join(headers) + " |\n") + f.write("|" + "|".join(["---"] * len(headers)) + "|\n") + + for scenario_name in sorted(by_scenario.keys(), key=lambda x: tuple(map(int, re.findall(r"\d+", x)))): + _, isl, osl = re.findall(r"\d+", scenario_name) + # cfg_part not used; just for readability. + for label, data in configs: + s = by_scenario[scenario_name].get(label) + if s is None: + if scenario_name in skipped and label in skipped[scenario_name]: + note = skipped[scenario_name][label] + f.write(f"| {scenario_name} | {isl} | {osl} | {label} | - | - | - | - | - | - | - | - | - | - | {note} |\n") + continue + cfg = s["config"] + m = s["metrics"] + status = slo_status(m["ttft_ms"]["p95"], m["tpot_ms"]["mean"], args.ttft_limit, args.tpot_limit) + gpu = m.get("gpu_memory") + f.write( + f"| {scenario_name} | {isl} | {osl} | {label} | {cfg['concurrency']} | " + f"{m['request_throughput']:.2f} | {m['output_token_throughput']:.2f} | " + f"{m['ttft_ms']['p95']:.2f} | {m['ttft_ms']['p99']:.2f} | " + f"{m['tpot_ms']['mean']:.2f} | {m['tpot_ms']['p95']:.2f} | {m['tpot_ms']['p99']:.2f} | " + f"{m['e2e_ms']['p99']:.2f} | {gpu_memory_str(gpu)} | {status} |\n" + ) + + # Best throughput per ISL/OSL. + f.write("\n## Best throughput per (ISL, OSL)\n\n") + f.write("| ISL | OSL | Best Config | Concurrency | OutTok/s | TTFT P95(ms) | TPOT Mean(ms) | SLO |\n") + f.write("|---:|---:|---|---:|---:|---:|---:|---:|\n") + best_by_shape: dict[tuple[int, int], tuple[float, str, dict]] = {} + for scenario_name, backends in by_scenario.items(): + _, isl, osl = re.findall(r"\d+", scenario_name) + isl_i, osl_i = int(isl), int(osl) + for label, s in backends.items(): + m = s["metrics"] + out_tok = m["output_token_throughput"] + if (isl_i, osl_i) not in best_by_shape or out_tok > best_by_shape[(isl_i, osl_i)][0]: + best_by_shape[(isl_i, osl_i)] = (out_tok, label, s) + for (isl_i, osl_i), (out_tok, label, s) in sorted(best_by_shape.items()): + m = s["metrics"] + status = slo_status(m["ttft_ms"]["p95"], m["tpot_ms"]["mean"], args.ttft_limit, args.tpot_limit) + f.write( + f"| {isl_i} | {osl_i} | {label} | {s['config']['concurrency']} | " + f"{out_tok:.2f} | {m['ttft_ms']['p95']:.2f} | {m['tpot_ms']['mean']:.2f} | {status} |\n" + ) + + f.write("\n## Notes\n\n") + f.write("- SLO check uses TTFT P95 and TPOT mean.\n") + f.write("- A PARTIAL indicates one of the two metrics is out of target; FAIL indicates both are out.\n") + f.write("- `Peak GPU mem` shows peak used / total MB and utilization percentage.\n") + f.write("- Optional (P) combinations that failed are marked as skipped/OOM and do not break the run.\n") + + print(f"Wrote comparison to {args.output}") + + +if __name__ == "__main__": + main() diff --git a/experiments/dsv4_h20_vllm_tp_dp_matrix/config.env b/experiments/dsv4_h20_vllm_tp_dp_matrix/config.env new file mode 100644 index 0000000..b9ffb7c --- /dev/null +++ b/experiments/dsv4_h20_vllm_tp_dp_matrix/config.env @@ -0,0 +1,76 @@ +# TP×DP matrix experiment for DeepSeek-V4-Flash on H20 (8 GPUs) using vLLM. +# Tests vLLM with three parallel configurations: +# TP=2, DP=4 -> 2 GPUs per replica, 4 replicas +# TP=4, DP=2 -> 4 GPUs per replica, 2 replicas +# TP=8, DP=1 -> 8 GPUs, no data parallelism + +EXPERIMENT="dsv4_h20_vllm_tp_dp_matrix" +MODEL_NAME="DeepSeek-V4-Flash" +MODEL_PATH="/data1/hf_models/DeepSeek-V4-Flash" +SERVED_MODEL_NAME="deepseek-v4-flash" + +VLLM_PORT="${VLLM_PORT:-30030}" + +# Python interpreter for orchestration scripts (parse_backend.py, compare.py, etc.) +# and the benchmark client. Defaults to the system python3 if the sglang venv +# does not exist on the host. +VENV_CLIENT="${VENV_CLIENT:-/data1/yy/sskj/envs/sglang}" + +# Run the benchmark client natively (0) or inside Docker (1). +USE_DOCKER_CLIENT="${USE_DOCKER_CLIENT:-1}" + +export CUDA_VISIBLE_DEVICES="${CUDA_VISIBLE_DEVICES:-0,1,2,3,4,5,6,7}" + +# Runtime working directory for logs, pid files, and tmp. Defaults to a local +# directory under this experiment so the benchmark is self-contained. +RUNTIME_BASE="${RUNTIME_BASE:-${SCRIPT_DIR}/runtime}" + +# Parallel configurations to test. Format: "TP DP" +declare -a PARALLEL_CONFIGS=( + "2 4" + "4 2" + "8 1" +) + +# vLLM server settings +MAX_MODEL_LEN="${MAX_MODEL_LEN:-1048576}" +MAX_NUM_SEQS="${MAX_NUM_SEQS:-128}" +GPU_MEMORY_UTILIZATION="${GPU_MEMORY_UTILIZATION:-0.9}" +KV_CACHE_DTYPE="${KV_CACHE_DTYPE:-fp8}" +BLOCK_SIZE="${BLOCK_SIZE:-256}" + +# Deployment switch. 0 = native vllm venv, 1 = Docker. +USE_DOCKER="${USE_DOCKER:-1}" +DOCKER_IMAGE="${DOCKER_IMAGE:-vllm/vllm-openai:latest}" + +# Benchmark client Docker image. vLLM's image does not include sglang.bench_serving, +# so the client runs inside the SGLang image and targets the vLLM backend via the +# OpenAI-compatible API. +DOCKER_CLIENT_IMAGE="${DOCKER_CLIENT_IMAGE:-lmsysorg/sglang:latest}" + +# Dataset used by sglang.bench_serving --dataset-name random. +# The random sampler needs a ShareGPT-style JSON file locally; it falls back to +# downloading from HuggingFace, which usually fails on offline H20 nodes. +DATASET_PATH="${DATASET_PATH:-/data1/yy/sskj/datasets/ShareGPT_V3_unfiltered_cleaned_split.json}" + +# Matrix and concurrency rules are defined in matrix.json by default. +MATRIX_FILE="${MATRIX_FILE:-${SCRIPT_DIR:-.}/matrix.json}" +MATRIX_MODE="${MATRIX_MODE:-Y}" + +# Sampling density for concurrency. +# 0 = use the default heuristic in generate_scenarios.py (6-8 points). +# 2 = only test the low and high endpoints. +export CONCURRENCY_SAMPLES="${CONCURRENCY_SAMPLES:-2}" + +# Per-scenario timeout to avoid hangs (seconds). +SCENARIO_TIMEOUT_S="${SCENARIO_TIMEOUT_S:-1800}" + +# GPU memory sampling interval (seconds). +GPU_MEM_SAMPLE_INTERVAL_S="${GPU_MEM_SAMPLE_INTERVAL_S:-1}" + +# Dry-run mode: if 1, only log the server args and scenario plan without starting +# any server or sending requests. +DRY_RUN="${DRY_RUN:-0}" + +# Per-config scenario limit for quick smoke tests. 0 = run all generated scenarios. +GRID_LIMIT="${GRID_LIMIT:-0}" diff --git a/experiments/dsv4_h20_vllm_tp_dp_matrix/generate_scenarios.py b/experiments/dsv4_h20_vllm_tp_dp_matrix/generate_scenarios.py new file mode 100644 index 0000000..2bfb152 --- /dev/null +++ b/experiments/dsv4_h20_vllm_tp_dp_matrix/generate_scenarios.py @@ -0,0 +1,98 @@ +#!/usr/bin/env python3 +"""Generate the scenario list for the TP×DP matrix experiment. + +Reads matrix.json and prints TSV lines: + + mark input_len output_len concurrency num_prompts + +mark is one of Y/P/N. The caller (run_bench.sh) decides how to treat each. +""" +import argparse +import json +import math +import os +from pathlib import Path + + +def sample_concurrency(low: int, high: int, target: int) -> list[int]: + """Return only the low and high concurrency values in [low, high]. + + For the TP×DP matrix we only need the two endpoints of the concurrency + range (e.g. 1 and 128 for ISL=1024). The `target` argument is kept for + API compatibility but is ignored. + """ + assert 1 <= low <= high, f"invalid concurrency range: {low}-{high}" + if low == high: + return [low] + return [low, high] + + +def generate_scenarios(matrix_path: Path, mode: str, target_samples: int) -> list[dict]: + with open(matrix_path, "r", encoding="utf-8") as f: + data = json.load(f) + + matrix = data["matrix"] + concurrency_cfg = data["concurrency"] + + scenarios = [] + for isl_str in sorted(matrix.keys(), key=int): + osl_map = matrix[isl_str] + low = concurrency_cfg[isl_str]["low"] + high = concurrency_cfg[isl_str]["high"] + concurrencies = sample_concurrency(low, high, target_samples) + + for osl_str in sorted(osl_map.keys(), key=int): + mark = osl_map[osl_str] + if mode == "Y" and mark != "Y": + continue + if mode == "Y+P" and mark not in ("Y", "P"): + continue + # mode == "all" keeps everything, including N. + + for conc in concurrencies: + scenarios.append( + { + "mark": mark, + "input_len": int(isl_str), + "output_len": int(osl_str), + "concurrency": conc, + "num_prompts": conc * 5, + } + ) + + return scenarios + + +def main() -> None: + parser = argparse.ArgumentParser() + parser.add_argument("--matrix", type=Path, default=Path("matrix.json")) + parser.add_argument("--mode", choices=["Y", "Y+P", "all"], default=None, + help="Scenario selection mode. Defaults to matrix.mode.") + parser.add_argument("--target-samples", type=int, default=0, + help="Target number of concurrency samples. 0 = heuristic (6-8).") + args = parser.parse_args() + + with open(args.matrix, "r", encoding="utf-8") as f: + data = json.load(f) + + mode = args.mode if args.mode else data.get("mode", "Y+P") + + target_samples = args.target_samples + if target_samples <= 0: + env_samples = os.getenv("CONCURRENCY_SAMPLES", "0") + try: + target_samples = int(env_samples) + except ValueError: + target_samples = 0 + if target_samples <= 0: + target_samples = 7 + + scenarios = generate_scenarios(args.matrix, mode, target_samples) + + print("mark\tinput_len\toutput_len\tconcurrency\tnum_prompts") + for s in scenarios: + print(f"{s['mark']}\t{s['input_len']}\t{s['output_len']}\t{s['concurrency']}\t{s['num_prompts']}") + + +if __name__ == "__main__": + main() diff --git a/experiments/dsv4_h20_vllm_tp_dp_matrix/matrix.json b/experiments/dsv4_h20_vllm_tp_dp_matrix/matrix.json new file mode 100644 index 0000000..34e4847 --- /dev/null +++ b/experiments/dsv4_h20_vllm_tp_dp_matrix/matrix.json @@ -0,0 +1,81 @@ +{ + "comment": "ISL/OSL matrix for dsv4_h20_vllm_tp_dp_matrix. Y=must test, P=optional (record N on failure), N=skip.", + "mode": "Y", + "comment": "Only mandatory (Y) combinations are tested; 1M ISL is excluded per user request.", + "matrix": { + "1024": { + "128": "Y", + "256": "Y", + "512": "Y", + "1024": "Y", + "2048": "Y", + "4096": "Y" + }, + "4096": { + "128": "Y", + "256": "Y", + "512": "Y", + "1024": "Y", + "2048": "Y", + "4096": "Y" + }, + "16384": { + "128": "Y", + "256": "Y", + "512": "Y", + "1024": "Y", + "2048": "Y", + "4096": "P" + }, + "65536": { + "128": "Y", + "256": "Y", + "512": "Y", + "1024": "Y", + "2048": "P", + "4096": "N" + }, + "131072": { + "128": "Y", + "256": "Y", + "512": "Y", + "1024": "P", + "2048": "N", + "4096": "N" + }, + "262144": { + "128": "Y", + "256": "Y", + "512": "P", + "1024": "N", + "2048": "N", + "4096": "N" + }, + "524288": { + "128": "Y", + "256": "P", + "512": "N", + "1024": "N", + "2048": "N", + "4096": "N" + }, + "1048576": { + "128": "Y", + "256": "P", + "512": "N", + "1024": "N", + "2048": "N", + "4096": "N" + } + }, + "concurrency": { + "1024": { "low": 1, "high": 128 }, + "4096": { "low": 1, "high": 64 }, + "16384": { "low": 1, "high": 32 }, + "65536": { "low": 1, "high": 8 }, + "131072": { "low": 1, "high": 4 }, + "262144": { "low": 1, "high": 2 }, + "524288": { "low": 1, "high": 2 }, + "1048576": { "low": 1, "high": 2 } + } +} diff --git a/experiments/dsv4_h20_vllm_tp_dp_matrix/run_adaptive_concurrency.sh b/experiments/dsv4_h20_vllm_tp_dp_matrix/run_adaptive_concurrency.sh new file mode 100755 index 0000000..ca23f6a --- /dev/null +++ b/experiments/dsv4_h20_vllm_tp_dp_matrix/run_adaptive_concurrency.sh @@ -0,0 +1,164 @@ +#!/usr/bin/env bash +# Find the Total-TPS saturation concurrency for each vLLM TP/DP/ISL/OSL shape. +set -Eeuo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +EXPERIMENT_NAME="$(basename "$SCRIPT_DIR")" + +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/lib.sh" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/platform.sh" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/config.env" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/adaptive_config.env" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/adaptive_bench_lib.sh" + +ENGINE="vllm" +ENGINE_PORT="$VLLM_PORT" +RESULT_BASE="${RESULT_BASE:-${SCRIPT_DIR}/adaptive_results}" +ACTIVE_ENGINE_SERVER_LOG="" + +if [[ -x "${VENV_CLIENT}/bin/python" ]]; then + PYTHON="${VENV_CLIENT}/bin/python" +else + PYTHON="$(command -v python3)" +fi + +DOCKER_IMAGE="${DOCKER_IMAGE:-vllm/vllm-openai:latest}" +DOCKER_CLIENT_IMAGE="${DOCKER_CLIENT_IMAGE:-lmsysorg/sglang:latest}" + +engine_is_healthy() { + curl --fail --silent --show-error --max-time 5 \ + "http://127.0.0.1:${ENGINE_PORT}/health" >/dev/null 2>&1 +} + +engine_stop_server() { + local tp="$1" + local dp="$2" + local pid_file="${RUNTIME_BASE}/${EXPERIMENT}_vllm_tp${tp}_dp${dp}.pid" + + if [[ -f "$pid_file" ]]; then + local pid + pid="$(cat "$pid_file")" + if kill -0 "$pid" 2>/dev/null; then + log "stopping vllm server pid=${pid} tp=${tp} dp=${dp}" + kill "$pid" 2>/dev/null || true + sleep 5 + kill -9 "$pid" 2>/dev/null || true + fi + rm -f "$pid_file" + fi + docker rm -f "${EXPERIMENT}_vllm_tp${tp}_dp${dp}" >/dev/null 2>&1 || true + ACTIVE_ENGINE_SERVER_LOG="" + sleep 2 +} + +engine_build_server_args() { + local tp="$1" + local dp="$2" + local -a args=( + vllm serve "$MODEL_PATH" + --trust-remote-code + --kv-cache-dtype "$KV_CACHE_DTYPE" + --block-size "$BLOCK_SIZE" + --tensor-parallel-size "$tp" + --gpu-memory-utilization "$GPU_MEMORY_UTILIZATION" + --max-model-len "$MAX_MODEL_LEN" + --max-num-seqs "$MAX_NUM_SEQS" + --no-enable-flashinfer-autotune + --host 0.0.0.0 + --port "$ENGINE_PORT" + ) + if (( dp > 1 )); then + args+=(--data-parallel-size "$dp") + fi + printf '%q ' "${args[@]}" +} + +engine_start_server() { + local tp="$1" + local dp="$2" + local outer_log="${ADAPTIVE_LOG_DIR}/vllm_tp${tp}_dp${dp}.server.outer.log" + log "starting vllm server tp=${tp} dp=${dp}" + bash "${SCRIPT_DIR}/start_vllm_dp.sh" "$tp" "$dp" >> "$outer_log" 2>&1 + if ! engine_is_healthy; then + log "ERROR: vllm health check failed tp=${tp} dp=${dp}" + return 1 + fi + ACTIVE_ENGINE_SERVER_LOG="$( + find "${RUNTIME_BASE}/logs" -maxdepth 1 -type f \ + -name "${EXPERIMENT}_vllm*tp${tp}_dp${dp}_*.log" \ + -printf '%T@ %p\n' 2>/dev/null | sort -nr | head -n 1 | cut -d' ' -f2- + )" + log "vllm server healthy tp=${tp} dp=${dp} log=${ACTIVE_ENGINE_SERVER_LOG:-unknown}" +} + +engine_detect_oom() { + local detail_log="$1" + local tp="$2" + local dp="$3" + local pattern='CUDA out of memory|torch\.OutOfMemoryError|OutOfMemory|out of memory|OOM|RESOURCE_EXHAUSTED|Failed to allocate memory' + local outer_log="${ADAPTIVE_LOG_DIR}/vllm_tp${tp}_dp${dp}.server.outer.log" + local -a logs=("$detail_log" "$outer_log") + if [[ -n "$ACTIVE_ENGINE_SERVER_LOG" ]]; then + logs+=("$ACTIVE_ENGINE_SERVER_LOG") + fi + grep -Eiq "$pattern" "${logs[@]}" 2>/dev/null +} + +engine_run_bench() { + local isl="$1" + local osl="$2" + local concurrency="$3" + local num_prompts="$4" + local output_file="$5" + local warmup_requests + warmup_requests="$(adaptive_warmup_request_count "$concurrency")" + local -a bench_args=( + --backend vllm + --host 127.0.0.1 + --port "$ENGINE_PORT" + --dataset-name "$BENCH_DATASET_NAME" + --random-input-len "$isl" + --random-output-len "$osl" + --random-range-ratio "$RANDOM_RANGE_RATIO" + --num-prompts "$num_prompts" + --max-concurrency "$concurrency" + --request-rate 10000 + --warmup-requests "$warmup_requests" + --output-file "$output_file" + --output-details + --disable-tqdm + ) + if [[ "$BENCH_DATASET_NAME" == "random" ]]; then + bench_args+=(--dataset-path "$DATASET_PATH") + else + bench_args+=(--tokenize-prompt) + fi + + if [[ "$USE_DOCKER_CLIENT" == "1" ]]; then + local -a volume_args=(-v "${MODEL_PATH}:${MODEL_PATH}:ro" -v "${RESULT_BASE}:${RESULT_BASE}") + if [[ "$BENCH_DATASET_NAME" == "random" ]]; then + volume_args+=(-v "${DATASET_PATH}:${DATASET_PATH}:ro") + fi + docker run --rm \ + --network host \ + "${volume_args[@]}" \ + -e PYTHONUNBUFFERED=1 \ + -e HF_HUB_OFFLINE=1 \ + -e TRANSFORMERS_OFFLINE=1 \ + "$DOCKER_CLIENT_IMAGE" \ + python -m sglang.bench_serving "${bench_args[@]}" + else + "$PYTHON" -m sglang.bench_serving "${bench_args[@]}" + fi +} + +export -f engine_run_bench +export ENGINE_PORT MODEL_PATH RESULT_BASE DOCKER_CLIENT_IMAGE USE_DOCKER_CLIENT +export BENCH_DATASET_NAME DATASET_PATH RANDOM_RANGE_RATIO BENCH_WARMUP_MAX_REQUESTS PYTHON + +adaptive_main "$@" diff --git a/experiments/dsv4_h20_vllm_tp_dp_matrix/run_adaptive_concurrency_add16.sh b/experiments/dsv4_h20_vllm_tp_dp_matrix/run_adaptive_concurrency_add16.sh new file mode 100755 index 0000000..482886d --- /dev/null +++ b/experiments/dsv4_h20_vllm_tp_dp_matrix/run_adaptive_concurrency_add16.sh @@ -0,0 +1,167 @@ +#!/usr/bin/env bash +# Find the Total-TPS saturation concurrency for each vLLM TP/DP/ISL/OSL shape. +set -Eeuo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +EXPERIMENT_NAME="$(basename "$SCRIPT_DIR")" + +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/lib.sh" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/platform.sh" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/config.env" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/adaptive_config.env" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/adaptive_bench_lib.sh" + +ENGINE="vllm" +ENGINE_PORT="$VLLM_PORT" +RESULT_BASE="${RESULT_BASE:-${SCRIPT_DIR}/adaptive_results}" +ACTIVE_ENGINE_SERVER_LOG="" + +if [[ -x "${VENV_CLIENT}/bin/python" ]]; then + PYTHON="${VENV_CLIENT}/bin/python" +else + PYTHON="$(command -v python3)" +fi + +DOCKER_IMAGE="${DOCKER_IMAGE:-vllm/vllm-openai:latest}" +DOCKER_CLIENT_IMAGE="${DOCKER_CLIENT_IMAGE:-lmsysorg/sglang:latest}" + +engine_is_healthy() { + curl --fail --silent --show-error --max-time 5 \ + "http://127.0.0.1:${ENGINE_PORT}/health" >/dev/null 2>&1 +} + +engine_stop_server() { + local tp="$1" + local dp="$2" + local pid_file="${RUNTIME_BASE}/${EXPERIMENT}_vllm_tp${tp}_dp${dp}.pid" + + if [[ -f "$pid_file" ]]; then + local pid + pid="$(cat "$pid_file")" + if kill -0 "$pid" 2>/dev/null; then + log "stopping vllm server pid=${pid} tp=${tp} dp=${dp}" + kill "$pid" 2>/dev/null || true + sleep 5 + kill -9 "$pid" 2>/dev/null || true + fi + rm -f "$pid_file" + fi + docker rm -f "${EXPERIMENT}_vllm_tp${tp}_dp${dp}" >/dev/null 2>&1 || true + ACTIVE_ENGINE_SERVER_LOG="" + sleep 2 +} + +engine_build_server_args() { + local tp="$1" + local dp="$2" + local -a args=( + vllm serve "$MODEL_PATH" + --trust-remote-code + --kv-cache-dtype "$KV_CACHE_DTYPE" + --block-size "$BLOCK_SIZE" + --tensor-parallel-size "$tp" + --gpu-memory-utilization "$GPU_MEMORY_UTILIZATION" + --max-model-len "$MAX_MODEL_LEN" + --max-num-seqs "$MAX_NUM_SEQS" + --no-enable-flashinfer-autotune + --host 0.0.0.0 + --port "$ENGINE_PORT" + ) + if (( dp > 1 )); then + args+=(--data-parallel-size "$dp") + fi + printf '%q ' "${args[@]}" +} + +engine_start_server() { + local tp="$1" + local dp="$2" + local outer_log="${ADAPTIVE_LOG_DIR}/vllm_tp${tp}_dp${dp}.server.outer.log" + log "starting vllm server tp=${tp} dp=${dp}" + bash "${SCRIPT_DIR}/start_vllm_dp.sh" "$tp" "$dp" >> "$outer_log" 2>&1 + if ! engine_is_healthy; then + log "ERROR: vllm health check failed tp=${tp} dp=${dp}" + return 1 + fi + ACTIVE_ENGINE_SERVER_LOG="$( + find "${RUNTIME_BASE}/logs" -maxdepth 1 -type f \ + -name "${EXPERIMENT}_vllm*tp${tp}_dp${dp}_*.log" \ + -printf '%T@ %p\n' 2>/dev/null | sort -nr | head -n 1 | cut -d' ' -f2- + )" + log "vllm server healthy tp=${tp} dp=${dp} log=${ACTIVE_ENGINE_SERVER_LOG:-unknown}" +} + +engine_detect_oom() { + local detail_log="$1" + local tp="$2" + local dp="$3" + local pattern='CUDA out of memory|torch\.OutOfMemoryError|OutOfMemory|out of memory|OOM|RESOURCE_EXHAUSTED|Failed to allocate memory' + local outer_log="${ADAPTIVE_LOG_DIR}/vllm_tp${tp}_dp${dp}.server.outer.log" + local -a logs=("$detail_log" "$outer_log") + if [[ -n "$ACTIVE_ENGINE_SERVER_LOG" ]]; then + logs+=("$ACTIVE_ENGINE_SERVER_LOG") + fi + grep -Eiq "$pattern" "${logs[@]}" 2>/dev/null +} + +engine_run_bench() { + local isl="$1" + local osl="$2" + local concurrency="$3" + local num_prompts="$4" + local output_file="$5" + local warmup_requests + warmup_requests="$(adaptive_warmup_request_count "$concurrency")" + local -a bench_args=( + --backend vllm + --host 127.0.0.1 + --port "$ENGINE_PORT" + --dataset-name "$BENCH_DATASET_NAME" + --random-input-len "$isl" + --random-output-len "$osl" + --random-range-ratio "$RANDOM_RANGE_RATIO" + --num-prompts "$num_prompts" + --max-concurrency "$concurrency" + --request-rate 10000 + --warmup-requests "$warmup_requests" + --output-file "$output_file" + --output-details + --disable-tqdm + ) + if [[ "$BENCH_DATASET_NAME" == "random" ]]; then + bench_args+=(--dataset-path "$DATASET_PATH") + else + bench_args+=(--tokenize-prompt) + fi + + if [[ "$USE_DOCKER_CLIENT" == "1" ]]; then + local -a volume_args=(-v "${MODEL_PATH}:${MODEL_PATH}:ro" -v "${RESULT_BASE}:${RESULT_BASE}") + if [[ "$BENCH_DATASET_NAME" == "random" ]]; then + volume_args+=(-v "${DATASET_PATH}:${DATASET_PATH}:ro") + fi + docker run --rm \ + --network host \ + "${volume_args[@]}" \ + -e PYTHONUNBUFFERED=1 \ + -e HF_HUB_OFFLINE=1 \ + -e TRANSFORMERS_OFFLINE=1 \ + "$DOCKER_CLIENT_IMAGE" \ + python -m sglang.bench_serving "${bench_args[@]}" + else + "$PYTHON" -m sglang.bench_serving "${bench_args[@]}" + fi +} + +export -f engine_run_bench +export ENGINE_PORT MODEL_PATH RESULT_BASE DOCKER_CLIENT_IMAGE USE_DOCKER_CLIENT +export BENCH_DATASET_NAME DATASET_PATH RANDOM_RANGE_RATIO BENCH_WARMUP_MAX_REQUESTS PYTHON + +export SEARCH_START_CONCURRENCY=16 +export SEARCH_ADDEND=16 + +adaptive_main "$@" diff --git a/experiments/dsv4_h20_vllm_tp_dp_matrix/run_bench.sh b/experiments/dsv4_h20_vllm_tp_dp_matrix/run_bench.sh new file mode 100755 index 0000000..7bb82c8 --- /dev/null +++ b/experiments/dsv4_h20_vllm_tp_dp_matrix/run_bench.sh @@ -0,0 +1,548 @@ +#!/usr/bin/env bash +# TP×DP matrix benchmark for DeepSeek-V4-Flash on vLLM (Docker). +set -Eeuo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +EXPERIMENT_NAME="$(basename "$SCRIPT_DIR")" + +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/lib.sh" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/../../scripts/common/platform.sh" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/config.env" + +# Export variables used inside functions that are called via bash -c subshells. +export DATASET_PATH MODEL_PATH RESULT_BASE \ + DOCKER_IMAGE DOCKER_CLIENT_IMAGE USE_DOCKER_CLIENT + +RUN_ID="${RUN_ID:-$(date '+%Y%m%d-%H%M%S')}" +RESULT_BASE="${SCRIPT_DIR}/results" +MATRIX_FILE="${MATRIX_FILE:-${SCRIPT_DIR}/matrix.json}" +MATRIX_MODE="${MATRIX_MODE:-Y}" +SCENARIO_TIMEOUT_S="${SCENARIO_TIMEOUT_S:-1800}" +GPU_MEM_SAMPLE_INTERVAL_S="${GPU_MEM_SAMPLE_INTERVAL_S:-1}" +DRY_RUN="${DRY_RUN:-0}" +GRID_LIMIT="${GRID_LIMIT:-0}" + +if [[ -x "${VENV_CLIENT}/bin/python" ]]; then + PYTHON="${VENV_CLIENT}/bin/python" +else + PYTHON="$(command -v python3)" +fi + +DOCKER_IMAGE="${DOCKER_IMAGE:-vllm/vllm-openai:latest}" +DOCKER_CLIENT_IMAGE="${DOCKER_CLIENT_IMAGE:-lmsysorg/sglang:latest}" + +log_dir_global="${RESULT_BASE}/${RUN_ID}/logs" +mkdir -p "$log_dir_global" +log_init "${log_dir_global}/orchestrator.log" + +log "experiment=${EXPERIMENT_NAME} run_id=${RUN_ID} platform=${PLATFORM} hardware=${HARDWARE}" +log "matrix_mode=${MATRIX_MODE} matrix_file=${MATRIX_FILE} dry_run=${DRY_RUN} grid_limit=${GRID_LIMIT}" + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + +is_server_healthy() { + curl --fail --silent --show-error --max-time 5 "http://127.0.0.1:${VLLM_PORT}/health" >/dev/null 2>&1 +} + +stop_server() { + local tp="$1" + local dp="$2" + + local pid_file="${RUNTIME_BASE}/${EXPERIMENT}_vllm_tp${tp}_dp${dp}.pid" + if [[ -f "$pid_file" ]]; then + local pid + pid="$(cat "$pid_file")" + if kill -0 "$pid" 2>/dev/null; then + log "stopping vllm server pid=${pid} (tp=${tp}, dp=${dp})" + kill "$pid" 2>/dev/null || true + sleep 5 + kill -9 "$pid" 2>/dev/null || true + fi + rm -f "$pid_file" + fi + + # Fallback: remove any Docker container started by this experiment. + docker rm -f "${EXPERIMENT}_vllm_tp${tp}_dp${dp}" >/dev/null 2>&1 || true + + # Fallback: kill any vllm serve processes for this model. + pkill -9 -f "vllm serve.*${MODEL_NAME}" 2>/dev/null || true + pkill -9 -f "vllm serve.*${MODEL_PATH}" 2>/dev/null || true + sleep 2 +} + +build_server_args() { + local tp="$1" + local dp="$2" + + local args=( + "vllm serve" "$MODEL_PATH" + --trust-remote-code + --kv-cache-dtype "$KV_CACHE_DTYPE" + --block-size "$BLOCK_SIZE" + --tensor-parallel-size "$tp" + --gpu-memory-utilization "$GPU_MEMORY_UTILIZATION" + --max-model-len "$MAX_MODEL_LEN" + --max-num-seqs "$MAX_NUM_SEQS" + --no-enable-flashinfer-autotune + --host 0.0.0.0 + --port "$VLLM_PORT" + ) + if [[ "$dp" -gt 1 ]]; then + args+=( + --data-parallel-size "$dp" + ) + fi + printf '%s ' "${args[@]}" +} + +start_server() { + local tp="$1" + local dp="$2" + + log "starting vllm server tp=${tp} dp=${dp}" + bash "${SCRIPT_DIR}/start_vllm_dp.sh" "$tp" "$dp" \ + >> "${log_dir_global}/vllm_tp${tp}_dp${dp}.server.outer.log" 2>&1 + + if ! is_server_healthy; then + log "error: vllm server tp=${tp} dp=${dp} failed health check on port ${VLLM_PORT}" + return 1 + fi + log "vllm server tp=${tp} dp=${dp} is healthy on port ${VLLM_PORT}" +} + +restart_server() { + local tp="$1" + local dp="$2" + log "restarting vllm server tp=${tp} dp=${dp} after non-OOM failure" + stop_server "$tp" "$dp" + sleep 10 + start_server "$tp" "$dp" +} + +run_bench_serving() { + # Run sglang.bench_serving either natively or inside the SGLang Docker image. + if [[ "${USE_DOCKER_CLIENT:-1}" == "1" ]]; then + local vol_args=() + vol_args+=("-v" "${MODEL_PATH}:${MODEL_PATH}:ro") + if [[ -n "${DATASET_PATH:-}" && -f "${DATASET_PATH}" ]]; then + vol_args+=("-v" "${DATASET_PATH}:${DATASET_PATH}:ro") + fi + vol_args+=("-v" "${RESULT_BASE}:${RESULT_BASE}") + docker run --rm \ + --network host \ + "${vol_args[@]}" \ + -e PYTHONUNBUFFERED=1 \ + -e HF_HUB_OFFLINE=1 \ + -e TRANSFORMERS_OFFLINE=1 \ + -e HF_DATASETS_OFFLINE=1 \ + "${DOCKER_CLIENT_IMAGE}" \ + python -m sglang.bench_serving "$@" + else + "$PYTHON" -m sglang.bench_serving "$@" + fi +} +export -f run_bench_serving + +run_warmup() { + local input_len="$1" + local output_len="$2" + + log "warming up (input=${input_len}, output=${output_len}, num=1)" + bash -c ' + run_bench_serving \ + --backend vllm \ + --host 127.0.0.1 \ + --port "'"$VLLM_PORT"'" \ + --dataset-name random \ + --dataset-path "'"$DATASET_PATH"'" \ + --random-input-len "'"$input_len"'" \ + --random-output-len "'"$output_len"'" \ + --num-prompts 1 \ + --max-concurrency 1 \ + --request-rate 10000 \ + --output-file /dev/null \ + --output-details \ + >> "'"${log_dir_global}/warmup.log"'" 2>&1 + ' + log "warmup completed" +} + +scenario_already_completed() { + local output_file="$1" + local expected="$2" + [[ -s "$output_file" ]] || return 1 + local completed + completed="$("$PYTHON" -c " +import json, sys +path = sys.argv[1] +try: + with open(path, 'r', encoding='utf-8') as f: + for line in f: + line = line.strip() + if line: + data = json.loads(line) + print(data.get('completed', 0)) + break +except Exception: + print(0) +" "$output_file")" + [[ "${completed:-0}" -ge "$expected" ]] +} + +scenario_already_processed() { + local result_root="$1" + local scenario_name="$2" + local json_path="${result_root}/results.json" + [[ -f "$json_path" ]] || return 1 + "$PYTHON" -c " +import json, sys +path, name = sys.argv[1], sys.argv[2] +try: + with open(path, 'r', encoding='utf-8') as f: + data = json.load(f) + for s in data.get('scenarios', []): + if s.get('name') == name: + if s.get('status') or s.get('metrics', {}).get('success', 0) > 0: + sys.exit(0) +except Exception: + pass +sys.exit(1) +" "$json_path" "$scenario_name" +} + +detect_oom() { + local detail_log="$1" + local server_outer_log="$2" + local pattern='CUDA out of memory|torch\.OutOfMemoryError|OutOfMemory|out of memory|OOM|RESOURCE_EXHAUSTED|Failed to allocate memory' + if grep -Eiq "$pattern" "$detail_log" "$server_outer_log" 2>/dev/null; then + return 0 + fi + return 1 +} + +start_gpu_monitor() { + local csv_path="$1" + mkdir -p "$(dirname "$csv_path")" + nvidia-smi \ + --query-gpu=timestamp,index,memory.used,memory.total,utilization.gpu \ + --format=csv \ + -l "$GPU_MEM_SAMPLE_INTERVAL_S" \ + > "$csv_path" 2>/dev/null & + echo $! +} + +stop_gpu_monitor() { + local pid="$1" + if kill -0 "$pid" 2>/dev/null; then + kill "$pid" 2>/dev/null || true + sleep 1 + kill -9 "$pid" 2>/dev/null || true + fi +} + +append_scenario_record() { + local result_root="$1" + local json_path="$result_root/results.json" + shift + local scenario_json + scenario_json="$("$PYTHON" -c " +import json, sys +pairs = [a.split('=', 1) for a in sys.argv[1:]] +d = {} +for k, v in pairs: + try: + d[k] = json.loads(v) + except json.JSONDecodeError: + d[k] = v +print(json.dumps(d, ensure_ascii=False)) +" "$@")" + PYTHON="$PYTHON" append_scenario_to_json "$json_path" "$scenario_json" +} + +record_skipped_csv() { + local csv_path="$1" + shift + # Args: key=value + local row + row="$("$PYTHON" -c " +import csv, json, sys, io +pairs = [a.split('=', 1) for a in sys.argv[1:]] +d = {} +for k, v in pairs: + try: + d[k] = json.loads(v) + except json.JSONDecodeError: + d[k] = v +buf = io.StringIO() +writer = csv.DictWriter(buf, fieldnames=['engine','tp','dp','mark','isl','osl','concurrency','status','reason','detail_log'], extrasaction='ignore') +writer.writerow(d) +print(buf.getvalue().strip()) +" "$@")" + echo "$row" >> "$csv_path" +} + +skip_remaining_scenarios() { + local result_root="$1" + local scenario_tsv="$2" + local start_index="$3" + local status="$4" + local reason="$5" + local tp="$6" + local dp="$7" + local skipped_csv="${RESULT_BASE}/${RUN_ID}/skipped_after_oom.csv" + + local i=0 + tail -n +2 "$scenario_tsv" | while IFS=$'\t' read -r mark isl osl conc num; do + if (( i < start_index )); then + i=$((i + 1)) + continue + fi + i=$((i + 1)) + local sname="c${conc}_i${isl}_o${osl}" + if scenario_already_processed "$result_root" "$sname"; then + continue + fi + append_scenario_record "$result_root" \ + "name=${sname}" \ + "config=$(jq -n --arg phase main --argjson c "$conc" --argjson i "$isl" --argjson o "$osl" --arg dataset random --argjson n "$num" '{phase: $phase, concurrency: $c, input_len: $i, output_len: $o, dataset: $dataset, num_prompts: $n}')" \ + "status=\"${status}\"" \ + "note=\"${reason}\"" + record_skipped_csv "$skipped_csv" \ + "engine=vllm" "tp=${tp}" "dp=${dp}" "mark=${mark}" "isl=${isl}" "osl=${osl}" "concurrency=${conc}" "status=${status}" "reason=${reason}" + done +} + +# --------------------------------------------------------------------------- +# Per-configuration runner +# --------------------------------------------------------------------------- + +run_parallel_config() { + local tp="$1" + local dp="$2" + local config_label="tp${tp}_dp${dp}" + local result_root="${RESULT_BASE}/${RUN_ID}/${config_label}" + local raw_dir="${result_root}/raw_outputs" + local gpu_log_dir="${result_root}/gpu_logs" + local phase_log_dir="${result_root}/logs" + mkdir -p "$raw_dir" "$gpu_log_dir" "$phase_log_dir" + + log "===== ${config_label} START =====" + + # Generate scenario list for this config. + local scenario_tsv="${result_root}/scenarios.tsv" + "$PYTHON" "${SCRIPT_DIR}/generate_scenarios.py" \ + --matrix "$MATRIX_FILE" \ + --mode "$MATRIX_MODE" \ + > "$scenario_tsv" + local total_scenarios + total_scenarios="$(tail -n +2 "$scenario_tsv" | wc -l)" + log "generated ${total_scenarios} scenarios for ${config_label}" + + # Write metadata. + ensure_result_root "$result_root" + write_metadata_json \ + "${result_root}/results.json" \ + "${EXPERIMENT_NAME}_${config_label}" \ + "$RUN_ID" \ + "$MODEL_PATH" \ + "vllm" \ + "vllm" \ + "$HARDWARE" \ + "$ACCELERATOR" \ + "$CHIP" \ + "experiments/${EXPERIMENT_NAME}/run_bench.sh" \ + "$DOCKER_IMAGE" \ + "H20 vLLM TP×DP matrix for DeepSeek-V4-Flash" + + local server_args_str + server_args_str="$(build_server_args "$tp" "$dp")" + jq --arg tp "$tp" --arg dp "$dp" --arg cuda "$CUDA_VISIBLE_DEVICES" --arg args "$server_args_str" \ + '.config = { + "tp": ($tp | tonumber), + "dp": ($dp | tonumber), + "cuda_visible_devices": $cuda, + "backend": "vllm", + "server_start_script": "experiments/'${EXPERIMENT_NAME}'/start_vllm_dp.sh", + "server_args": $args + }' "${result_root}/results.json" > "${result_root}/results.json.tmp" && \ + mv "${result_root}/results.json.tmp" "${result_root}/results.json" + + if [[ "$DRY_RUN" == "1" ]]; then + log "DRY_RUN: would start server with args: ${server_args_str}" + local line + tail -n +2 "$scenario_tsv" | while IFS=$'\t' read -r mark isl osl conc num; do + log "DRY_RUN: ${config_label} scenario mark=${mark} c=${conc} i=${isl} o=${osl} n=${num}" + done + log "===== ${config_label} DONE (dry run) =====" + return 0 + fi + + # Initialize skipped_after_oom.csv for this run. + local skipped_csv="${RESULT_BASE}/${RUN_ID}/skipped_after_oom.csv" + if [[ ! -f "$skipped_csv" ]]; then + echo "engine,tp,dp,mark,isl,osl,concurrency,status,reason,detail_log" > "$skipped_csv" + fi + + # Start server once for this TP×DP config. + if ! start_server "$tp" "$dp"; then + log "ERROR: ${config_label} failed to start; skipping all scenarios" + skip_remaining_scenarios "$result_root" "$scenario_tsv" 0 "SKIPPED_SERVICE_START_FAILED" "service failed to start" "$tp" "$dp" + log "===== ${config_label} DONE =====" + return 0 + fi + + # Warmup with a small prompt before the first scenario. + run_warmup 1024 128 || true + + # Read scenarios into an array so we can skip remaining entries on failure. + local -a scenarios=() + while IFS= read -r line; do + scenarios+=("$line") + done < <(tail -n +2 "$scenario_tsv") + + local i mark isl osl conc num + local output_file detail_log gpu_csv sname bench_rc + for (( i = 0; i < ${#scenarios[@]}; i++ )); do + IFS=$'\t' read -r mark isl osl conc num <<< "${scenarios[$i]}" + + if [[ "$GRID_LIMIT" -gt 0 && "$i" -ge "$GRID_LIMIT" ]]; then + log "GRID_LIMIT=${GRID_LIMIT} reached; skipping remaining scenarios" + skip_remaining_scenarios "$result_root" "$scenario_tsv" "$i" "SKIPPED_GRID_LIMIT" "GRID_LIMIT reached" "$tp" "$dp" + break + fi + + sname="c${conc}_i${isl}_o${osl}" + output_file="${raw_dir}/vllm_main_${conc}_${isl}_${osl}.jsonl" + detail_log="${phase_log_dir}/vllm_${config_label}_${sname}.log" + gpu_csv="${gpu_log_dir}/gpu_mem_${conc}_${isl}_${osl}.csv" + + if scenario_already_completed "$output_file" "$num" || scenario_already_processed "$result_root" "$sname"; then + log "skipping already-processed ${config_label} scenario: ${sname}" + continue + fi + + log "running ${config_label} scenario: mark=${mark} c=${conc} i=${isl} o=${osl} n=${num}" + + local gpu_pid + gpu_pid="$(start_gpu_monitor "$gpu_csv")" + + bench_rc=0 + timeout "$SCENARIO_TIMEOUT_S" bash -c ' + run_bench_serving \ + --backend vllm \ + --host 127.0.0.1 \ + --port "'"$VLLM_PORT"'" \ + --dataset-name random \ + --dataset-path "'"$DATASET_PATH"'" \ + --random-input-len "'"$isl"'" \ + --random-output-len "'"$osl"'" \ + --num-prompts "'"$num"'" \ + --max-concurrency "'"$conc"'" \ + --request-rate 10000 \ + --output-file "'"$output_file"'" \ + --output-details \ + > "'"$detail_log"'" 2>&1 + ' || bench_rc=$? + + stop_gpu_monitor "$gpu_pid" + + if [[ "$bench_rc" -eq 0 ]]; then + log "finished ${config_label} scenario: output=${output_file}" + append_scenario_record "$result_root" \ + "name=${sname}" \ + "config=$(jq -n --arg phase main --argjson c "$conc" --argjson i "$isl" --argjson o "$osl" --arg dataset random --argjson n "$num" '{phase: $phase, concurrency: $c, input_len: $i, output_len: $o, dataset: $dataset, num_prompts: $n}')" \ + "status=\"completed\"" \ + "note=\"benchmark finished successfully\"" + continue + fi + + # Failure handling. + if detect_oom "$detail_log" "${log_dir_global}/vllm_tp${tp}_dp${dp}.server.outer.log"; then + log "ERROR: ${config_label} scenario ${sname} triggered OOM; stopping config" + append_scenario_record "$result_root" \ + "name=${sname}" \ + "config=$(jq -n --arg phase main --argjson c "$conc" --argjson i "$isl" --argjson o "$osl" --arg dataset random --argjson n "$num" '{phase: $phase, concurrency: $c, input_len: $i, output_len: $o, dataset: $dataset, num_prompts: $n}')" \ + "status=\"OOM\"" \ + "note=\"detected CUDA out-of-memory\"" + record_skipped_csv "$skipped_csv" \ + "engine=vllm" "tp=${tp}" "dp=${dp}" "mark=${mark}" "isl=${isl}" "osl=${osl}" "concurrency=${conc}" "status=OOM" "reason=detected CUDA out-of-memory" "detail_log=${detail_log}" + stop_server "$tp" "$dp" + skip_remaining_scenarios "$result_root" "$scenario_tsv" "$((i + 1))" "SKIPPED_AFTER_OOM" "previous case OOM" "$tp" "$dp" + break + fi + + log "ERROR: ${config_label} scenario ${sname} failed (rc=${bench_rc}); see ${detail_log}" + if [[ "$mark" == "P" ]]; then + log "optional (P) scenario failed; recording as skipped and continuing" + append_scenario_record "$result_root" \ + "name=${sname}" \ + "config=$(jq -n --arg phase main --argjson c "$conc" --argjson i "$isl" --argjson o "$osl" --arg dataset random --argjson n "$num" '{phase: $phase, concurrency: $c, input_len: $i, output_len: $o, dataset: $dataset, num_prompts: $n}')" \ + "status=\"skipped_optional\"" \ + "note=\"optional scenario failed (rc=${bench_rc})\"" + record_skipped_csv "$skipped_csv" \ + "engine=vllm" "tp=${tp}" "dp=${dp}" "mark=${mark}" "isl=${isl}" "osl=${osl}" "concurrency=${conc}" "status=skipped_optional" "reason=optional scenario failed (rc=${bench_rc})" "detail_log=${detail_log}" + continue + fi + + # Mandatory scenario failed but not OOM: try to restart the server. + if restart_server "$tp" "$dp"; then + run_warmup 1024 128 || true + log "resuming ${config_label} after server restart" + continue + fi + + log "ERROR: ${config_label} server restart failed; skipping remaining scenarios" + append_scenario_record "$result_root" \ + "name=${sname}" \ + "config=$(jq -n --arg phase main --argjson c "$conc" --argjson i "$isl" --argjson o "$osl" --arg dataset random --argjson n "$num" '{phase: $phase, concurrency: $c, input_len: $i, output_len: $o, dataset: $dataset, num_prompts: $n}')" \ + "status=\"FAILED\"" \ + "note=\"scenario failed and server restart failed (rc=${bench_rc})\"" + record_skipped_csv "$skipped_csv" \ + "engine=vllm" "tp=${tp}" "dp=${dp}" "mark=${mark}" "isl=${isl}" "osl=${osl}" "concurrency=${conc}" "status=FAILED" "reason=scenario failed and server restart failed" "detail_log=${detail_log}" + skip_remaining_scenarios "$result_root" "$scenario_tsv" "$((i + 1))" "SKIPPED_RESTART_FAILED" "server restart failed" "$tp" "$dp" + break + done + + stop_server "$tp" "$dp" + + # Parse results. + log "parsing ${config_label} results" + "$PYTHON" "${SCRIPT_DIR}/../../scripts/common/parse_backend.py" "$result_root" --backend vllm \ + >> "${phase_log_dir}/parse.log" 2>&1 || { + log "WARNING: parser failed for ${config_label}; see ${phase_log_dir}/parse.log" + } + + log "===== ${config_label} DONE =====" +} + +# --------------------------------------------------------------------------- +# Main +# --------------------------------------------------------------------------- + +# Cleanup any leftovers. +for cfg in "${PARALLEL_CONFIGS[@]}"; do + read -r tp dp <<< "$cfg" + stop_server "$tp" "$dp" +done + +# Run each parallel configuration. +for cfg in "${PARALLEL_CONFIGS[@]}"; do + read -r tp dp <<< "$cfg" + run_parallel_config "$tp" "$dp" +done + +# Generate cross-configuration comparison. +log "generating comparison report" +"$PYTHON" "${SCRIPT_DIR}/compare.py" \ + --run-root "${RESULT_BASE}/${RUN_ID}" \ + --output "${RESULT_BASE}/${RUN_ID}/comparison.md" \ + >> "${log_dir_global}/compare.log" 2>&1 || { + log "WARNING: comparison script failed; see ${log_dir_global}/compare.log" + } + +log "all results saved to ${RESULT_BASE}/${RUN_ID}" diff --git a/experiments/dsv4_h20_vllm_tp_dp_matrix/start_vllm_docker.sh b/experiments/dsv4_h20_vllm_tp_dp_matrix/start_vllm_docker.sh new file mode 100755 index 0000000..bab8154 --- /dev/null +++ b/experiments/dsv4_h20_vllm_tp_dp_matrix/start_vllm_docker.sh @@ -0,0 +1,102 @@ +#!/usr/bin/env bash +# Start vLLM server in Docker for a given TPxDP configuration. +# Usage: start_vllm_docker.sh +# +# Uses the vllm/vllm-openai image and the same argument set as the +# bare-metal start script. The container is removed automatically on stop. +set -e + +TP="${1}" +DP="${2}" + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/config.env" + +RUNTIME_BASE="${RUNTIME_BASE:-${SCRIPT_DIR}/runtime}" +mkdir -p "${RUNTIME_BASE}/logs" "${RUNTIME_BASE}/tmp" + +IMAGE="${DOCKER_IMAGE:-vllm/vllm-openai:latest}" +PORT="${VLLM_PORT:-30030}" +NAME="${EXPERIMENT}_vllm_tp${TP}_dp${DP}" +PID_FILE="${RUNTIME_BASE}/${EXPERIMENT}_vllm_tp${TP}_dp${DP}.pid" + +LOG="${RUNTIME_BASE}/logs/${EXPERIMENT}_vllm_docker_tp${TP}_dp${DP}_$(date +%Y%m%d_%H%M%S).log" +rm -f "$PID_FILE" + +# Clean up any stale container with the same name. +docker rm -f "$NAME" >/dev/null 2>&1 || true + +SERVER_ARGS=( + "$MODEL_PATH" + --trust-remote-code + --kv-cache-dtype "$KV_CACHE_DTYPE" + --block-size "$BLOCK_SIZE" + --tensor-parallel-size "$TP" + --gpu-memory-utilization "$GPU_MEMORY_UTILIZATION" + --max-model-len "$MAX_MODEL_LEN" + --max-num-seqs "$MAX_NUM_SEQS" + --no-enable-flashinfer-autotune + --host 0.0.0.0 + --port "$PORT" +) + +if [[ "$DP" -gt 1 ]]; then + SERVER_ARGS+=( + --data-parallel-size "$DP" + ) +fi + +SERVER_ARGS_STR="${SERVER_ARGS[*]}" + +echo "=== Starting vLLM server in Docker (TP=${TP}, DP=${DP}) ===" +echo "Image: $IMAGE" +echo "Model: $MODEL_PATH" +echo "Container name: $NAME" +echo "Host port: $PORT" +echo "Command: vllm ${SERVER_ARGS_STR}" +echo "Log: $LOG" + +# Run docker in the foreground so that killing the host process stops the +# container (the --rm flag ensures cleanup). nohup lets us background it and +# capture the host PID in the same way as the bare-metal start script. +nohup docker run --rm \ + --name "$NAME" \ + --gpus all \ + --ipc host \ + --shm-size 16g \ + --ulimit memlock=-1 \ + -p "${PORT}:${PORT}" \ + -v "${MODEL_PATH}:${MODEL_PATH}:ro" \ + -v "${RUNTIME_BASE}/tmp:/tmp" \ + -e CUDA_VISIBLE_DEVICES="${CUDA_VISIBLE_DEVICES}" \ + -e PYTHONUNBUFFERED=1 \ + -e HF_HUB_OFFLINE=1 \ + -e TRANSFORMERS_OFFLINE=1 \ + "$IMAGE" \ + "${SERVER_ARGS[@]}" \ + > "$LOG" 2>&1 & + +PID=$! +echo $PID > "$PID_FILE" +echo "PID: $PID" +echo "Waiting for health on port ${PORT}..." + +for i in $(seq 1 240); do + if curl --fail --silent --show-error --max-time 5 "http://127.0.0.1:${PORT}/health" >/dev/null 2>&1; then + echo "vLLM server is ready at http://127.0.0.1:${PORT}" + echo "Log: $LOG" + exit 0 + fi + if ! kill -0 $PID 2>/dev/null; then + echo "ERROR: Docker vLLM server exited early" + tail -200 "$LOG" + exit 1 + fi + echo "Waiting... ($i/240)" + sleep 5 +done + +echo "ERROR: Docker vLLM server not healthy after 240 retries" +tail -200 "$LOG" +exit 1 diff --git a/experiments/dsv4_h20_vllm_tp_dp_matrix/start_vllm_dp.sh b/experiments/dsv4_h20_vllm_tp_dp_matrix/start_vllm_dp.sh new file mode 100755 index 0000000..741ca9c --- /dev/null +++ b/experiments/dsv4_h20_vllm_tp_dp_matrix/start_vllm_dp.sh @@ -0,0 +1,87 @@ +#!/usr/bin/env bash +# Start vLLM server for a given TP×DP configuration. +# Usage: start_vllm_dp.sh +# +# By default this delegates to the Docker start script because the experiment +# is intended to run vLLM inside a container. Set USE_DOCKER=0 to use the +# local VENV_VLLM environment instead. +set -e + +TP="${1}" +DP="${2}" + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +# shellcheck source=/dev/null +source "${SCRIPT_DIR}/config.env" + +if [[ "${USE_DOCKER:-1}" == "1" ]]; then + exec "${SCRIPT_DIR}/start_vllm_docker.sh" "$@" +fi + +RUNTIME_BASE="${RUNTIME_BASE:-${SCRIPT_DIR}/runtime}" +mkdir -p "${RUNTIME_BASE}/logs" "${RUNTIME_BASE}/tmp" + +VENV="${VENV_VLLM}" +export PATH="$VENV/bin:$PATH" +export PYTHONUNBUFFERED=1 +export TMPDIR="${RUNTIME_BASE}/tmp" +export CUDA_VISIBLE_DEVICES="${CUDA_VISIBLE_DEVICES}" + +LOG="${RUNTIME_BASE}/logs/${EXPERIMENT}_vllm_tp${TP}_dp${DP}_$(date +%Y%m%d_%H%M%S).log" +PID_FILE="${RUNTIME_BASE}/${EXPERIMENT}_vllm_tp${TP}_dp${DP}.pid" + +rm -f "$PID_FILE" + +SERVER_ARGS=( + vllm serve "$MODEL_PATH" + --trust-remote-code + --kv-cache-dtype "$KV_CACHE_DTYPE" + --block-size "$BLOCK_SIZE" + --tensor-parallel-size "$TP" + --gpu-memory-utilization "$GPU_MEMORY_UTILIZATION" + --max-model-len "$MAX_MODEL_LEN" + --max-num-seqs "$MAX_NUM_SEQS" + --no-enable-flashinfer-autotune + --host 0.0.0.0 + --port "$VLLM_PORT" +) + +if [[ "$DP" -gt 1 ]]; then + SERVER_ARGS+=( + --data-parallel-size "$DP" + ) +fi + +SERVER_ARGS_STR="${SERVER_ARGS[*]}" + +echo "=== Starting vLLM server (TP=${TP}, DP=${DP}) ===" +echo "Model: $MODEL_PATH" +echo "Port: $VLLM_PORT" +echo "Command: $SERVER_ARGS_STR" +echo "Log: $LOG" + +nohup "${SERVER_ARGS[@]}" > "$LOG" 2>&1 & + +PID=$! +echo $PID > "$PID_FILE" +echo "PID: $PID" +echo "Waiting for health on port ${VLLM_PORT}..." + +for i in $(seq 1 240); do + if curl --fail --silent --show-error --max-time 5 "http://127.0.0.1:${VLLM_PORT}/health" >/dev/null 2>&1; then + echo "vLLM server is ready at http://127.0.0.1:${VLLM_PORT}" + echo "Log: $LOG" + exit 0 + fi + if ! kill -0 $PID 2>/dev/null; then + echo "ERROR: vLLM server exited early" + tail -200 "$LOG" + exit 1 + fi + echo "Waiting... ($i/240)" + sleep 5 +done + +echo "ERROR: vLLM server not healthy after 240 retries" +tail -200 "$LOG" +exit 1 diff --git a/platforms/nvidia_h20.env b/platforms/nvidia_h20.env new file mode 100644 index 0000000..cf635db --- /dev/null +++ b/platforms/nvidia_h20.env @@ -0,0 +1,24 @@ +# Platform configuration for NVIDIA H20 +# Source this file via scripts/common/platform.sh + +CHIP="nvidia_h20" +ACCELERATOR="NVIDIA H20" +HARDWARE="8x NVIDIA H20 96GB" +ENGINE="vllm" + +# Device selection +DEVICE_SELECT_ENV="CUDA_VISIBLE_DEVICES=0,1,2,3,4,5,6,7" +CUDA_VISIBLE_DEVICES="0,1,2,3,4,5,6,7" + +# Default serving port and model root on the host +DEFAULT_PORT="30004" +MODEL_ROOT="/data1/hf_models" + +# Virtual environments on the host (used by native H20 scripts). +# Override these if your H20 machine uses different paths. +VENV_VLLM="${VENV_VLLM:-/data1/yy/sskj/envs/vllm}" +VENV_SGLANG="${VENV_SGLANG:-/data1/yy/sskj/envs/sglang}" +VENV_CLIENT="${VENV_CLIENT:-/data1/yy/sskj/envs/sglang}" + +# Default server start script for the legacy benchmark grid. +SERVER_START_SCRIPT="${SERVER_START_SCRIPT:-${ROOT_DIR}/scripts/start_dsv4_dspark_8card.sh}" diff --git a/scripts/common/platform.sh b/scripts/common/platform.sh index b20f88f..8fea032 100755 --- a/scripts/common/platform.sh +++ b/scripts/common/platform.sh @@ -20,7 +20,13 @@ if [[ -z "${PLATFORM:-}" ]]; then elif lspci 2>/dev/null | grep -qiE "XPU|Kunlun"; then PLATFORM="kunlun_p800" elif command -v nvidia-smi >/dev/null 2>&1; then - PLATFORM="nvidia_h200" + _GPU_NAME=$(nvidia-smi --query-gpu=name --format=csv,noheader 2>/dev/null | head -n 1 || true) + if [[ -n "$_GPU_NAME" ]] && grep -qi "H20" <<< "$_GPU_NAME"; then + PLATFORM="nvidia_h20" + else + PLATFORM="nvidia_h200" + fi + unset _GPU_NAME else echo "ERROR: Could not auto-detect PLATFORM. Set PLATFORM env var explicitly." >&2 exit 1