Tiny add CPU resource monitoring for overload diagnosis (#16852)
This commit is contained in:
@@ -37,6 +37,7 @@ from sglang.srt.managers.io_struct import (
|
|||||||
)
|
)
|
||||||
from sglang.srt.managers.schedule_batch import Req, RequestStage
|
from sglang.srt.managers.schedule_batch import Req, RequestStage
|
||||||
from sglang.srt.managers.scheduler import run_scheduler_process
|
from sglang.srt.managers.scheduler import run_scheduler_process
|
||||||
|
from sglang.srt.metrics.cpu_monitor import start_cpu_monitor_thread
|
||||||
from sglang.srt.server_args import (
|
from sglang.srt.server_args import (
|
||||||
DP_ATTENTION_HANDSHAKE_PORT_DELTA,
|
DP_ATTENTION_HANDSHAKE_PORT_DELTA,
|
||||||
PortArgs,
|
PortArgs,
|
||||||
@@ -207,6 +208,9 @@ class DataParallelController:
|
|||||||
test_stuck_time=envs.SGLANG_TEST_STUCK_DP_CONTROLLER.get(),
|
test_stuck_time=envs.SGLANG_TEST_STUCK_DP_CONTROLLER.get(),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
if server_args.enable_metrics:
|
||||||
|
start_cpu_monitor_thread("data_parallel_controller")
|
||||||
|
|
||||||
def send_to_all_workers(self, obj):
|
def send_to_all_workers(self, obj):
|
||||||
for worker in self.workers:
|
for worker in self.workers:
|
||||||
worker.send_pyobj(obj)
|
worker.send_pyobj(obj)
|
||||||
|
|||||||
@@ -34,6 +34,7 @@ from sglang.srt.managers.io_struct import (
|
|||||||
FreezeGCReq,
|
FreezeGCReq,
|
||||||
)
|
)
|
||||||
from sglang.srt.managers.multi_tokenizer_mixin import MultiHttpWorkerDetokenizerMixin
|
from sglang.srt.managers.multi_tokenizer_mixin import MultiHttpWorkerDetokenizerMixin
|
||||||
|
from sglang.srt.metrics.cpu_monitor import start_cpu_monitor_thread
|
||||||
from sglang.srt.server_args import PortArgs, ServerArgs
|
from sglang.srt.server_args import PortArgs, ServerArgs
|
||||||
from sglang.srt.utils import (
|
from sglang.srt.utils import (
|
||||||
configure_logger,
|
configure_logger,
|
||||||
@@ -87,6 +88,9 @@ class DetokenizerManager(MultiHttpWorkerDetokenizerMixin):
|
|||||||
# Init running status
|
# Init running status
|
||||||
self.init_running_status(server_args)
|
self.init_running_status(server_args)
|
||||||
|
|
||||||
|
if server_args.enable_metrics:
|
||||||
|
start_cpu_monitor_thread("detokenizer")
|
||||||
|
|
||||||
# Init dispatcher
|
# Init dispatcher
|
||||||
self.init_request_dispatcher()
|
self.init_request_dispatcher()
|
||||||
|
|
||||||
|
|||||||
@@ -80,6 +80,7 @@ from sglang.srt.managers.tokenizer_manager_multiitem_mixin import (
|
|||||||
TokenizerManagerMultiItemMixin,
|
TokenizerManagerMultiItemMixin,
|
||||||
)
|
)
|
||||||
from sglang.srt.metrics.collector import TokenizerMetricsCollector
|
from sglang.srt.metrics.collector import TokenizerMetricsCollector
|
||||||
|
from sglang.srt.metrics.cpu_monitor import start_cpu_monitor_thread
|
||||||
from sglang.srt.sampling.sampling_params import SamplingParams
|
from sglang.srt.sampling.sampling_params import SamplingParams
|
||||||
from sglang.srt.server_args import (
|
from sglang.srt.server_args import (
|
||||||
PortArgs,
|
PortArgs,
|
||||||
@@ -216,6 +217,9 @@ class TokenizerManager(TokenizerCommunicatorMixin, TokenizerManagerMultiItemMixi
|
|||||||
# Init metric collector and watchdog
|
# Init metric collector and watchdog
|
||||||
self.init_metric_collector_watchdog()
|
self.init_metric_collector_watchdog()
|
||||||
|
|
||||||
|
if self.enable_metrics:
|
||||||
|
start_cpu_monitor_thread("tokenizer")
|
||||||
|
|
||||||
# Init request dispatcher
|
# Init request dispatcher
|
||||||
self.init_request_dispatcher()
|
self.init_request_dispatcher()
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,31 @@
|
|||||||
|
import threading
|
||||||
|
import time
|
||||||
|
|
||||||
|
import psutil
|
||||||
|
|
||||||
|
|
||||||
|
def start_cpu_monitor_thread(component: str, interval: float = 5.0) -> threading.Thread:
|
||||||
|
from prometheus_client import Counter
|
||||||
|
|
||||||
|
cpu_seconds_total = Counter(
|
||||||
|
name="sglang:process_cpu_seconds_total",
|
||||||
|
documentation="Total CPU time consumed by this process (user + system)",
|
||||||
|
labelnames=["component"],
|
||||||
|
)
|
||||||
|
|
||||||
|
def monitor():
|
||||||
|
process = psutil.Process()
|
||||||
|
last_times = process.cpu_times()
|
||||||
|
|
||||||
|
while True:
|
||||||
|
time.sleep(interval)
|
||||||
|
curr_times = process.cpu_times()
|
||||||
|
delta = (curr_times.user - last_times.user) + (
|
||||||
|
curr_times.system - last_times.system
|
||||||
|
)
|
||||||
|
cpu_seconds_total.labels(component=component).inc(delta)
|
||||||
|
last_times = curr_times
|
||||||
|
|
||||||
|
t = threading.Thread(target=monitor, daemon=True)
|
||||||
|
t.start()
|
||||||
|
return t
|
||||||
@@ -0,0 +1,38 @@
|
|||||||
|
import time
|
||||||
|
import unittest
|
||||||
|
|
||||||
|
from sglang.test.ci.ci_register import register_cpu_ci
|
||||||
|
|
||||||
|
register_cpu_ci(est_time=60, suite="default", nightly=True)
|
||||||
|
|
||||||
|
|
||||||
|
class TestCpuMonitor(unittest.TestCase):
|
||||||
|
def test_cpu_monitor(self):
|
||||||
|
from prometheus_client import REGISTRY
|
||||||
|
|
||||||
|
from sglang.srt.metrics.cpu_monitor import start_cpu_monitor_thread
|
||||||
|
|
||||||
|
thread = start_cpu_monitor_thread("test", interval=0.1)
|
||||||
|
self.assertTrue(thread.is_alive())
|
||||||
|
self.assertTrue(thread.daemon)
|
||||||
|
|
||||||
|
end_time = time.monotonic() + 0.3
|
||||||
|
while time.monotonic() < end_time:
|
||||||
|
_ = sum(i * i for i in range(1000))
|
||||||
|
time.sleep(0.2)
|
||||||
|
|
||||||
|
value = None
|
||||||
|
for metric in REGISTRY.collect():
|
||||||
|
for sample in metric.samples:
|
||||||
|
if (
|
||||||
|
sample.name == "sglang:process_cpu_seconds_total"
|
||||||
|
and sample.labels.get("component") == "test"
|
||||||
|
):
|
||||||
|
value = sample.value
|
||||||
|
print(f"sglang:process_cpu_seconds_total = {value}")
|
||||||
|
self.assertIsNotNone(value)
|
||||||
|
self.assertGreater(value, 0)
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
@@ -180,6 +180,7 @@ class TestEnableMetrics(CustomTestCase):
|
|||||||
("sglang:realtime_tokens_total", {"mode": "decode"}),
|
("sglang:realtime_tokens_total", {"mode": "decode"}),
|
||||||
("sglang:gpu_execution_seconds_total", {"category": "forward_extend"}),
|
("sglang:gpu_execution_seconds_total", {"category": "forward_extend"}),
|
||||||
("sglang:gpu_execution_seconds_total", {"category": "forward_decode"}),
|
("sglang:gpu_execution_seconds_total", {"category": "forward_decode"}),
|
||||||
|
("sglang:process_cpu_seconds_total", {"component": "tokenizer"}),
|
||||||
]
|
]
|
||||||
_check_metrics_positive(self, metrics, metrics_to_check)
|
_check_metrics_positive(self, metrics, metrics_to_check)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user