InferScale-Sim / src /inferscale /optimizer.py
ArchitSharma's picture
Commiting v0.3
44745f2
Raw
History Blame Contribute Delete
11.2 kB
from __future__ import annotations
from copy import deepcopy
from statistics import mean, pstdev
from .models import SimulationConfig
from .simulator import run_simulation
def evaluate_rate(base: SimulationConfig, rate: float, repetitions: int = 3) -> dict:
runs = []
for rep in range(repetitions):
cfg = deepcopy(base)
cfg.request_rate_rps = rate
cfg.seed = base.seed + rep * 101
runs.append(run_simulation(cfg.to_dict()))
attainments = [run["summary"]["slo_attainment"] for run in runs]
goodputs = [run["summary"]["goodput_rps"] for run in runs]
ttfts = [run["latency"]["ttft_ms"]["p95"] for run in runs]
e2es = [run["latency"]["e2e_ms"]["p95"] for run in runs]
unfinished_values = [run["summary"]["requests_unfinished"] for run in runs]
passed = all(
attainment >= base.slo_attainment_target and unfinished == 0
for attainment, unfinished in zip(attainments, unfinished_values, strict=True)
)
return {
"rate_rps": rate,
"passed": passed,
"slo_attainment": mean(attainments),
"slo_attainment_min": min(attainments),
"slo_attainment_max": max(attainments),
"slo_attainment_std": pstdev(attainments) if len(attainments) > 1 else 0.0,
"goodput_rps": mean(goodputs),
"goodput_std": pstdev(goodputs) if len(goodputs) > 1 else 0.0,
"p95_ttft_ms": mean(ttfts),
"p95_e2e_ms": mean(e2es),
"mean_unfinished": mean(unfinished_values),
"max_unfinished": max(unfinished_values),
"repetitions": repetitions,
"criterion": "all-repetitions-meet-target-and-drain",
"target": base.slo_attainment_target,
}
def capacity_search(
config: dict,
min_rate: float = 0.25,
max_rate: float = 32.0,
iterations: int = 8,
repetitions: int = 2,
headroom: float = 0.20,
) -> dict:
base = SimulationConfig.from_dict(config)
low = max(0.01, min_rate)
high = max(low * 1.01, max_rate)
trace: list[dict] = []
low_eval = evaluate_rate(base, low, repetitions)
trace.append(low_eval)
if not low_eval["passed"]:
return {
"status": "no_feasible_rate",
"capacity_rps": 0.0,
"recommended_rps": 0.0,
"headroom": headroom,
"criterion": "all-repetitions-meet-target-and-drain",
"trace": trace,
}
high_eval = evaluate_rate(base, high, repetitions)
trace.append(high_eval)
if high_eval["passed"]:
return {
"status": "upper_bound_still_feasible",
"capacity_rps": high,
"recommended_rps": high * (1.0 - headroom),
"headroom": headroom,
"criterion": "all-repetitions-meet-target-and-drain",
"trace": sorted(trace, key=lambda item: item["rate_rps"]),
}
best = low
for _ in range(max(iterations, 1)):
mid = (low + high) / 2.0
result = evaluate_rate(base, mid, repetitions)
trace.append(result)
if result["passed"]:
best = mid
low = mid
else:
high = mid
return {
"status": "ok",
"capacity_rps": best,
"recommended_rps": best * (1.0 - headroom),
"headroom": headroom,
"criterion": "all-repetitions-meet-target-and-drain",
"trace": sorted(trace, key=lambda item: item["rate_rps"]),
}
def compare_schedulers(config: dict, schedulers: list[str] | None = None) -> dict:
schedulers = schedulers or [
"static_fcfs",
"continuous_fcfs",
"continuous_sjf",
"continuous_slo",
"chunked_slo",
]
base = SimulationConfig.from_dict(config)
base.topology = "colocated"
rows = []
for scheduler in schedulers:
cfg = deepcopy(base)
cfg.scheduler = scheduler
result = run_simulation(cfg.to_dict())
rows.append(
{
"scheduler": scheduler,
"request_throughput_rps": result["summary"]["request_throughput_rps"],
"goodput_rps": result["summary"]["goodput_rps"],
"slo_attainment": result["summary"]["slo_attainment"],
"p95_ttft_ms": result["latency"]["ttft_ms"]["p95"],
"p95_e2e_ms": result["latency"]["e2e_ms"]["p95"],
"peak_kv_utilization": result["resource"]["peak_kv_utilization"],
"unfinished": result["summary"]["requests_unfinished"],
"bottleneck": result["diagnostics"]["label"],
}
)
rows.sort(key=lambda row: (row["slo_attainment"], row["goodput_rps"]), reverse=True)
return {"rows": rows}
def compare_topologies(config: dict) -> dict:
"""Compare colocated and P/D-disaggregated serving on one deterministic trace."""
base = SimulationConfig.from_dict(config)
rows = []
scenarios = [
("colocated", False),
("colocated", True),
("disaggregated_pd", False),
("disaggregated_pd", True),
]
for topology, cache_enabled in scenarios:
cfg = deepcopy(base)
cfg.topology = topology
cfg.prefix_cache_enabled = cache_enabled
if topology == "disaggregated_pd" and cfg.scheduler == "static_fcfs":
cfg.scheduler = "continuous_fcfs"
result = run_simulation(cfg.to_dict())
resource = result["resource"]
accelerator_instances = int(resource.get("accelerator_instances", 1))
goodput_rps = result["summary"]["goodput_rps"]
rows.append(
{
"scenario": f"{topology}{'_cache' if cache_enabled else ''}",
"topology": topology,
"prefix_cache": cache_enabled,
"accelerator_instances": accelerator_instances,
"goodput_rps": goodput_rps,
"goodput_per_accelerator": goodput_rps / max(accelerator_instances, 1),
"request_throughput_rps": result["summary"]["request_throughput_rps"],
"slo_attainment": result["summary"]["slo_attainment"],
"p95_ttft_ms": result["latency"]["ttft_ms"]["p95"],
"p95_e2e_ms": result["latency"]["e2e_ms"]["p95"],
"peak_kv_utilization": resource.get("peak_kv_utilization", 0.0),
"prefix_hit_rate": resource.get("prefix_cache_hit_rate", 0.0),
"prefill_tokens_saved": resource.get("prefill_tokens_saved", 0),
"p95_transfer_ms": resource.get("p95_transfer_ms", 0.0),
"prefill_busy_fraction": resource.get("prefill_busy_fraction", 0.0),
"decode_busy_fraction": resource.get("decode_busy_fraction", result["summary"].get("busy_fraction", 0.0)),
"transfer_busy_fraction": resource.get("transfer_busy_fraction", 0.0),
"bottleneck": result["diagnostics"]["label"],
}
)
rows.sort(key=lambda row: (row["slo_attainment"], row["goodput_rps"]), reverse=True)
return {"rows": rows}
def _is_dominated(candidate: dict, rows: list[dict], throughput_key: str = "goodput_rps") -> bool:
for other in rows:
if other is candidate:
continue
no_worse = (
other[throughput_key] >= candidate[throughput_key]
and other["p95_ttft_ms"] <= candidate["p95_ttft_ms"]
)
strictly_better = (
other[throughput_key] > candidate[throughput_key]
or other["p95_ttft_ms"] < candidate["p95_ttft_ms"]
)
if no_worse and strictly_better:
return True
return False
def design_space_search(config: dict, include_disaggregated: bool = True) -> dict:
"""Small browser-safe design-space sweep with a goodput/TTFT Pareto frontier.
The sweep is intentionally bounded: its purpose is to expose configuration
interactions interactively, not to claim exhaustive optimization.
"""
base = SimulationConfig.from_dict(config)
rows: list[dict] = []
colocated_candidates = []
for scheduler in ["continuous_fcfs", "continuous_slo", "chunked_slo"]:
for batch in [8, 16, 32]:
colocated_candidates.append((scheduler, batch, False))
# Prefix reuse is added as a separate systems dimension for the SLO-aware
# scheduler at the user's current batch size.
colocated_candidates.append(("continuous_slo", base.max_batch_size, True))
for scheduler, batch, cache in colocated_candidates:
cfg = deepcopy(base)
cfg.topology = "colocated"
cfg.scheduler = scheduler
cfg.max_batch_size = batch
cfg.prefix_cache_enabled = cache
result = run_simulation(cfg.to_dict())
rows.append(_design_row(result, cfg, f"colocated / {scheduler} / batch {batch}{' / cache' if cache else ''}"))
if include_disaggregated:
for prefill_workers, decode_workers in [(1, 1), (1, 2), (2, 1)]:
for cache in [False, True]:
cfg = deepcopy(base)
cfg.topology = "disaggregated_pd"
cfg.scheduler = "continuous_slo"
cfg.prefill_workers = prefill_workers
cfg.decode_workers = decode_workers
cfg.prefix_cache_enabled = cache
result = run_simulation(cfg.to_dict())
label = f"P/D {prefill_workers}P:{decode_workers}D{' / cache' if cache else ''}"
rows.append(_design_row(result, cfg, label))
for row in rows:
row["pareto"] = not _is_dominated(row, rows, "goodput_rps")
row["efficiency_pareto"] = not _is_dominated(row, rows, "goodput_per_accelerator")
rows.sort(key=lambda r: (not r["pareto"], not r["slo_pass"], -r["goodput_rps"], r["p95_ttft_ms"]))
return {
"rows": rows,
"pareto_count": sum(1 for r in rows if r["pareto"]),
"efficiency_pareto_count": sum(1 for r in rows if r["efficiency_pareto"]),
"candidate_count": len(rows),
"objectives": [
"performance: maximize goodput / minimize p95 TTFT",
"efficiency: maximize goodput per accelerator / minimize p95 TTFT",
],
}
def _design_row(result: dict, cfg: SimulationConfig, label: str) -> dict:
resource = result["resource"]
accelerator_instances = int(resource.get("accelerator_instances", 1))
goodput_rps = result["summary"]["goodput_rps"]
return {
"label": label,
"topology": cfg.topology,
"scheduler": cfg.scheduler,
"max_batch_size": cfg.max_batch_size,
"prefix_cache": cfg.prefix_cache_enabled,
"accelerator_instances": accelerator_instances,
"goodput_rps": goodput_rps,
"goodput_per_accelerator": goodput_rps / max(accelerator_instances, 1),
"slo_attainment": result["summary"]["slo_attainment"],
"slo_pass": result["summary"]["slo_attainment"] >= cfg.slo_attainment_target and result["summary"]["requests_unfinished"] == 0,
"p95_ttft_ms": result["latency"]["ttft_ms"]["p95"],
"p95_e2e_ms": result["latency"]["e2e_ms"]["p95"],
"peak_kv_utilization": resource.get("peak_kv_utilization", 0.0),
"bottleneck": result["diagnostics"]["label"],
}