Move idle-metrics logging to SchedulerMetricsReporter (#25631)
This commit is contained in:
@@ -199,9 +199,6 @@ from sglang.srt.managers.scheduler_output_processor_mixin import (
|
|||||||
)
|
)
|
||||||
from sglang.srt.managers.scheduler_pp_mixin import SchedulerPPMixin
|
from sglang.srt.managers.scheduler_pp_mixin import SchedulerPPMixin
|
||||||
from sglang.srt.managers.scheduler_recv_skipper import SchedulerRecvSkipper
|
from sglang.srt.managers.scheduler_recv_skipper import SchedulerRecvSkipper
|
||||||
from sglang.srt.managers.scheduler_runtime_checker_mixin import (
|
|
||||||
SchedulerRuntimeCheckerMixin,
|
|
||||||
)
|
|
||||||
from sglang.srt.managers.utils import GenerationBatchResult, validate_input_length
|
from sglang.srt.managers.utils import GenerationBatchResult, validate_input_length
|
||||||
from sglang.srt.mem_cache import kv_cache_builder
|
from sglang.srt.mem_cache import kv_cache_builder
|
||||||
from sglang.srt.mem_cache.common import maybe_cache_unfinished_req, release_kv_cache
|
from sglang.srt.mem_cache.common import maybe_cache_unfinished_req, release_kv_cache
|
||||||
@@ -361,7 +358,6 @@ class Scheduler(
|
|||||||
SchedulerDisaggregationDecodeMixin,
|
SchedulerDisaggregationDecodeMixin,
|
||||||
SchedulerDisaggregationPrefillMixin,
|
SchedulerDisaggregationPrefillMixin,
|
||||||
SchedulerMultiplexMixin,
|
SchedulerMultiplexMixin,
|
||||||
SchedulerRuntimeCheckerMixin,
|
|
||||||
SchedulerPPMixin,
|
SchedulerPPMixin,
|
||||||
SchedulerDllmMixin,
|
SchedulerDllmMixin,
|
||||||
SchedulerMlxOverlapMixin,
|
SchedulerMlxOverlapMixin,
|
||||||
@@ -3188,7 +3184,7 @@ class Scheduler(
|
|||||||
self.invariant_checker._check_tree_cache()
|
self.invariant_checker._check_tree_cache()
|
||||||
|
|
||||||
# metrics every 30s
|
# metrics every 30s
|
||||||
self._maybe_log_idle_metrics()
|
self.metrics_reporter._maybe_log_idle_metrics()
|
||||||
|
|
||||||
# kv event publishing
|
# kv event publishing
|
||||||
self.kv_events_publisher.publish_kv_events()
|
self.kv_events_publisher.publish_kv_events()
|
||||||
|
|||||||
@@ -961,3 +961,46 @@ class SchedulerMetricsReporter:
|
|||||||
if ENABLE_METRICS_DEVICE_TIMER:
|
if ENABLE_METRICS_DEVICE_TIMER:
|
||||||
self._device_timer_window_batch_count = 0
|
self._device_timer_window_batch_count = 0
|
||||||
self.fwd_occupancy = float("nan")
|
self.fwd_occupancy = float("nan")
|
||||||
|
|
||||||
|
def _maybe_log_idle_metrics(self):
|
||||||
|
"""Collect and log metrics every 30 seconds during idle."""
|
||||||
|
if (
|
||||||
|
not self.current_scheduler_metrics_enabled
|
||||||
|
or time.perf_counter() <= self.metrics_collector.last_log_time + 30
|
||||||
|
):
|
||||||
|
return
|
||||||
|
|
||||||
|
self.scheduler.pool_stats_observer.get_pool_stats().update_scheduler_stats(
|
||||||
|
self.stats
|
||||||
|
)
|
||||||
|
self.stats.num_streaming_sessions = (
|
||||||
|
self.scheduler.pool_stats_observer.streaming_session_count()
|
||||||
|
)
|
||||||
|
self.stats.streaming_session_held_tokens = (
|
||||||
|
self.scheduler.pool_stats_observer.session_held_tokens()
|
||||||
|
)
|
||||||
|
|
||||||
|
priority_enabled = self.scheduler.enable_priority_scheduling
|
||||||
|
self.stats.num_running_reqs = QueueCount.from_reqs(
|
||||||
|
self.scheduler.running_batch.reqs, priority_enabled
|
||||||
|
)
|
||||||
|
self.stats.gen_throughput = 0
|
||||||
|
self.stats.num_queue_reqs = QueueCount.from_reqs(
|
||||||
|
self.scheduler.waiting_queue, priority_enabled
|
||||||
|
)
|
||||||
|
self.stats.num_grammar_queue_reqs = len(self.scheduler.grammar_manager)
|
||||||
|
if self.scheduler.disaggregation_mode == DisaggregationMode.PREFILL:
|
||||||
|
self.stats.num_prefill_bootstrap_queue_reqs = QueueCount.from_reqs(
|
||||||
|
self.scheduler.disagg_prefill_bootstrap_queue.queue, priority_enabled
|
||||||
|
)
|
||||||
|
self.stats.num_prefill_inflight_queue_reqs = QueueCount.from_reqs(
|
||||||
|
self.scheduler.disagg_prefill_inflight_queue, priority_enabled
|
||||||
|
)
|
||||||
|
if self.scheduler.disaggregation_mode == DisaggregationMode.DECODE:
|
||||||
|
self.stats.num_decode_prealloc_queue_reqs = QueueCount.from_reqs(
|
||||||
|
self.scheduler.disagg_decode_prealloc_queue.queue, priority_enabled
|
||||||
|
)
|
||||||
|
self.stats.num_decode_transfer_queue_reqs = QueueCount.from_reqs(
|
||||||
|
self.scheduler.disagg_decode_transfer_queue.queue, priority_enabled
|
||||||
|
)
|
||||||
|
self.metrics_collector.log_stats(self.stats)
|
||||||
|
|||||||
@@ -1,56 +0,0 @@
|
|||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
import logging
|
|
||||||
import time
|
|
||||||
from typing import TYPE_CHECKING
|
|
||||||
|
|
||||||
from sglang.srt.disaggregation.utils import DisaggregationMode
|
|
||||||
from sglang.srt.observability.metrics_collector import QueueCount
|
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
|
||||||
from sglang.srt.managers.scheduler import Scheduler
|
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
|
||||||
|
|
||||||
|
|
||||||
class SchedulerRuntimeCheckerMixin:
|
|
||||||
def _maybe_log_idle_metrics(self: Scheduler):
|
|
||||||
"""Collect and log metrics every 30 seconds during idle."""
|
|
||||||
if (
|
|
||||||
not self.current_scheduler_metrics_enabled
|
|
||||||
or time.perf_counter() <= self.metrics_collector.last_log_time + 30
|
|
||||||
):
|
|
||||||
return
|
|
||||||
|
|
||||||
self.pool_stats_observer.get_pool_stats().update_scheduler_stats(self.stats)
|
|
||||||
self.stats.num_streaming_sessions = (
|
|
||||||
self.pool_stats_observer.streaming_session_count()
|
|
||||||
)
|
|
||||||
self.stats.streaming_session_held_tokens = (
|
|
||||||
self.pool_stats_observer.session_held_tokens()
|
|
||||||
)
|
|
||||||
|
|
||||||
priority_enabled = self.enable_priority_scheduling
|
|
||||||
self.stats.num_running_reqs = QueueCount.from_reqs(
|
|
||||||
self.running_batch.reqs, priority_enabled
|
|
||||||
)
|
|
||||||
self.stats.gen_throughput = 0
|
|
||||||
self.stats.num_queue_reqs = QueueCount.from_reqs(
|
|
||||||
self.waiting_queue, priority_enabled
|
|
||||||
)
|
|
||||||
self.stats.num_grammar_queue_reqs = len(self.grammar_manager)
|
|
||||||
if self.disaggregation_mode == DisaggregationMode.PREFILL:
|
|
||||||
self.stats.num_prefill_bootstrap_queue_reqs = QueueCount.from_reqs(
|
|
||||||
self.disagg_prefill_bootstrap_queue.queue, priority_enabled
|
|
||||||
)
|
|
||||||
self.stats.num_prefill_inflight_queue_reqs = QueueCount.from_reqs(
|
|
||||||
self.disagg_prefill_inflight_queue, priority_enabled
|
|
||||||
)
|
|
||||||
if self.disaggregation_mode == DisaggregationMode.DECODE:
|
|
||||||
self.stats.num_decode_prealloc_queue_reqs = QueueCount.from_reqs(
|
|
||||||
self.disagg_decode_prealloc_queue.queue, priority_enabled
|
|
||||||
)
|
|
||||||
self.stats.num_decode_transfer_queue_reqs = QueueCount.from_reqs(
|
|
||||||
self.disagg_decode_transfer_queue.queue, priority_enabled
|
|
||||||
)
|
|
||||||
self.metrics_collector.log_stats(self.stats)
|
|
||||||
Reference in New Issue
Block a user