diff --git a/python/sglang/srt/managers/data_parallel_controller.py b/python/sglang/srt/managers/data_parallel_controller.py index 4f297a32d..9790649a9 100644 --- a/python/sglang/srt/managers/data_parallel_controller.py +++ b/python/sglang/srt/managers/data_parallel_controller.py @@ -37,6 +37,7 @@ from sglang.srt.managers.io_struct import ( ) from sglang.srt.managers.schedule_batch import Req, RequestStage 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 ( DP_ATTENTION_HANDSHAKE_PORT_DELTA, PortArgs, @@ -207,6 +208,9 @@ class DataParallelController: 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): for worker in self.workers: worker.send_pyobj(obj) diff --git a/python/sglang/srt/managers/detokenizer_manager.py b/python/sglang/srt/managers/detokenizer_manager.py index 60a317b40..33cbacfa2 100644 --- a/python/sglang/srt/managers/detokenizer_manager.py +++ b/python/sglang/srt/managers/detokenizer_manager.py @@ -34,6 +34,7 @@ from sglang.srt.managers.io_struct import ( FreezeGCReq, ) 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.utils import ( configure_logger, @@ -87,6 +88,9 @@ class DetokenizerManager(MultiHttpWorkerDetokenizerMixin): # Init running status self.init_running_status(server_args) + if server_args.enable_metrics: + start_cpu_monitor_thread("detokenizer") + # Init dispatcher self.init_request_dispatcher() diff --git a/python/sglang/srt/managers/tokenizer_manager.py b/python/sglang/srt/managers/tokenizer_manager.py index 37ad13227..a433a0597 100644 --- a/python/sglang/srt/managers/tokenizer_manager.py +++ b/python/sglang/srt/managers/tokenizer_manager.py @@ -80,6 +80,7 @@ from sglang.srt.managers.tokenizer_manager_multiitem_mixin import ( TokenizerManagerMultiItemMixin, ) 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.server_args import ( PortArgs, @@ -216,6 +217,9 @@ class TokenizerManager(TokenizerCommunicatorMixin, TokenizerManagerMultiItemMixi # Init metric collector and watchdog self.init_metric_collector_watchdog() + if self.enable_metrics: + start_cpu_monitor_thread("tokenizer") + # Init request dispatcher self.init_request_dispatcher() diff --git a/python/sglang/srt/metrics/cpu_monitor.py b/python/sglang/srt/metrics/cpu_monitor.py new file mode 100644 index 000000000..783ed2f18 --- /dev/null +++ b/python/sglang/srt/metrics/cpu_monitor.py @@ -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 diff --git a/test/registered/metrics/test_cpu_monitor.py b/test/registered/metrics/test_cpu_monitor.py new file mode 100644 index 000000000..d12142f64 --- /dev/null +++ b/test/registered/metrics/test_cpu_monitor.py @@ -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() diff --git a/test/registered/metrics/test_metrics.py b/test/registered/metrics/test_metrics.py index 31a38b8d0..9b0f9eee8 100644 --- a/test/registered/metrics/test_metrics.py +++ b/test/registered/metrics/test_metrics.py @@ -180,6 +180,7 @@ class TestEnableMetrics(CustomTestCase): ("sglang:realtime_tokens_total", {"mode": "decode"}), ("sglang:gpu_execution_seconds_total", {"category": "forward_extend"}), ("sglang:gpu_execution_seconds_total", {"category": "forward_decode"}), + ("sglang:process_cpu_seconds_total", {"component": "tokenizer"}), ] _check_metrics_positive(self, metrics, metrics_to_check)