Tiny add scheduler status logging (#16872)
This commit is contained in:
@@ -163,6 +163,8 @@ class Envs:
|
|||||||
SGLANG_LOG_MS = EnvBool(False)
|
SGLANG_LOG_MS = EnvBool(False)
|
||||||
SGLANG_DISABLE_REQUEST_LOGGING = EnvBool(False)
|
SGLANG_DISABLE_REQUEST_LOGGING = EnvBool(False)
|
||||||
SGLANG_LOG_REQUEST_EXCEEDED_MS = EnvInt(-1)
|
SGLANG_LOG_REQUEST_EXCEEDED_MS = EnvInt(-1)
|
||||||
|
SGLANG_LOG_SCHEDULER_STATUS_TARGET = EnvStr("")
|
||||||
|
SGLANG_LOG_SCHEDULER_STATUS_INTERVAL = EnvFloat(60.0)
|
||||||
|
|
||||||
# SGLang CI
|
# SGLang CI
|
||||||
SGLANG_IS_IN_CI = EnvBool(False)
|
SGLANG_IS_IN_CI = EnvBool(False)
|
||||||
|
|||||||
@@ -20,6 +20,7 @@ from sglang.srt.metrics.collector import (
|
|||||||
)
|
)
|
||||||
from sglang.srt.utils import get_bool_env_var
|
from sglang.srt.utils import get_bool_env_var
|
||||||
from sglang.srt.utils.device_timer import DeviceTimer
|
from sglang.srt.utils.device_timer import DeviceTimer
|
||||||
|
from sglang.srt.utils.scheduler_status_logger import SchedulerStatusLogger
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from sglang.srt.managers.scheduler import EmbeddingBatchResult, Scheduler
|
from sglang.srt.managers.scheduler import EmbeddingBatchResult, Scheduler
|
||||||
@@ -112,6 +113,8 @@ class SchedulerMetricsMixin:
|
|||||||
if self.enable_kv_cache_events:
|
if self.enable_kv_cache_events:
|
||||||
self.init_kv_events(self.server_args.kv_events_config)
|
self.init_kv_events(self.server_args.kv_events_config)
|
||||||
|
|
||||||
|
self.scheduler_status_logger = SchedulerStatusLogger.maybe_create()
|
||||||
|
|
||||||
def init_kv_events(self: Scheduler, kv_events_config: Optional[str]):
|
def init_kv_events(self: Scheduler, kv_events_config: Optional[str]):
|
||||||
if self.enable_kv_cache_events:
|
if self.enable_kv_cache_events:
|
||||||
self.kv_event_publisher = EventPublisherFactory.create(
|
self.kv_event_publisher = EventPublisherFactory.create(
|
||||||
@@ -447,6 +450,9 @@ class SchedulerMetricsMixin:
|
|||||||
dp_cooperation_info=batch.dp_cooperation_info,
|
dp_cooperation_info=batch.dp_cooperation_info,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
if x := self.scheduler_status_logger:
|
||||||
|
x.maybe_dump(batch, self.waiting_queue)
|
||||||
|
|
||||||
def log_batch_result_stats(
|
def log_batch_result_stats(
|
||||||
self: Scheduler,
|
self: Scheduler,
|
||||||
batch: ScheduleBatch,
|
batch: ScheduleBatch,
|
||||||
|
|||||||
@@ -0,0 +1,49 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import time
|
||||||
|
from typing import TYPE_CHECKING, List, Optional
|
||||||
|
|
||||||
|
import torch.distributed as dist
|
||||||
|
|
||||||
|
from sglang.srt.environ import envs
|
||||||
|
from sglang.srt.utils.log_utils import create_log_targets, log_json
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
from sglang.srt.managers.schedule_batch import Req, ScheduleBatch
|
||||||
|
|
||||||
|
|
||||||
|
class SchedulerStatusLogger:
|
||||||
|
def __init__(self, targets: List[str], dump_interval: float):
|
||||||
|
self.loggers = create_log_targets(targets=targets, name_prefix=__name__)
|
||||||
|
self.dump_interval = dump_interval
|
||||||
|
self.last_dump_time = 0.0
|
||||||
|
self.rank = dist.get_rank() if dist.is_initialized() else 0
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def maybe_create() -> Optional["SchedulerStatusLogger"]:
|
||||||
|
target = envs.SGLANG_LOG_SCHEDULER_STATUS_TARGET.get()
|
||||||
|
if not target:
|
||||||
|
return None
|
||||||
|
|
||||||
|
return SchedulerStatusLogger(
|
||||||
|
targets=[t.strip() for t in target.split(",") if t.strip()],
|
||||||
|
dump_interval=envs.SGLANG_LOG_SCHEDULER_STATUS_INTERVAL.get(),
|
||||||
|
)
|
||||||
|
|
||||||
|
def maybe_dump(
|
||||||
|
self, running_batch: "ScheduleBatch", waiting_queue: List["Req"]
|
||||||
|
) -> None:
|
||||||
|
now = time.time()
|
||||||
|
if now - self.last_dump_time < self.dump_interval:
|
||||||
|
return
|
||||||
|
|
||||||
|
self.last_dump_time = now
|
||||||
|
log_json(
|
||||||
|
self.loggers,
|
||||||
|
"scheduler.status",
|
||||||
|
{
|
||||||
|
"rank": self.rank,
|
||||||
|
"running_rids": [r.rid for r in running_batch.reqs],
|
||||||
|
"queued_rids": [r.rid for r in waiting_queue],
|
||||||
|
},
|
||||||
|
)
|
||||||
@@ -0,0 +1,73 @@
|
|||||||
|
import json
|
||||||
|
import os
|
||||||
|
import shutil
|
||||||
|
import tempfile
|
||||||
|
import time
|
||||||
|
import unittest
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import requests
|
||||||
|
|
||||||
|
from sglang.srt.utils import kill_process_tree
|
||||||
|
from sglang.test.ci.ci_register import register_cuda_ci
|
||||||
|
from sglang.test.test_utils import (
|
||||||
|
DEFAULT_TIMEOUT_FOR_SERVER_LAUNCH,
|
||||||
|
DEFAULT_URL_FOR_TEST,
|
||||||
|
CustomTestCase,
|
||||||
|
popen_launch_server,
|
||||||
|
)
|
||||||
|
|
||||||
|
register_cuda_ci(est_time=120, suite="nightly-1-gpu", nightly=True)
|
||||||
|
|
||||||
|
|
||||||
|
class TestSchedulerStatusLogger(CustomTestCase):
|
||||||
|
@classmethod
|
||||||
|
def setUpClass(cls):
|
||||||
|
cls.temp_dir = tempfile.mkdtemp()
|
||||||
|
cls.addClassCleanup(shutil.rmtree, cls.temp_dir)
|
||||||
|
env = os.environ.copy()
|
||||||
|
env["SGLANG_LOG_SCHEDULER_STATUS_TARGET"] = cls.temp_dir
|
||||||
|
env["SGLANG_LOG_SCHEDULER_STATUS_INTERVAL"] = "1"
|
||||||
|
cls.process = popen_launch_server(
|
||||||
|
"Qwen/Qwen3-0.6B",
|
||||||
|
DEFAULT_URL_FOR_TEST,
|
||||||
|
timeout=DEFAULT_TIMEOUT_FOR_SERVER_LAUNCH,
|
||||||
|
other_args=["--skip-server-warmup"],
|
||||||
|
env=env,
|
||||||
|
)
|
||||||
|
cls.addClassCleanup(kill_process_tree, cls.process.pid)
|
||||||
|
|
||||||
|
def test_scheduler_status_dump(self):
|
||||||
|
response = requests.post(
|
||||||
|
DEFAULT_URL_FOR_TEST + "/generate",
|
||||||
|
json={
|
||||||
|
"text": "Hello",
|
||||||
|
"sampling_params": {"max_new_tokens": 8, "temperature": 0},
|
||||||
|
},
|
||||||
|
timeout=30,
|
||||||
|
)
|
||||||
|
self.assertEqual(response.status_code, 200)
|
||||||
|
|
||||||
|
time.sleep(2)
|
||||||
|
|
||||||
|
events = list(_find_log_events(self.temp_dir, "scheduler.status"))
|
||||||
|
print(f"{events=}")
|
||||||
|
self.assertGreater(len(events), 0, "scheduler.status event not found")
|
||||||
|
data = events[0]
|
||||||
|
for field in ["timestamp", "rank", "running_rids", "queued_rids"]:
|
||||||
|
self.assertIn(field, data)
|
||||||
|
self.assertIsInstance(data["running_rids"], list)
|
||||||
|
self.assertIsInstance(data["queued_rids"], list)
|
||||||
|
|
||||||
|
|
||||||
|
def _find_log_events(log_dir: str, event_name: str):
|
||||||
|
for f in Path(log_dir).glob("*.log"):
|
||||||
|
for line in f.read_text().splitlines():
|
||||||
|
if line.startswith("{"):
|
||||||
|
data = json.loads(line)
|
||||||
|
if data.get("event") == event_name:
|
||||||
|
yield data
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
Reference in New Issue
Block a user