forked from Karylab-cklius/vllm
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
da742158e1 | ||
|
|
53c8d72f71 | ||
|
|
0ccb2ef093 | ||
|
|
bb529e2e47 |
@@ -44,6 +44,7 @@ from vllm.distributed.stateless_coordinator import StatelessGroupCoordinator
|
||||
from vllm.distributed.utils import StatelessProcessGroup
|
||||
from vllm.logger import init_logger
|
||||
from vllm.model_executor.models.interfaces import MixtureOfExperts
|
||||
from vllm.v1.metrics.stats import EplbMetricsStats
|
||||
|
||||
from .async_worker import start_async_worker
|
||||
from .policy import EPLB_POLICIES, AbstractEplbPolicy, DefaultEplbPolicy
|
||||
@@ -311,6 +312,20 @@ class EplbState:
|
||||
newly started EP ranks may not have physical experts
|
||||
mapped yet.
|
||||
"""
|
||||
self.last_eplb_stats: EplbMetricsStats | None = None
|
||||
"""
|
||||
Most recent EPLB balancedness stats for Prometheus export.
|
||||
Computed locally (no inter-rank sync) each step.
|
||||
"""
|
||||
self.rearrangements_since_last_report: int = 0
|
||||
"""
|
||||
Number of rearrangements since the last stats were consumed.
|
||||
"""
|
||||
self.last_rearrangement_seconds: float = 0.0
|
||||
"""
|
||||
Duration of the most recent rearrangement in seconds.
|
||||
"""
|
||||
|
||||
if self.device.type == "cuda":
|
||||
self.cuda_device_index = self.device.index
|
||||
if self.cuda_device_index is None and torch.cuda.is_available():
|
||||
@@ -571,18 +586,22 @@ class EplbState:
|
||||
.float()
|
||||
)
|
||||
|
||||
# Compute balancedness ratio:
|
||||
# for each layer:
|
||||
# (mean load across ranks) / (max load across ranks)
|
||||
avg_tokens_tensor = num_tokens_per_rank.mean(dim=0).sum(dim=0)
|
||||
max_tokens_tensor = num_tokens_per_rank.max(dim=0).values.sum(dim=0)
|
||||
# Compute per-layer balancedness ratio:
|
||||
# for each layer: (mean across ranks) / (max across ranks)
|
||||
# then average across layers.
|
||||
# dim=-1 is the rank dimension.
|
||||
avg_per_layer = num_tokens_per_rank.mean(dim=-1)
|
||||
max_per_layer = num_tokens_per_rank.max(dim=-1).values
|
||||
per_layer_balance = torch.where(
|
||||
max_per_layer > 0,
|
||||
avg_per_layer / max_per_layer,
|
||||
torch.ones_like(max_per_layer),
|
||||
)
|
||||
balancedness = float(per_layer_balance.mean().item())
|
||||
|
||||
# Just to make type checker happy
|
||||
tokens_tensors: list[float] = torch.stack(
|
||||
[avg_tokens_tensor, max_tokens_tensor]
|
||||
).tolist()
|
||||
avg_tokens, max_tokens = tokens_tensors
|
||||
balancedness = avg_tokens / max_tokens if max_tokens > 0 else 0.0
|
||||
# Summary stats for logging
|
||||
avg_tokens = float(avg_per_layer.sum().item())
|
||||
max_tokens = float(max_per_layer.sum().item())
|
||||
|
||||
if ep_group.rank() == 0:
|
||||
logger.info(
|
||||
@@ -598,6 +617,55 @@ class EplbState:
|
||||
- self.expert_rearrangement_step,
|
||||
)
|
||||
|
||||
# Per-layer breakdown: worst/best layers,
|
||||
# per-rank token counts for the worst layer
|
||||
worst_layer = int(per_layer_balance.argmin().item())
|
||||
best_layer = int(per_layer_balance.argmax().item())
|
||||
worst_balance = float(per_layer_balance[worst_layer].item())
|
||||
best_balance = float(per_layer_balance[best_layer].item())
|
||||
|
||||
worst_layer_ranks = num_tokens_per_rank[worst_layer]
|
||||
worst_min_rank = int(worst_layer_ranks.argmin().item())
|
||||
worst_max_rank = int(worst_layer_ranks.argmax().item())
|
||||
|
||||
logger.info(
|
||||
"EPLB balance breakdown: "
|
||||
"worst_layer=%d (balance=%.4f, "
|
||||
"min_rank=%d[%.0f], max_rank=%d[%.0f]), "
|
||||
"best_layer=%d (balance=%.4f), "
|
||||
"num_layers=%d",
|
||||
worst_layer,
|
||||
worst_balance,
|
||||
worst_min_rank,
|
||||
float(worst_layer_ranks[worst_min_rank].item()),
|
||||
worst_max_rank,
|
||||
float(worst_layer_ranks[worst_max_rank].item()),
|
||||
best_layer,
|
||||
best_balance,
|
||||
num_tokens_per_rank.shape[0],
|
||||
)
|
||||
|
||||
# Log replica distribution for debug
|
||||
replica_count = eplb_model_state.logical_replica_count
|
||||
if replica_count is not None and replica_count.numel() > 0:
|
||||
rc_float = replica_count.float()
|
||||
logger.debug(
|
||||
"EPLB replica stats (layer avg): "
|
||||
"min=%.1f, max=%.1f, mean=%.2f, "
|
||||
"num_with_replicas=%d/%d",
|
||||
float(rc_float.min().item()),
|
||||
float(rc_float.max().item()),
|
||||
float(rc_float.mean().item()),
|
||||
int((rc_float > 1).any(dim=0).sum().item()),
|
||||
replica_count.shape[-1],
|
||||
)
|
||||
|
||||
# Compute local balancedness stats for Prometheus (no inter-rank sync).
|
||||
# Uses only the driver rank's expert_load_pass which records routing
|
||||
# decisions for all physical experts across all EP ranks.
|
||||
if not is_dummy:
|
||||
self._compute_local_load_stats()
|
||||
|
||||
# Update the expert load sliding window
|
||||
if not is_dummy:
|
||||
for eplb_model_state in self.model_states.values():
|
||||
@@ -639,6 +707,7 @@ class EplbState:
|
||||
return
|
||||
self.expert_rearrangement_step = 0
|
||||
self.rearrange()
|
||||
self.rearrangements_since_last_report += 1
|
||||
|
||||
def rearrange(
|
||||
self,
|
||||
@@ -674,6 +743,33 @@ class EplbState:
|
||||
)
|
||||
|
||||
# Map the physical expert load to global logical experts
|
||||
if is_main_rank:
|
||||
# Log window utilization diagnostics
|
||||
nonzero_slots = sum(
|
||||
int((ms.expert_load_window.sum(dim=(1, 2)) > 0).sum().item())
|
||||
for ms in self.model_states.values()
|
||||
)
|
||||
logger.info(
|
||||
"EPLB window state: window_step=%d/%d, "
|
||||
"rearrangement_step=%d/%d, "
|
||||
"nonzero_window_slots=%d/%d",
|
||||
self.expert_load_window_step,
|
||||
self.expert_load_window_size,
|
||||
self.expert_rearrangement_step,
|
||||
self.expert_rearrangement_step_interval,
|
||||
nonzero_slots,
|
||||
self.expert_load_window_size,
|
||||
)
|
||||
if self.expert_load_window_size > self.expert_rearrangement_step_interval:
|
||||
logger.warning(
|
||||
"EPLB: window_size (%d) > step_interval (%d). "
|
||||
"Stale window entries from before the last "
|
||||
"rearrangement will be converted with the current "
|
||||
"physical->logical mapping, which may be incorrect. "
|
||||
"Consider setting window_size <= step_interval.",
|
||||
self.expert_load_window_size,
|
||||
self.expert_rearrangement_step_interval,
|
||||
)
|
||||
global_expert_load_windows = []
|
||||
for eplb_model_state in self.model_states.values():
|
||||
expert_load_window = eplb_model_state.expert_load_window[
|
||||
@@ -736,6 +832,35 @@ class EplbState:
|
||||
for eplb_model_state, global_expert_load_window in zip(
|
||||
self.model_states.values(), global_expert_load_windows
|
||||
):
|
||||
if is_main_rank:
|
||||
# Log load statistics the algorithm will use
|
||||
load = global_expert_load_window.float()
|
||||
load_per_layer = load.sum(dim=-1)
|
||||
logger.info(
|
||||
"EPLB rearrange input: "
|
||||
"num_replicas=%d, num_groups=%d, "
|
||||
"num_nodes=%d, num_gpus=%d, "
|
||||
"total_load_per_layer: "
|
||||
"min=%.0f, max=%.0f, mean=%.0f",
|
||||
num_replicas,
|
||||
num_groups,
|
||||
num_nodes,
|
||||
num_gpus,
|
||||
float(load_per_layer.min().item()),
|
||||
float(load_per_layer.max().item()),
|
||||
float(load_per_layer.mean().item()),
|
||||
)
|
||||
# Top-5 hottest experts (averaged across layers)
|
||||
avg_load = load.mean(dim=0)
|
||||
top5_vals, top5_ids = avg_load.topk(min(5, avg_load.shape[0]))
|
||||
logger.info(
|
||||
"EPLB top-5 hottest logical experts (avg across layers): %s",
|
||||
", ".join(
|
||||
f"e{int(eid)}={float(val):.0f}"
|
||||
for eid, val in zip(top5_ids.tolist(), top5_vals.tolist())
|
||||
),
|
||||
)
|
||||
|
||||
if not self.is_async or is_profile:
|
||||
# Get new expert mappings for the model
|
||||
(
|
||||
@@ -751,6 +876,54 @@ class EplbState:
|
||||
eplb_model_state.physical_to_logical_map,
|
||||
)
|
||||
|
||||
if is_main_rank and not is_profile:
|
||||
# Log what the algorithm decided
|
||||
old_p2l = eplb_model_state.physical_to_logical_map
|
||||
new_p2l = new_physical_to_logical_map.to(old_p2l.device)
|
||||
changed_slots = int((old_p2l != new_p2l).sum().item())
|
||||
total_slots = old_p2l.numel()
|
||||
rc = new_logical_replica_count.float()
|
||||
logger.info(
|
||||
"EPLB rearrange result: "
|
||||
"changed_slots=%d/%d (%.1f%%), "
|
||||
"replica_count: "
|
||||
"min=%.0f, max=%.0f, mean=%.2f",
|
||||
changed_slots,
|
||||
total_slots,
|
||||
100.0 * changed_slots / max(total_slots, 1),
|
||||
float(rc.min().item()),
|
||||
float(rc.max().item()),
|
||||
float(rc.mean().item()),
|
||||
)
|
||||
|
||||
# Simulate new per-rank load to preview
|
||||
# balancedness
|
||||
new_rc = new_logical_replica_count.to(load.device).float()
|
||||
per_expert_load = load / new_rc.clamp(min=1)
|
||||
phys_load = per_expert_load.gather(
|
||||
dim=-1,
|
||||
index=new_p2l.to(load.device).long(),
|
||||
)
|
||||
per_rank_load = phys_load.reshape(
|
||||
phys_load.shape[0], num_gpus, -1
|
||||
).sum(dim=-1)
|
||||
avg_rl = per_rank_load.mean(dim=-1)
|
||||
max_rl = per_rank_load.max(dim=-1).values
|
||||
predicted_balance = torch.where(
|
||||
max_rl > 0,
|
||||
avg_rl / max_rl,
|
||||
torch.ones_like(max_rl),
|
||||
)
|
||||
logger.info(
|
||||
"EPLB predicted post-rearrange "
|
||||
"balancedness: mean=%.4f, "
|
||||
"min=%.4f (layer %d), max=%.4f",
|
||||
float(predicted_balance.mean().item()),
|
||||
float(predicted_balance.min().item()),
|
||||
int(predicted_balance.argmin().item()),
|
||||
float(predicted_balance.max().item()),
|
||||
)
|
||||
|
||||
# Update expert weights
|
||||
rearrange_expert_weights_inplace(
|
||||
eplb_model_state.physical_to_logical_map,
|
||||
@@ -801,6 +974,8 @@ class EplbState:
|
||||
end_event.record()
|
||||
end_event.synchronize()
|
||||
gpu_elapsed = start_event.elapsed_time(end_event) / 1000.0
|
||||
if not is_profile:
|
||||
self.last_rearrangement_seconds = gpu_elapsed
|
||||
logger.info(
|
||||
"Rearranged experts %s in %.2f s.",
|
||||
" (profile) " if is_profile else " ",
|
||||
@@ -1023,6 +1198,40 @@ class EplbState:
|
||||
offset += shape[0]
|
||||
return all_reduce_list
|
||||
|
||||
def _compute_local_load_stats(self) -> None:
|
||||
"""Compute per-layer, per-EP-rank token counts from this rank's view.
|
||||
|
||||
No inter-rank communication. Each rank's expert_load_pass records
|
||||
how many of THIS rank's tokens were routed to each physical expert.
|
||||
Physical experts are partitioned across EP ranks, so reshaping by
|
||||
rank gives per-destination-rank token loads.
|
||||
|
||||
These per-rank counts are directly summable across DP ranks in
|
||||
Prometheus to get the global load distribution.
|
||||
"""
|
||||
ep_group = get_ep_group().device_group
|
||||
num_ranks = ep_group.size()
|
||||
|
||||
# Use the first model's expert_load_pass (main model, not drafter)
|
||||
eplb_model_state = next(iter(self.model_states.values()))
|
||||
expert_load = eplb_model_state.expert_load_pass
|
||||
# expert_load: (num_moe_layers, num_physical_experts)
|
||||
num_layers = expert_load.shape[0]
|
||||
|
||||
# Reshape to (num_moe_layers, num_ranks, experts_per_rank)
|
||||
# and sum per-rank token loads
|
||||
per_rank = expert_load.reshape(num_layers, num_ranks, -1).sum(dim=2).float()
|
||||
# per_rank: (num_moe_layers, num_ep_ranks)
|
||||
|
||||
rearrangements = self.rearrangements_since_last_report
|
||||
self.rearrangements_since_last_report = 0
|
||||
|
||||
self.last_eplb_stats = EplbMetricsStats(
|
||||
tokens_per_ep_rank=per_rank.cpu().tolist(),
|
||||
rearrangements=rearrangements,
|
||||
last_rearrangement_seconds=self.last_rearrangement_seconds,
|
||||
)
|
||||
|
||||
def _sync_load_pass(self) -> list[torch.Tensor]:
|
||||
"""
|
||||
Sync the expert load pass across all ranks for log stats.
|
||||
|
||||
@@ -0,0 +1,123 @@
|
||||
# SPDX-License-Identifier: Apache-2.0
|
||||
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import prometheus_client
|
||||
|
||||
from vllm.config import ParallelConfig
|
||||
from vllm.v1.metrics.stats import EplbMetricsStats
|
||||
|
||||
|
||||
def _make_per_engine(
|
||||
metric: prometheus_client.Gauge | prometheus_client.Counter,
|
||||
per_engine_labelvalues: dict[int, list[object]],
|
||||
) -> dict[int, prometheus_client.Gauge | prometheus_client.Counter]:
|
||||
return {
|
||||
idx: metric.labels(*labelvalues)
|
||||
for idx, labelvalues in per_engine_labelvalues.items()
|
||||
}
|
||||
|
||||
|
||||
class EplbProm:
|
||||
"""Record EPLB load metrics in Prometheus.
|
||||
|
||||
Each EP rank independently reports how many tokens it routed to each
|
||||
destination EP rank, per MoE layer. These gauges are directly summable
|
||||
across DP ranks in Prometheus/Grafana — no inter-rank synchronization:
|
||||
|
||||
# Global load per EP rank per layer
|
||||
sum by (layer_idx, dst_ep_rank) (vllm:eplb_tokens_routed_to_ep_rank)
|
||||
|
||||
# Imbalance ratio per layer
|
||||
max by (layer_idx) (sum by (...) (...))
|
||||
/ avg by (layer_idx) (sum by (...) (...))
|
||||
|
||||
The ``layer_idx`` and ``dst_ep_rank`` labels are added on top of the
|
||||
standard per-engine labels (model_name, etc.).
|
||||
"""
|
||||
|
||||
_gauge_cls = prometheus_client.Gauge
|
||||
_counter_cls = prometheus_client.Counter
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
parallel_config: ParallelConfig,
|
||||
labelnames: list[str],
|
||||
per_engine_labelvalues: dict[int, list[object]],
|
||||
):
|
||||
self.enabled = parallel_config.enable_eplb
|
||||
if not self.enabled:
|
||||
return
|
||||
|
||||
# Per-layer, per-destination-rank token count gauge.
|
||||
# Extra labels: layer_idx, dst_ep_rank
|
||||
extended_labels = [*labelnames, "layer_idx", "dst_ep_rank"]
|
||||
self._tokens_gauge = self._gauge_cls(
|
||||
name="vllm:eplb_tokens_routed_to_ep_rank",
|
||||
documentation=(
|
||||
"Tokens routed to each destination EP rank per MoE layer "
|
||||
"(from this DP rank's perspective). Summable across DP ranks."
|
||||
),
|
||||
multiprocess_mode="mostrecent",
|
||||
labelnames=extended_labels,
|
||||
)
|
||||
self._per_engine_labelvalues = per_engine_labelvalues
|
||||
|
||||
# Rearrangement counter and timing (shared labels, no extra dims)
|
||||
counter_rearrangements = self._counter_cls(
|
||||
name="vllm:eplb_rearrangements_total",
|
||||
documentation="Total number of EPLB expert rearrangements.",
|
||||
labelnames=labelnames,
|
||||
)
|
||||
self.counter_rearrangements = _make_per_engine(
|
||||
counter_rearrangements, per_engine_labelvalues
|
||||
)
|
||||
|
||||
gauge_rearrangement_seconds = self._gauge_cls(
|
||||
name="vllm:eplb_rearrangement_seconds",
|
||||
documentation=(
|
||||
"Duration of the most recent EPLB expert rearrangement in seconds."
|
||||
),
|
||||
multiprocess_mode="mostrecent",
|
||||
labelnames=labelnames,
|
||||
)
|
||||
self.gauge_rearrangement_seconds = _make_per_engine(
|
||||
gauge_rearrangement_seconds, per_engine_labelvalues
|
||||
)
|
||||
|
||||
# Cache of labeled gauge children keyed by
|
||||
# (engine_idx, layer_idx, dst_ep_rank) to avoid repeated .labels()
|
||||
self._tokens_children: dict[
|
||||
tuple[int, int, int],
|
||||
prometheus_client.Gauge,
|
||||
] = {}
|
||||
|
||||
def _get_tokens_child(
|
||||
self, engine_idx: int, layer_idx: int, dst_ep_rank: int
|
||||
) -> prometheus_client.Gauge:
|
||||
key = (engine_idx, layer_idx, dst_ep_rank)
|
||||
child = self._tokens_children.get(key)
|
||||
if child is None:
|
||||
base_labels = self._per_engine_labelvalues[engine_idx]
|
||||
child = self._tokens_gauge.labels(
|
||||
*base_labels, str(layer_idx), str(dst_ep_rank)
|
||||
)
|
||||
self._tokens_children[key] = child
|
||||
return child
|
||||
|
||||
def observe(self, eplb_stats: EplbMetricsStats, engine_idx: int = 0):
|
||||
if not self.enabled:
|
||||
return
|
||||
|
||||
# Set per-layer, per-EP-rank token counts
|
||||
for layer_idx, rank_counts in enumerate(eplb_stats.tokens_per_ep_rank):
|
||||
for dst_rank, count in enumerate(rank_counts):
|
||||
self._get_tokens_child(engine_idx, layer_idx, dst_rank).set(count)
|
||||
|
||||
# Rearrangement metrics
|
||||
if eplb_stats.rearrangements > 0:
|
||||
self.counter_rearrangements[engine_idx].inc(eplb_stats.rearrangements)
|
||||
self.gauge_rearrangement_seconds[engine_idx].set(
|
||||
eplb_stats.last_rearrangement_seconds
|
||||
)
|
||||
@@ -50,7 +50,7 @@ from vllm.v1.core.sched.utils import check_stop, remove_all
|
||||
from vllm.v1.engine import EngineCoreEventType, EngineCoreOutput, EngineCoreOutputs
|
||||
from vllm.v1.kv_cache_interface import KVCacheConfig, MambaSpec
|
||||
from vllm.v1.metrics.perf import ModelMetrics, PerfStats
|
||||
from vllm.v1.metrics.stats import PrefixCacheStats, SchedulerStats
|
||||
from vllm.v1.metrics.stats import EplbMetricsStats, PrefixCacheStats, SchedulerStats
|
||||
from vllm.v1.outputs import DraftTokenIds, KVConnectorOutput, ModelRunnerOutput
|
||||
from vllm.v1.request import Request, RequestStatus, StreamingUpdate
|
||||
from vllm.v1.spec_decode.metrics import SpecDecodingStats
|
||||
@@ -1287,6 +1287,7 @@ class Scheduler(SchedulerInterface):
|
||||
num_nans_in_logits = model_runner_output.num_nans_in_logits
|
||||
kv_connector_output = model_runner_output.kv_connector_output
|
||||
cudagraph_stats = model_runner_output.cudagraph_stats
|
||||
eplb_stats = model_runner_output.eplb_stats
|
||||
|
||||
perf_stats: PerfStats | None = None
|
||||
if self.perf_metrics and self.perf_metrics.is_enabled():
|
||||
@@ -1517,7 +1518,11 @@ class Scheduler(SchedulerInterface):
|
||||
|
||||
if (
|
||||
stats := self.make_stats(
|
||||
spec_decoding_stats, kv_connector_stats, cudagraph_stats, perf_stats
|
||||
spec_decoding_stats,
|
||||
kv_connector_stats,
|
||||
cudagraph_stats,
|
||||
perf_stats,
|
||||
eplb_stats,
|
||||
)
|
||||
) is not None:
|
||||
# Return stats to only one of the front-ends.
|
||||
@@ -1876,6 +1881,7 @@ class Scheduler(SchedulerInterface):
|
||||
kv_connector_stats: KVConnectorStats | None = None,
|
||||
cudagraph_stats: CUDAGraphStat | None = None,
|
||||
perf_stats: PerfStats | None = None,
|
||||
eplb_stats: EplbMetricsStats | None = None,
|
||||
) -> SchedulerStats | None:
|
||||
if not self.log_stats:
|
||||
return None
|
||||
@@ -1906,6 +1912,7 @@ class Scheduler(SchedulerInterface):
|
||||
kv_connector_stats=connector_stats_payload,
|
||||
cudagraph_stats=cudagraph_stats,
|
||||
perf_stats=perf_stats,
|
||||
eplb_stats=eplb_stats,
|
||||
)
|
||||
|
||||
def _get_encoder_cache_usage(self) -> float:
|
||||
|
||||
@@ -12,6 +12,7 @@ from prometheus_client import Counter, Gauge, Histogram
|
||||
import vllm.envs as envs
|
||||
from vllm.compilation.cuda_graph import CUDAGraphLogging
|
||||
from vllm.config import SupportsMetricsInfo, VllmConfig
|
||||
from vllm.distributed.eplb.metrics import EplbProm
|
||||
from vllm.distributed.kv_transfer.kv_connector.v1.metrics import (
|
||||
KVConnectorLogging,
|
||||
KVConnectorPrometheus,
|
||||
@@ -393,6 +394,7 @@ class PrometheusStatLogger(AggregateStatLoggerBase):
|
||||
_spec_decoding_cls = SpecDecodingProm
|
||||
_kv_connector_cls = KVConnectorPrometheus
|
||||
_perf_metrics_cls = PerfMetricsProm
|
||||
_eplb_cls = EplbProm
|
||||
|
||||
def __init__(
|
||||
self, vllm_config: VllmConfig, engine_indexes: list[int] | None = None
|
||||
@@ -428,6 +430,9 @@ class PrometheusStatLogger(AggregateStatLoggerBase):
|
||||
self.perf_metrics_prom = self._perf_metrics_cls(
|
||||
vllm_config, labelnames, per_engine_labelvalues
|
||||
)
|
||||
self.eplb_prom = self._eplb_cls(
|
||||
vllm_config.parallel_config, labelnames, per_engine_labelvalues
|
||||
)
|
||||
|
||||
#
|
||||
# Scheduler state
|
||||
@@ -1072,6 +1077,9 @@ class PrometheusStatLogger(AggregateStatLoggerBase):
|
||||
if scheduler_stats.perf_stats is not None:
|
||||
self.perf_metrics_prom.observe(scheduler_stats.perf_stats, engine_idx)
|
||||
|
||||
if scheduler_stats.eplb_stats is not None:
|
||||
self.eplb_prom.observe(scheduler_stats.eplb_stats, engine_idx)
|
||||
|
||||
if (
|
||||
self.kv_cache_metrics_enabled
|
||||
and scheduler_stats.kv_cache_eviction_events
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
|
||||
import time
|
||||
|
||||
from vllm.distributed.eplb.metrics import EplbProm
|
||||
from vllm.distributed.kv_transfer.kv_connector.v1.metrics import KVConnectorPrometheus
|
||||
from vllm.v1.metrics.loggers import PrometheusStatLogger
|
||||
from vllm.v1.metrics.perf import PerfMetricsProm
|
||||
@@ -190,6 +191,17 @@ class RayPerfMetricsProm(PerfMetricsProm):
|
||||
_counter_cls = RayCounterWrapper
|
||||
|
||||
|
||||
class RayEplbProm(EplbProm):
|
||||
"""
|
||||
RayEplbProm is used by RayMetrics to log Ray metrics.
|
||||
Provides the same EPLB load metrics as EplbProm but uses
|
||||
Ray's util.metrics library.
|
||||
"""
|
||||
|
||||
_gauge_cls = RayGaugeWrapper
|
||||
_counter_cls = RayCounterWrapper
|
||||
|
||||
|
||||
class RayPrometheusStatLogger(PrometheusStatLogger):
|
||||
"""RayPrometheusStatLogger uses Ray metrics instead."""
|
||||
|
||||
@@ -199,6 +211,7 @@ class RayPrometheusStatLogger(PrometheusStatLogger):
|
||||
_spec_decoding_cls = RaySpecDecodingProm
|
||||
_kv_connector_cls = RayKVConnectorPrometheus
|
||||
_perf_metrics_cls = RayPerfMetricsProm
|
||||
_eplb_cls = RayEplbProm
|
||||
|
||||
@staticmethod
|
||||
def _unregister_vllm_metrics():
|
||||
|
||||
@@ -167,6 +167,26 @@ class KVCacheEvictionEvent:
|
||||
reuse_gaps_seconds: tuple[float, ...]
|
||||
|
||||
|
||||
@dataclass
|
||||
class EplbMetricsStats:
|
||||
"""EPLB load stats computed per step for Prometheus export.
|
||||
|
||||
Each EP rank independently reports how many tokens it routed to each
|
||||
destination EP rank, per MoE layer. These per-rank gauges are directly
|
||||
summable across DP ranks in Prometheus/Grafana without inter-rank
|
||||
synchronization:
|
||||
|
||||
sum by (layer_idx, dst_ep_rank) (vllm:eplb_tokens_routed_to_ep_rank)
|
||||
"""
|
||||
|
||||
# tokens_per_ep_rank[layer][dst_rank] = token count
|
||||
# Shape semantics: (num_moe_layers, ep_size)
|
||||
tokens_per_ep_rank: list[list[float]]
|
||||
|
||||
rearrangements: int
|
||||
last_rearrangement_seconds: float
|
||||
|
||||
|
||||
@dataclass
|
||||
class SchedulerStats:
|
||||
"""Stats associated with the scheduler."""
|
||||
@@ -196,6 +216,8 @@ class SchedulerStats:
|
||||
|
||||
perf_stats: PerfStats | None = None
|
||||
|
||||
eplb_stats: EplbMetricsStats | None = None
|
||||
|
||||
|
||||
@dataclass
|
||||
class RequestStateStats:
|
||||
|
||||
@@ -15,9 +15,11 @@ from vllm.v1.core.sched.output import SchedulerOutput
|
||||
if TYPE_CHECKING:
|
||||
from vllm.distributed.kv_events import KVConnectorKVEvents
|
||||
from vllm.distributed.kv_transfer.kv_connector.v1.metrics import KVConnectorStats
|
||||
from vllm.v1.metrics.stats import EplbMetricsStats
|
||||
else:
|
||||
KVConnectorStats = object
|
||||
KVConnectorKVEvents = object
|
||||
EplbMetricsStats = object
|
||||
|
||||
|
||||
class LogprobsLists(NamedTuple):
|
||||
@@ -247,6 +249,9 @@ class ModelRunnerOutput:
|
||||
# information related to cudagraph execution
|
||||
cudagraph_stats: CUDAGraphStat | None = None
|
||||
|
||||
# EPLB balancedness stats
|
||||
eplb_stats: "EplbMetricsStats | None" = None
|
||||
|
||||
|
||||
# ModelRunnerOutput wrapper for async scheduling.
|
||||
class AsyncModelRunnerOutput(ABC):
|
||||
|
||||
@@ -3865,6 +3865,10 @@ class GPUModelRunner(
|
||||
else:
|
||||
logger.error("RoutedExpertsCapturer not initialized.")
|
||||
|
||||
eplb_stats = (
|
||||
self.eplb_state.last_eplb_stats if self.eplb_state is not None else None
|
||||
)
|
||||
|
||||
output = ModelRunnerOutput(
|
||||
req_ids=req_ids_output_copy,
|
||||
req_id_to_index=req_id_to_index_output_copy,
|
||||
@@ -3877,6 +3881,7 @@ class GPUModelRunner(
|
||||
else None,
|
||||
num_nans_in_logits=num_nans_in_logits,
|
||||
cudagraph_stats=cudagraph_stats,
|
||||
eplb_stats=eplb_stats,
|
||||
)
|
||||
|
||||
if not self.use_async_scheduling:
|
||||
|
||||
Reference in New Issue
Block a user