From 8441857fe953ef3323f9d2b31597b1d3a092ca77 Mon Sep 17 00:00:00 2001 From: Ye Zhang Date: Mon, 7 Jul 2025 14:13:05 +0800 Subject: [PATCH 01/12] feat: core metric impl Signed-off-by: Ye Zhang rebase in progress; onto 82d03ca97 --- tensorrt_llm/executor/postproc_worker.py | 19 ++++-- tensorrt_llm/executor/proxy.py | 26 ++++++- tensorrt_llm/executor/result.py | 61 +++++++++++++++-- tensorrt_llm/executor/worker.py | 26 +++++-- tensorrt_llm/llmapi/llm_args.py | 5 ++ tensorrt_llm/metrics/__init__.py | 7 ++ tensorrt_llm/metrics/collector.py | 86 ++++++++++++++++++++++++ tensorrt_llm/metrics/enums.py | 15 +++++ tensorrt_llm/sampling_params.py | 7 +- tensorrt_llm/serve/openai_server.py | 27 ++++++++ tensorrt_llm/utils/__init__.py | 0 tensorrt_llm/utils/utils.py | 22 ++++++ 12 files changed, 280 insertions(+), 21 deletions(-) create mode 100644 tensorrt_llm/metrics/__init__.py create mode 100644 tensorrt_llm/metrics/collector.py create mode 100644 tensorrt_llm/metrics/enums.py create mode 100644 tensorrt_llm/utils/__init__.py create mode 100644 tensorrt_llm/utils/utils.py diff --git a/tensorrt_llm/executor/postproc_worker.py b/tensorrt_llm/executor/postproc_worker.py index 2e5a3cd29673..c358168e58c9 100644 --- a/tensorrt_llm/executor/postproc_worker.py +++ b/tensorrt_llm/executor/postproc_worker.py @@ -3,7 +3,7 @@ from collections import deque from dataclasses import dataclass from typing import (TYPE_CHECKING, Any, Callable, Dict, List, NamedTuple, - Optional) + Optional, Union) import zmq import zmq.asyncio @@ -18,7 +18,7 @@ if TYPE_CHECKING: from .result import (DetokenizedGenerationResultBase, GenerationResult, - GenerationResultBase) + GenerationResultBase, ResponseWrapper) __all__ = [ "PostprocWorker", @@ -57,7 +57,7 @@ class PostprocWorker: @dataclass class Input: - rsp: "tllm.Response" + rsp: Union["tllm.Response", "ResponseWrapper"] # The information necessary for creating a GenerationResult in the first Input for each request sampling_params: Optional[SamplingParams] = None @@ -69,6 +69,7 @@ class Output(NamedTuple): res: Any is_final: bool error: str = "" + metrics: Optional[dict[str, float]] = None def __init__( self, @@ -118,7 +119,10 @@ def default_record_creator( streaming=inp.streaming, tokenizer=tokenizer) - async def _handle_input(self, input: "PostprocWorker.Input") -> Any: + async def _handle_input( + self, + input: Union["PostprocWorker.Input", "ResponseWrapper"] + ) -> [Any, Optional[dict[str, float]]]: ''' Handle a single response from await_response worker. ''' if input.rsp.result.context_logits is not None or \ input.rsp.result.generation_logits is not None: @@ -139,6 +143,7 @@ async def _handle_input(self, input: "PostprocWorker.Input") -> Any: record._handle_response(input.rsp) # inplace # Left the result_handler determine the final output dtype. # NOTE: This will change the CompletionOutput._postprocess_result + metrics_dict = record.metrics_dict if postproc_params := record.postproc_params: result_handler, args = postproc_params.post_processor, postproc_params.postproc_args args.tokenizer = self._tokenizer @@ -150,7 +155,7 @@ async def _handle_input(self, input: "PostprocWorker.Input") -> Any: # TODO: Keep only the diff token_ids and text in streaming mode when # result_handler is not set - return out + return out, metrics_dict async def _batched_put(self): ''' Batched IPC send. ''' @@ -173,8 +178,8 @@ async def handle_single_input(inp: PostprocWorker.Input, client_id = inp.rsp.client_id is_final = inp.rsp.result.is_final if is_llm_response( inp.rsp) else True - res = await self._handle_input(inp) - batch.append(PostprocWorker.Output(client_id, res, is_final)) + res, metrics = await self._handle_input(inp) + batch.append(PostprocWorker.Output(client_id=client_id, res=res, is_final=is_final, metrics=metrics)) if is_final: self._records.pop(client_id) diff --git a/tensorrt_llm/executor/proxy.py b/tensorrt_llm/executor/proxy.py index 1cb86dfdff71..703001f6a3cd 100644 --- a/tensorrt_llm/executor/proxy.py +++ b/tensorrt_llm/executor/proxy.py @@ -10,8 +10,10 @@ import zmq.asyncio from tensorrt_llm.logger import logger +from tensorrt_llm.metrics.collector import MetricsCollector from .._utils import customized_gc_thresholds, mpi_rank, nvtx_range_debug +from ..utils.utils import set_prometheus_multiproc_dir from ..llmapi.mpi_session import (MpiCommSession, MpiPoolSession, MpiSession, RemoteMpiCommSessionClient) from ..llmapi.tracer import enable_llm_tracer, get_tracer, global_tracer @@ -20,9 +22,8 @@ print_colored_debug) from .executor import GenerationExecutor from .ipc import FusedIpcQueue, IpcQueue -from .postproc_worker import PostprocWorkerConfig -from .request import CancellingRequest, GenerationRequest -from .result import GenerationResult, IterationResult +from .postproc_worker import PostprocWorker, PostprocWorkerConfig +from .result import GenerationResult, IterationResult, ResponseWrapper from .utils import (ErrorResponse, IntraProcessQueue, WorkerCommIpcAddrs, create_mpi_comm_session, get_spawn_proxy_process_env, is_llm_response, print_alive_threads) @@ -60,6 +61,9 @@ def __init__( self.workers_started = False self.worker_cls = worker_cls + set_prometheus_multiproc_dir() + self.metrics_collector = MetricsCollector({"model_name": "undefined", "engine_type": "undefined"}) + mpi_process_pre_spawned: bool = get_spawn_proxy_process_env() if mpi_session is None: @@ -178,6 +182,22 @@ def process_res(res): else: queue.put(res) + # log metrics to prometheus when response is finished and return_perf_metrics is on + if isinstance(res, ResponseWrapper): + if res._response and res._response.result and res._response.result.is_final: + self._results[client_id]._handle_response(res) + if metrics_dict := self._results[client_id].metrics_dict: + if finish_reason := metrics_dict.get(MetricsCollector.labelname_finish_reason): + self.metrics_collector.log_request_success(1, {MetricsCollector.labelname_finish_reason: + finish_reason}) + self.metrics_collector.log_histogram(metrics_dict) + if isinstance(res, PostprocWorker.Output): + if metrics_dict := res.metrics: + if finish_reason := metrics_dict.get(MetricsCollector.labelname_finish_reason): + self.metrics_collector.log_request_success(1, {MetricsCollector.labelname_finish_reason: + finish_reason}) + self.metrics_collector.log_histogram(metrics_dict) + if (is_llm_response(res) and res.result.is_final) or isinstance( res, ErrorResponse): self._results.pop(client_id) diff --git a/tensorrt_llm/executor/result.py b/tensorrt_llm/executor/result.py index 0408a6c757c0..5efe89b6d3fd 100644 --- a/tensorrt_llm/executor/result.py +++ b/tensorrt_llm/executor/result.py @@ -15,6 +15,7 @@ from ..disaggregated_params import DisaggregatedParams from ..llmapi.tracer import global_tracer from ..llmapi.utils import AsyncQueue +from ..metrics import SupportedMetricNames, RequestEventTiming, MetricsCollector from ..sampling_params import LogprobParams, SamplingParams from .utils import ErrorResponse, has_event_loop, is_llm_response @@ -50,14 +51,18 @@ class LogProbsResult(NamedTuple): class ResponseWrapper: - """Wrapper of runtime response with optional outputs computed post runtime. + """ + 1. Wrapper of runtime response with optional outputs computed post runtime. + 2. a workaround to pass around RequestPerfMetrics. """ def __init__(self, response: Union["PostprocWorker.Output", tllm.Response], - logprobs: Optional[LogProbsResult] = None): + logprobs: Optional[LogProbsResult] = None, + request_perf_metrics: Optional[dict[str, float]] = None): self._response = response self.logprobs = logprobs + self.request_perf_metrics = request_perf_metrics @property def _is_llm_response(self): @@ -68,6 +73,14 @@ def __getattr__(self, name): response = object.__getattribute__(self, '_response') return getattr(response, name) + def __getstate__(self): + return (self._response, self.logprobs, self.request_perf_metrics) + + def __setstate__(self, state): + self._response = state[0] + self.logprobs = state[1] + self.request_perf_metrics = state[2] + @dataclass(slots=True) class CompletionOutput: @@ -146,6 +159,7 @@ def __init__(self, self.disaggregated_params = None self.decoding_iter = 0 self._done = False + self.metrics_dict = {} if has_event_loop(): self.aqueue = AsyncQueue() @@ -201,7 +215,8 @@ def _handle_sequence(self, finish_reasons, response_tensors, sequence_index, - logprobs_result=None): + logprobs_result=None, + req_perf_metrics_dict: Optional[dict[str, float]] = None): """ Handle a single sequence in the response. """ seq_idx = sequence_index @@ -271,6 +286,7 @@ def _handle_sequence(self, else: raise ValueError( f"Unknown finish reason: {finish_reasons[src_idx]}") + self.record_stats(output, req_perf_metrics_dict) @nvtx_range_debug("handle_response", color="red", @@ -278,7 +294,9 @@ def _handle_sequence(self, def _handle_response(self, response: Union["PostprocWorker.Output", tllm.Response, ResponseWrapper, ErrorResponse]): + req_perf_metrics_dict = None if isinstance(response, ResponseWrapper): + req_perf_metrics_dict = response.request_perf_metrics logprobs_result = response.logprobs response = response._response else: @@ -322,11 +340,11 @@ def _handle_response(self, if self.sampling_params.use_beam_search: for beam_idx, _ in enumerate(response_result.output_token_ids): self._handle_sequence(finish_reasons, response_result, - beam_idx, logprobs_result) + beam_idx, logprobs_result, req_perf_metrics_dict) else: self._handle_sequence(finish_reasons, response_result, response_result.sequence_index, - logprobs_result) + logprobs_result, req_perf_metrics_dict) if response_result.context_logits is not None: self._context_logits = response_result.context_logits @@ -342,6 +360,16 @@ def _handle_response(self, else: raise ValueError(f"Unknown response type: {response}") + def record_stats(self, output: CompletionOutput, stats: Optional[dict[str, float]] = None): + if not stats: + return + metrics_stats = {} + if output.finish_reason: + metrics_stats.update({MetricsCollector.labelname_finish_reason: output.finish_reason}) + processed_metrics_stat = _process_req_perf_metrics(stats, len(output.token_ids)) + if processed_metrics_stat: + metrics_stats.update(processed_metrics_stat) + self.metrics_dict = metrics_stats class DetokenizedGenerationResultBase(GenerationResultBase): ''' The base class for the generation result with detokenization support. ''' @@ -688,3 +716,26 @@ def _topk_logprobs(logits: torch.Tensor, top_k: int, return LogProbsResult(prompt=prompt_logprobs, generation=generation_logprobs) + +def _process_req_perf_metrics(req_perf_metrics_dict: Optional[dict[str, float]], + output_length: int) -> Optional[dict[str, float]]: + stat = {} + if not req_perf_metrics_dict: + return stat + ttft = req_perf_metrics_dict.get(RequestEventTiming.FIRST_TOKEN_TIME, 0) - \ + req_perf_metrics_dict.get(RequestEventTiming.ARRIVAL_TIME, 0) + e2e = req_perf_metrics_dict.get(RequestEventTiming.LAST_TOKEN_TIME, 0) - \ + req_perf_metrics_dict.get(RequestEventTiming.ARRIVAL_TIME, 0) + request_queue_time = req_perf_metrics_dict.get(RequestEventTiming.FIRST_SCHEDULED_TIME, 0) - \ + req_perf_metrics_dict.get(RequestEventTiming.ARRIVAL_TIME, 0) + stat = {SupportedMetricNames.TTFT: ttft, + SupportedMetricNames.E2E: e2e, + SupportedMetricNames.REQUEST_QUEUE_TIME: request_queue_time} + if output_length > 1: + tpot = (req_perf_metrics_dict.get(RequestEventTiming.LAST_TOKEN_TIME, 0) - + req_perf_metrics_dict.get(RequestEventTiming.FIRST_TOKEN_TIME, 0)) / (output_length - 1) + stat.update({SupportedMetricNames.TPOT: tpot}) + stat = dict( + filter(lambda item: item[1] > 0, stat.items()) + ) + return stat diff --git a/tensorrt_llm/executor/worker.py b/tensorrt_llm/executor/worker.py index db8d84fcc898..75874e0a5591 100644 --- a/tensorrt_llm/executor/worker.py +++ b/tensorrt_llm/executor/worker.py @@ -25,6 +25,7 @@ clear_sched_affinity, print_colored_debug, print_traceback_on_error) from ..lora_manager import LoraConfig, LoraManager +from ..metrics import RequestEventTiming from ..prompt_adapter_manager import PromptAdapterManager from ..runtime import ModelConfig from ..runtime.model_runner import _engine_config_to_model_config @@ -472,7 +473,7 @@ def _deduce_max_tokens(request: GenerationRequest, request.sampling_params.end_id, pad_id=request.sampling_params.pad_id, output_config=request.sampling_params._get_output_config( - is_pytorch_backend=self._is_pytorch_backend), + is_pytorch_backend=self._is_pytorch_backend, return_perf_metrics=request.return_perf_metrics), # Beam search enforces return_all_generated_tokens=True regardless of the passed value return_all_generated_tokens=False, # convert python config into pybind config @@ -901,8 +902,9 @@ def handle_for_worker(self, responses: List[tllm.Response]) -> None: logprobs_result = _get_logprobs(self.worker, response, self.worker._is_pytorch_backend) - if logprobs_result: - response = ResponseWrapper(response, logprobs_result) + req_perf_metrics = _get_req_perf_metrics(response) + if logprobs_result or req_perf_metrics: + response = ResponseWrapper(response, logprobs_result, req_perf_metrics) # For AsyncQueue.sync_q, we will batch the events to avoid too many # event notifications, thus put without wait here. @@ -942,8 +944,9 @@ def handle_for_ipc_batched(self, responses: List[tllm.Response]) -> None: else: logprobs_result = _get_logprobs(self.worker, response, self.worker._is_pytorch_backend) - if logprobs_result: - response = ResponseWrapper(response, logprobs_result) + req_perf_metrics = _get_req_perf_metrics(response) + if logprobs_result or req_perf_metrics: + response = ResponseWrapper(response, logprobs_result, req_perf_metrics) _send_rsp(self.worker, response, @@ -1051,3 +1054,16 @@ def _send_rsp( worker._pop_result(response.client_id) else: raise ValueError(f"Unknown response type: {response}") + +def _get_req_perf_metrics(response: tllm.Response) -> Optional[dict[str, float]]: + req_perf_metrics, metrics_dict = None, {} + if response.result: + req_perf_metrics = response.result.request_perf_metrics + if req_perf_metrics and req_perf_metrics.timing_metrics: + metrics_dict = { + RequestEventTiming.ARRIVAL_TIME: req_perf_metrics.timing_metrics.arrival_time.total_seconds(), + RequestEventTiming.FIRST_TOKEN_TIME: req_perf_metrics.timing_metrics.first_token_time.total_seconds(), + RequestEventTiming.FIRST_SCHEDULED_TIME: + req_perf_metrics.timing_metrics.first_scheduled_time.total_seconds(), + RequestEventTiming.LAST_TOKEN_TIME: req_perf_metrics.timing_metrics.last_token_time.total_seconds()} + return metrics_dict diff --git a/tensorrt_llm/llmapi/llm_args.py b/tensorrt_llm/llmapi/llm_args.py index 1169a779be6e..258edbabeff0 100644 --- a/tensorrt_llm/llmapi/llm_args.py +++ b/tensorrt_llm/llmapi/llm_args.py @@ -1311,6 +1311,11 @@ class BaseLlmArgs(StrictBaseModel): status="deprecated", ) + return_perf_metrics: Optional[bool] = Field( + default=False, + description="optionally returns perf metrics", + ) + _parallel_config: Optional[object] = PrivateAttr(default=None) _model_format: Optional[_ModelFormatKind] = PrivateAttr(default=None) _speculative_model: Optional[str] = PrivateAttr(default=None) diff --git a/tensorrt_llm/metrics/__init__.py b/tensorrt_llm/metrics/__init__.py new file mode 100644 index 000000000000..4650a84709d5 --- /dev/null +++ b/tensorrt_llm/metrics/__init__.py @@ -0,0 +1,7 @@ +from .collector import * +from .enums import * +__all__ = [ + "SupportedMetricNames", + "MetricsCollector", + "RequestEventTiming" +] \ No newline at end of file diff --git a/tensorrt_llm/metrics/collector.py b/tensorrt_llm/metrics/collector.py new file mode 100644 index 000000000000..23616b83ae2b --- /dev/null +++ b/tensorrt_llm/metrics/collector.py @@ -0,0 +1,86 @@ +"""Utilities for Prometheus Metrics Collection.""" + +import time +from typing import Dict, Optional, Union + +from .enums import SupportedMetricNames + + +class MetricsCollector: + labelname_finish_reason = "finished_reason" + def __init__(self, labels: Dict[str, str]) -> None: + from prometheus_client import Counter, Histogram + self.last_log_time = time.time() + self.labels = labels + + self.finish_reason_label = {MetricsCollector.labelname_finish_reason: "unknown"} + self.labels_with_finished_reason = {**self.labels, **self.finish_reason_label} + + self.counter_request_success = Counter( + name="request_success_total", + documentation="Count of successfully processed requests.", + labelnames=self.labels_with_finished_reason.keys()) + + self.histogram_e2e_time_request = Histogram( + name="e2e_request_latency_seconds", + documentation="Histogram of end to end request latency in seconds.", + buckets=[ + 0.3, 0.5, 0.8, 1.0, 1.5, 2.0, 2.5, 5.0, 10.0, 15.0, 20.0, 30.0, + 40.0, 50.0, 60.0, 120.0, 240.0, 480.0, 960.0, 1920.0, 7680.0 + ], + labelnames=self.labels.keys()) + + self.histogram_time_to_first_token = Histogram( + name="time_to_first_token_seconds", + documentation="Histogram of time to first token in seconds.", + buckets=[ + 0.001, 0.005, 0.01, 0.02, 0.04, 0.06, 0.08, 0.1, 0.25, 0.5, + 0.75, 1.0, 2.5, 5.0, 7.5, 10.0, 20.0, 40.0, 80.0, 160.0, 640.0, + 2560.0 + ], + labelnames=self.labels.keys()) + + self.histogram_time_per_output_token = Histogram( + name="time_per_output_token_seconds", + documentation="Histogram of time per output token in seconds.", + buckets=[ + 0.01, 0.025, 0.05, 0.075, 0.1, 0.15, 0.2, 0.3, 0.4, 0.5, 0.75, + 1.0, 2.5, 5.0, 7.5, 10.0, 20.0, 40.0, 80.0 + ], + labelnames=self.labels.keys()) + + self.histogram_queue_time_request = Histogram( + name="request_queue_time_seconds", + documentation="Histogram of time spent in WAITING phase for request.", + buckets = [ + 0.3, 0.5, 0.8, 1.0, 1.5, 2.0, 2.5, 5.0, 10.0, 15.0, 20.0, 30.0, + 40.0, 50.0, 60.0, 120.0, 240.0, 480.0, 960.0, 1920.0, 7680.0], + labelnames=self.labels.keys()) + + def _label_merge(self, labels: Dict[str, str]) -> Dict[str, str]: + if labels is None or len(labels) == 0: + return self.labels + return {**self.labels, **labels} + + def _log_counter(self, counter, labels: Dict[str, str], data: Union[int, float]) -> None: + # Convenience function for logging to counter. + counter.labels(**self._label_merge(labels)).inc(data) + + def _log_histogram(self, histogram, data: Union[int, float]) -> None: + # Convenience function for logging to histogram. + histogram.labels(**self.labels).observe(data) + + def log_request_success(self, data: Union[int, float], labels: Dict[str, str]) -> None: + self._log_counter(self.counter_request_success, labels, data) + self.last_log_time = time.time() + + def log_histogram(self, data: Optional[dict[str, float]]) -> None: + if e2e := data.get(SupportedMetricNames.E2E, 0): + self._log_histogram(self.histogram_e2e_time_request, e2e) + if ttft := data.get(SupportedMetricNames.TTFT, 0): + self._log_histogram(self.histogram_time_to_first_token, ttft) + if tpot := data.get(SupportedMetricNames.TPOT, 0): + self._log_histogram(self.histogram_time_per_output_token, tpot) + if request_queue_time := data.get(SupportedMetricNames.REQUEST_QUEUE_TIME, 0): + self._log_histogram(self.histogram_queue_time_request, request_queue_time) + self.last_log_time = time.time() diff --git a/tensorrt_llm/metrics/enums.py b/tensorrt_llm/metrics/enums.py new file mode 100644 index 000000000000..1bab966106cc --- /dev/null +++ b/tensorrt_llm/metrics/enums.py @@ -0,0 +1,15 @@ +from enum import Enum + + +class SupportedMetricNames(Enum): + TTFT = "ttft" + TPOT = "tpot" + E2E = "e2e" + REQUEST_QUEUE_TIME = "request_queue_time" + + +class RequestEventTiming(Enum): + ARRIVAL_TIME = "arrival_time" + FIRST_TOKEN_TIME = "first_token_time" + FIRST_SCHEDULED_TIME = "first_scheduled_time" + LAST_TOKEN_TIME = "last_token_time" diff --git a/tensorrt_llm/sampling_params.py b/tensorrt_llm/sampling_params.py index 361c0fc0c0f4..f37d357bbee4 100644 --- a/tensorrt_llm/sampling_params.py +++ b/tensorrt_llm/sampling_params.py @@ -437,7 +437,11 @@ def _get_sampling_config(self) -> tllme.SamplingConfig: return tllme.SamplingConfig(**llmapi_to_rt_param_map) - def _get_output_config(self, is_pytorch_backend: bool = False) -> tllme.OutputConfig: + def _get_output_config( + self, + is_pytorch_backend: bool = False, + return_perf_metrics: Optional[bool] = False + ) -> tllme.OutputConfig: sampling_param_fields = set(dir(SamplingParams)) fields = [ f @@ -451,6 +455,7 @@ def _get_output_config(self, is_pytorch_backend: bool = False) -> tllme.OutputCo config_kwargs["return_log_probs"] = bool(self.logprobs) else: config_kwargs["return_log_probs"] = self._return_log_probs + config_kwargs["return_perf_metrics"] = return_perf_metrics return tllme.OutputConfig(**config_kwargs) diff --git a/tensorrt_llm/serve/openai_server.py b/tensorrt_llm/serve/openai_server.py index d90578ce36b3..f4001c0c2802 100644 --- a/tensorrt_llm/serve/openai_server.py +++ b/tensorrt_llm/serve/openai_server.py @@ -1,6 +1,7 @@ #!/usr/bin/env python import asyncio import os +import re import signal import traceback from contextlib import asynccontextmanager @@ -13,6 +14,7 @@ from fastapi import FastAPI, Request from fastapi.exceptions import RequestValidationError from fastapi.responses import JSONResponse, Response, StreamingResponse +from starlette.routing import Mount from transformers import AutoConfig, AutoProcessor from tensorrt_llm._tensorrt_engine import LLM @@ -151,6 +153,31 @@ def register_routes(self): self.app.add_api_route("/v1/chat/completions", self.openai_chat, methods=["POST"]) + # register /prometheus/metrics + self.mount_metrics() + + def mount_metrics(self): + # Lazy import for prometheus multiprocessing. + # We need to set PROMETHEUS_MULTIPROC_DIR environment variable + # before prometheus_client is imported. + # See https://prometheus.github.io/client_python/multiprocess/ + from prometheus_client import (CollectorRegistry, make_asgi_app, + multiprocess) + from prometheus_fastapi_instrumentator import Instrumentator + registry = CollectorRegistry() + multiprocess.MultiProcessCollector(registry) + Instrumentator( + should_group_status_codes=False, + should_respect_env_var=True, + excluded_handlers=[ + ".*" + ], + registry=registry, + ).add().instrument(self.app).expose(self.app) + metrics_app = make_asgi_app(registry=registry) + metrics_route = Mount("/prometheus/metrics", metrics_app) + metrics_route.path_regex = re.compile("^/prometheus/metrics(?P.*)$") + self.app.routes.append(metrics_route) async def health(self) -> Response: return Response(status_code=200) diff --git a/tensorrt_llm/utils/__init__.py b/tensorrt_llm/utils/__init__.py new file mode 100644 index 000000000000..e69de29bb2d1 diff --git a/tensorrt_llm/utils/utils.py b/tensorrt_llm/utils/utils.py new file mode 100644 index 000000000000..c52718603bc4 --- /dev/null +++ b/tensorrt_llm/utils/utils.py @@ -0,0 +1,22 @@ +import os +import tempfile + +from tensorrt_llm.logger import logger + + +def set_prometheus_multiproc_dir(): + # Set prometheus multiprocess directory + # sglang uses prometheus multiprocess mode + # we need to set this before importing prometheus_client + # https://prometheus.github.io/client_python/multiprocess/ + global prometheus_multiproc_dir + + if "PROMETHEUS_MULTIPROC_DIR" in os.environ: + logger.info("User set PROMETHEUS_MULTIPROC_DIR detected.") + prometheus_multiproc_dir = tempfile.TemporaryDirectory( + dir=os.environ["PROMETHEUS_MULTIPROC_DIR"] + ) + else: + prometheus_multiproc_dir = tempfile.TemporaryDirectory() + os.environ["PROMETHEUS_MULTIPROC_DIR"] = prometheus_multiproc_dir.name + logger.info(f"PROMETHEUS_MULTIPROC_DIR: {os.environ['PROMETHEUS_MULTIPROC_DIR']}") From f79d80c614073cebc16ebca66102c87b81c92d2f Mon Sep 17 00:00:00 2001 From: Ye Zhang Date: Wed, 9 Jul 2025 20:56:17 +0800 Subject: [PATCH 02/12] fix: pre-commit checks Signed-off-by: Ye Zhang rebase in progress; onto 82d03ca97 --- tensorrt_llm/executor/postproc_worker.py | 9 ++++-- tensorrt_llm/executor/proxy.py | 27 +++++++++++----- tensorrt_llm/executor/result.py | 39 +++++++++++++++--------- tensorrt_llm/executor/worker.py | 27 +++++++++++----- tensorrt_llm/metrics/__init__.py | 7 ++--- tensorrt_llm/metrics/collector.py | 30 ++++++++++++------ tensorrt_llm/utils/utils.py | 6 ++-- 7 files changed, 95 insertions(+), 50 deletions(-) diff --git a/tensorrt_llm/executor/postproc_worker.py b/tensorrt_llm/executor/postproc_worker.py index c358168e58c9..7dff32891859 100644 --- a/tensorrt_llm/executor/postproc_worker.py +++ b/tensorrt_llm/executor/postproc_worker.py @@ -120,8 +120,7 @@ def default_record_creator( tokenizer=tokenizer) async def _handle_input( - self, - input: Union["PostprocWorker.Input", "ResponseWrapper"] + self, input: Union["PostprocWorker.Input", "ResponseWrapper"] ) -> [Any, Optional[dict[str, float]]]: ''' Handle a single response from await_response worker. ''' if input.rsp.result.context_logits is not None or \ @@ -179,7 +178,11 @@ async def handle_single_input(inp: PostprocWorker.Input, is_final = inp.rsp.result.is_final if is_llm_response( inp.rsp) else True res, metrics = await self._handle_input(inp) - batch.append(PostprocWorker.Output(client_id=client_id, res=res, is_final=is_final, metrics=metrics)) + batch.append( + PostprocWorker.Output(client_id=client_id, + res=res, + is_final=is_final, + metrics=metrics)) if is_final: self._records.pop(client_id) diff --git a/tensorrt_llm/executor/proxy.py b/tensorrt_llm/executor/proxy.py index 703001f6a3cd..0f27fa0b9289 100644 --- a/tensorrt_llm/executor/proxy.py +++ b/tensorrt_llm/executor/proxy.py @@ -13,13 +13,13 @@ from tensorrt_llm.metrics.collector import MetricsCollector from .._utils import customized_gc_thresholds, mpi_rank, nvtx_range_debug -from ..utils.utils import set_prometheus_multiproc_dir from ..llmapi.mpi_session import (MpiCommSession, MpiPoolSession, MpiSession, RemoteMpiCommSessionClient) from ..llmapi.tracer import enable_llm_tracer, get_tracer, global_tracer from ..llmapi.utils import (AsyncQueue, ManagedThread, _SyncQueue, enable_llm_debug, print_colored, print_colored_debug) +from ..utils.utils import set_prometheus_multiproc_dir from .executor import GenerationExecutor from .ipc import FusedIpcQueue, IpcQueue from .postproc_worker import PostprocWorker, PostprocWorkerConfig @@ -62,7 +62,10 @@ def __init__( self.worker_cls = worker_cls set_prometheus_multiproc_dir() - self.metrics_collector = MetricsCollector({"model_name": "undefined", "engine_type": "undefined"}) + self.metrics_collector = MetricsCollector({ + "model_name": "undefined", + "engine_type": "undefined" + }) mpi_process_pre_spawned: bool = get_spawn_proxy_process_env() @@ -187,15 +190,23 @@ def process_res(res): if res._response and res._response.result and res._response.result.is_final: self._results[client_id]._handle_response(res) if metrics_dict := self._results[client_id].metrics_dict: - if finish_reason := metrics_dict.get(MetricsCollector.labelname_finish_reason): - self.metrics_collector.log_request_success(1, {MetricsCollector.labelname_finish_reason: - finish_reason}) + if finish_reason := metrics_dict.get( + MetricsCollector.labelname_finish_reason): + self.metrics_collector.log_request_success( + 1, { + MetricsCollector.labelname_finish_reason: + finish_reason + }) self.metrics_collector.log_histogram(metrics_dict) if isinstance(res, PostprocWorker.Output): if metrics_dict := res.metrics: - if finish_reason := metrics_dict.get(MetricsCollector.labelname_finish_reason): - self.metrics_collector.log_request_success(1, {MetricsCollector.labelname_finish_reason: - finish_reason}) + if finish_reason := metrics_dict.get( + MetricsCollector.labelname_finish_reason): + self.metrics_collector.log_request_success( + 1, { + MetricsCollector.labelname_finish_reason: + finish_reason + }) self.metrics_collector.log_histogram(metrics_dict) if (is_llm_response(res) and res.result.is_final) or isinstance( diff --git a/tensorrt_llm/executor/result.py b/tensorrt_llm/executor/result.py index 5efe89b6d3fd..d2620849ca29 100644 --- a/tensorrt_llm/executor/result.py +++ b/tensorrt_llm/executor/result.py @@ -15,7 +15,7 @@ from ..disaggregated_params import DisaggregatedParams from ..llmapi.tracer import global_tracer from ..llmapi.utils import AsyncQueue -from ..metrics import SupportedMetricNames, RequestEventTiming, MetricsCollector +from ..metrics import MetricsCollector, RequestEventTiming, SupportedMetricNames from ..sampling_params import LogprobParams, SamplingParams from .utils import ErrorResponse, has_event_loop, is_llm_response @@ -216,7 +216,8 @@ def _handle_sequence(self, response_tensors, sequence_index, logprobs_result=None, - req_perf_metrics_dict: Optional[dict[str, float]] = None): + req_perf_metrics_dict: Optional[dict[str, + float]] = None): """ Handle a single sequence in the response. """ seq_idx = sequence_index @@ -340,7 +341,8 @@ def _handle_response(self, if self.sampling_params.use_beam_search: for beam_idx, _ in enumerate(response_result.output_token_ids): self._handle_sequence(finish_reasons, response_result, - beam_idx, logprobs_result, req_perf_metrics_dict) + beam_idx, logprobs_result, + req_perf_metrics_dict) else: self._handle_sequence(finish_reasons, response_result, response_result.sequence_index, @@ -360,17 +362,24 @@ def _handle_response(self, else: raise ValueError(f"Unknown response type: {response}") - def record_stats(self, output: CompletionOutput, stats: Optional[dict[str, float]] = None): + def record_stats(self, + output: CompletionOutput, + stats: Optional[dict[str, float]] = None): if not stats: return metrics_stats = {} if output.finish_reason: - metrics_stats.update({MetricsCollector.labelname_finish_reason: output.finish_reason}) - processed_metrics_stat = _process_req_perf_metrics(stats, len(output.token_ids)) + metrics_stats.update({ + MetricsCollector.labelname_finish_reason: + output.finish_reason + }) + processed_metrics_stat = _process_req_perf_metrics( + stats, len(output.token_ids)) if processed_metrics_stat: metrics_stats.update(processed_metrics_stat) self.metrics_dict = metrics_stats + class DetokenizedGenerationResultBase(GenerationResultBase): ''' The base class for the generation result with detokenization support. ''' # import once and avoid cyclic import @@ -717,6 +726,7 @@ def _topk_logprobs(logits: torch.Tensor, top_k: int, return LogProbsResult(prompt=prompt_logprobs, generation=generation_logprobs) + def _process_req_perf_metrics(req_perf_metrics_dict: Optional[dict[str, float]], output_length: int) -> Optional[dict[str, float]]: stat = {} @@ -728,14 +738,15 @@ def _process_req_perf_metrics(req_perf_metrics_dict: Optional[dict[str, float]], req_perf_metrics_dict.get(RequestEventTiming.ARRIVAL_TIME, 0) request_queue_time = req_perf_metrics_dict.get(RequestEventTiming.FIRST_SCHEDULED_TIME, 0) - \ req_perf_metrics_dict.get(RequestEventTiming.ARRIVAL_TIME, 0) - stat = {SupportedMetricNames.TTFT: ttft, - SupportedMetricNames.E2E: e2e, - SupportedMetricNames.REQUEST_QUEUE_TIME: request_queue_time} + stat = { + SupportedMetricNames.TTFT: ttft, + SupportedMetricNames.E2E: e2e, + SupportedMetricNames.REQUEST_QUEUE_TIME: request_queue_time + } if output_length > 1: - tpot = (req_perf_metrics_dict.get(RequestEventTiming.LAST_TOKEN_TIME, 0) - - req_perf_metrics_dict.get(RequestEventTiming.FIRST_TOKEN_TIME, 0)) / (output_length - 1) + tpot = (req_perf_metrics_dict.get( + RequestEventTiming.LAST_TOKEN_TIME, 0) - req_perf_metrics_dict.get( + RequestEventTiming.FIRST_TOKEN_TIME, 0)) / (output_length - 1) stat.update({SupportedMetricNames.TPOT: tpot}) - stat = dict( - filter(lambda item: item[1] > 0, stat.items()) - ) + stat = dict(filter(lambda item: item[1] > 0, stat.items())) return stat diff --git a/tensorrt_llm/executor/worker.py b/tensorrt_llm/executor/worker.py index 75874e0a5591..74affffad3b4 100644 --- a/tensorrt_llm/executor/worker.py +++ b/tensorrt_llm/executor/worker.py @@ -473,7 +473,8 @@ def _deduce_max_tokens(request: GenerationRequest, request.sampling_params.end_id, pad_id=request.sampling_params.pad_id, output_config=request.sampling_params._get_output_config( - is_pytorch_backend=self._is_pytorch_backend, return_perf_metrics=request.return_perf_metrics), + is_pytorch_backend=self._is_pytorch_backend, + return_perf_metrics=request.return_perf_metrics), # Beam search enforces return_all_generated_tokens=True regardless of the passed value return_all_generated_tokens=False, # convert python config into pybind config @@ -904,7 +905,8 @@ def handle_for_worker(self, responses: List[tllm.Response]) -> None: self.worker._is_pytorch_backend) req_perf_metrics = _get_req_perf_metrics(response) if logprobs_result or req_perf_metrics: - response = ResponseWrapper(response, logprobs_result, req_perf_metrics) + response = ResponseWrapper(response, logprobs_result, + req_perf_metrics) # For AsyncQueue.sync_q, we will batch the events to avoid too many # event notifications, thus put without wait here. @@ -946,7 +948,8 @@ def handle_for_ipc_batched(self, responses: List[tllm.Response]) -> None: self.worker._is_pytorch_backend) req_perf_metrics = _get_req_perf_metrics(response) if logprobs_result or req_perf_metrics: - response = ResponseWrapper(response, logprobs_result, req_perf_metrics) + response = ResponseWrapper(response, logprobs_result, + req_perf_metrics) _send_rsp(self.worker, response, @@ -1055,15 +1058,23 @@ def _send_rsp( else: raise ValueError(f"Unknown response type: {response}") -def _get_req_perf_metrics(response: tllm.Response) -> Optional[dict[str, float]]: + +def _get_req_perf_metrics( + response: tllm.Response) -> Optional[dict[str, float]]: req_perf_metrics, metrics_dict = None, {} if response.result: req_perf_metrics = response.result.request_perf_metrics if req_perf_metrics and req_perf_metrics.timing_metrics: metrics_dict = { - RequestEventTiming.ARRIVAL_TIME: req_perf_metrics.timing_metrics.arrival_time.total_seconds(), - RequestEventTiming.FIRST_TOKEN_TIME: req_perf_metrics.timing_metrics.first_token_time.total_seconds(), + RequestEventTiming.ARRIVAL_TIME: + req_perf_metrics.timing_metrics.arrival_time.total_seconds(), + RequestEventTiming.FIRST_TOKEN_TIME: + req_perf_metrics.timing_metrics.first_token_time.total_seconds( + ), RequestEventTiming.FIRST_SCHEDULED_TIME: - req_perf_metrics.timing_metrics.first_scheduled_time.total_seconds(), - RequestEventTiming.LAST_TOKEN_TIME: req_perf_metrics.timing_metrics.last_token_time.total_seconds()} + req_perf_metrics.timing_metrics.first_scheduled_time. + total_seconds(), + RequestEventTiming.LAST_TOKEN_TIME: + req_perf_metrics.timing_metrics.last_token_time.total_seconds() + } return metrics_dict diff --git a/tensorrt_llm/metrics/__init__.py b/tensorrt_llm/metrics/__init__.py index 4650a84709d5..339982c2d61d 100644 --- a/tensorrt_llm/metrics/__init__.py +++ b/tensorrt_llm/metrics/__init__.py @@ -1,7 +1,4 @@ from .collector import * from .enums import * -__all__ = [ - "SupportedMetricNames", - "MetricsCollector", - "RequestEventTiming" -] \ No newline at end of file + +__all__ = ["SupportedMetricNames", "MetricsCollector", "RequestEventTiming"] diff --git a/tensorrt_llm/metrics/collector.py b/tensorrt_llm/metrics/collector.py index 23616b83ae2b..abbc7a50a7d1 100644 --- a/tensorrt_llm/metrics/collector.py +++ b/tensorrt_llm/metrics/collector.py @@ -8,13 +8,19 @@ class MetricsCollector: labelname_finish_reason = "finished_reason" + def __init__(self, labels: Dict[str, str]) -> None: from prometheus_client import Counter, Histogram self.last_log_time = time.time() self.labels = labels - self.finish_reason_label = {MetricsCollector.labelname_finish_reason: "unknown"} - self.labels_with_finished_reason = {**self.labels, **self.finish_reason_label} + self.finish_reason_label = { + MetricsCollector.labelname_finish_reason: "unknown" + } + self.labels_with_finished_reason = { + **self.labels, + **self.finish_reason_label + } self.counter_request_success = Counter( name="request_success_total", @@ -51,10 +57,12 @@ def __init__(self, labels: Dict[str, str]) -> None: self.histogram_queue_time_request = Histogram( name="request_queue_time_seconds", - documentation="Histogram of time spent in WAITING phase for request.", - buckets = [ + documentation= + "Histogram of time spent in WAITING phase for request.", + buckets=[ 0.3, 0.5, 0.8, 1.0, 1.5, 2.0, 2.5, 5.0, 10.0, 15.0, 20.0, 30.0, - 40.0, 50.0, 60.0, 120.0, 240.0, 480.0, 960.0, 1920.0, 7680.0], + 40.0, 50.0, 60.0, 120.0, 240.0, 480.0, 960.0, 1920.0, 7680.0 + ], labelnames=self.labels.keys()) def _label_merge(self, labels: Dict[str, str]) -> Dict[str, str]: @@ -62,7 +70,8 @@ def _label_merge(self, labels: Dict[str, str]) -> Dict[str, str]: return self.labels return {**self.labels, **labels} - def _log_counter(self, counter, labels: Dict[str, str], data: Union[int, float]) -> None: + def _log_counter(self, counter, labels: Dict[str, str], + data: Union[int, float]) -> None: # Convenience function for logging to counter. counter.labels(**self._label_merge(labels)).inc(data) @@ -70,7 +79,8 @@ def _log_histogram(self, histogram, data: Union[int, float]) -> None: # Convenience function for logging to histogram. histogram.labels(**self.labels).observe(data) - def log_request_success(self, data: Union[int, float], labels: Dict[str, str]) -> None: + def log_request_success(self, data: Union[int, float], + labels: Dict[str, str]) -> None: self._log_counter(self.counter_request_success, labels, data) self.last_log_time = time.time() @@ -81,6 +91,8 @@ def log_histogram(self, data: Optional[dict[str, float]]) -> None: self._log_histogram(self.histogram_time_to_first_token, ttft) if tpot := data.get(SupportedMetricNames.TPOT, 0): self._log_histogram(self.histogram_time_per_output_token, tpot) - if request_queue_time := data.get(SupportedMetricNames.REQUEST_QUEUE_TIME, 0): - self._log_histogram(self.histogram_queue_time_request, request_queue_time) + if request_queue_time := data.get( + SupportedMetricNames.REQUEST_QUEUE_TIME, 0): + self._log_histogram(self.histogram_queue_time_request, + request_queue_time) self.last_log_time = time.time() diff --git a/tensorrt_llm/utils/utils.py b/tensorrt_llm/utils/utils.py index c52718603bc4..6d946014ad0a 100644 --- a/tensorrt_llm/utils/utils.py +++ b/tensorrt_llm/utils/utils.py @@ -14,9 +14,9 @@ def set_prometheus_multiproc_dir(): if "PROMETHEUS_MULTIPROC_DIR" in os.environ: logger.info("User set PROMETHEUS_MULTIPROC_DIR detected.") prometheus_multiproc_dir = tempfile.TemporaryDirectory( - dir=os.environ["PROMETHEUS_MULTIPROC_DIR"] - ) + dir=os.environ["PROMETHEUS_MULTIPROC_DIR"]) else: prometheus_multiproc_dir = tempfile.TemporaryDirectory() os.environ["PROMETHEUS_MULTIPROC_DIR"] = prometheus_multiproc_dir.name - logger.info(f"PROMETHEUS_MULTIPROC_DIR: {os.environ['PROMETHEUS_MULTIPROC_DIR']}") + logger.info( + f"PROMETHEUS_MULTIPROC_DIR: {os.environ['PROMETHEUS_MULTIPROC_DIR']}") From 39f1822b08e1cd07fdbb74f4de9840dcf1bb8a28 Mon Sep 17 00:00:00 2001 From: Ye Zhang Date: Thu, 17 Jul 2025 14:33:17 +0800 Subject: [PATCH 03/12] fix: consider multi-response when calculating tpot Signed-off-by: Ye Zhang rebase in progress; onto 82d03ca97 --- tensorrt_llm/_utils.py | 18 ++++++++++++++++++ tensorrt_llm/executor/proxy.py | 22 ++++------------------ tensorrt_llm/executor/result.py | 11 +++++++---- tensorrt_llm/llmapi/llm_args.py | 4 ++-- tensorrt_llm/metrics/collector.py | 7 +++++++ tensorrt_llm/sampling_params.py | 4 +--- tensorrt_llm/utils/__init__.py | 0 tensorrt_llm/utils/utils.py | 22 ---------------------- 8 files changed, 39 insertions(+), 49 deletions(-) delete mode 100644 tensorrt_llm/utils/__init__.py delete mode 100644 tensorrt_llm/utils/utils.py diff --git a/tensorrt_llm/_utils.py b/tensorrt_llm/_utils.py index 337caae51266..7c6179148435 100644 --- a/tensorrt_llm/_utils.py +++ b/tensorrt_llm/_utils.py @@ -20,6 +20,7 @@ import math import os import struct +import tempfile import trace import weakref from contextlib import contextmanager @@ -1112,3 +1113,20 @@ def is_multi_device_enable(): the number of devices """ return local_mpi_size() > 1 + + +def set_prometheus_multiproc_dir() -> object: + # Set prometheus multiprocess directory + # sglang uses prometheus multiprocess mode + # we need to set this before importing prometheus_client + # https://prometheus.github.io/client_python/multiprocess/ + global prometheus_multiproc_dir + if "PROMETHEUS_MULTIPROC_DIR" in os.environ: + logger.info("User set PROMETHEUS_MULTIPROC_DIR detected.") + prometheus_multiproc_dir = tempfile.TemporaryDirectory( + dir=os.environ["PROMETHEUS_MULTIPROC_DIR"]) + else: + prometheus_multiproc_dir = tempfile.TemporaryDirectory() + os.environ["PROMETHEUS_MULTIPROC_DIR"] = prometheus_multiproc_dir.name + logger.info( + f"PROMETHEUS_MULTIPROC_DIR: {os.environ['PROMETHEUS_MULTIPROC_DIR']}") diff --git a/tensorrt_llm/executor/proxy.py b/tensorrt_llm/executor/proxy.py index 0f27fa0b9289..9eeb7a4350dd 100644 --- a/tensorrt_llm/executor/proxy.py +++ b/tensorrt_llm/executor/proxy.py @@ -12,14 +12,14 @@ from tensorrt_llm.logger import logger from tensorrt_llm.metrics.collector import MetricsCollector -from .._utils import customized_gc_thresholds, mpi_rank, nvtx_range_debug +from .._utils import (customized_gc_thresholds, mpi_rank, nvtx_range_debug, + set_prometheus_multiproc_dir) from ..llmapi.mpi_session import (MpiCommSession, MpiPoolSession, MpiSession, RemoteMpiCommSessionClient) from ..llmapi.tracer import enable_llm_tracer, get_tracer, global_tracer from ..llmapi.utils import (AsyncQueue, ManagedThread, _SyncQueue, enable_llm_debug, print_colored, print_colored_debug) -from ..utils.utils import set_prometheus_multiproc_dir from .executor import GenerationExecutor from .ipc import FusedIpcQueue, IpcQueue from .postproc_worker import PostprocWorker, PostprocWorkerConfig @@ -190,24 +190,10 @@ def process_res(res): if res._response and res._response.result and res._response.result.is_final: self._results[client_id]._handle_response(res) if metrics_dict := self._results[client_id].metrics_dict: - if finish_reason := metrics_dict.get( - MetricsCollector.labelname_finish_reason): - self.metrics_collector.log_request_success( - 1, { - MetricsCollector.labelname_finish_reason: - finish_reason - }) - self.metrics_collector.log_histogram(metrics_dict) + self.metrics_collector.log_metrics_dict(metrics_dict) if isinstance(res, PostprocWorker.Output): if metrics_dict := res.metrics: - if finish_reason := metrics_dict.get( - MetricsCollector.labelname_finish_reason): - self.metrics_collector.log_request_success( - 1, { - MetricsCollector.labelname_finish_reason: - finish_reason - }) - self.metrics_collector.log_histogram(metrics_dict) + self.metrics_collector.log_metrics_dict(metrics_dict) if (is_llm_response(res) and res.result.is_final) or isinstance( res, ErrorResponse): diff --git a/tensorrt_llm/executor/result.py b/tensorrt_llm/executor/result.py index d2620849ca29..2a5efff3d32a 100644 --- a/tensorrt_llm/executor/result.py +++ b/tensorrt_llm/executor/result.py @@ -374,7 +374,7 @@ def record_stats(self, output.finish_reason }) processed_metrics_stat = _process_req_perf_metrics( - stats, len(output.token_ids)) + stats, len(output.token_ids), self.sampling_params.n > 1) if processed_metrics_stat: metrics_stats.update(processed_metrics_stat) self.metrics_dict = metrics_stats @@ -727,8 +727,11 @@ def _topk_logprobs(logits: torch.Tensor, top_k: int, generation=generation_logprobs) -def _process_req_perf_metrics(req_perf_metrics_dict: Optional[dict[str, float]], - output_length: int) -> Optional[dict[str, float]]: +def _process_req_perf_metrics( + req_perf_metrics_dict: Optional[dict[str, float]], + output_length: int, + is_multiple_response: Optional[bool] = False +) -> Optional[dict[str, float]]: stat = {} if not req_perf_metrics_dict: return stat @@ -743,7 +746,7 @@ def _process_req_perf_metrics(req_perf_metrics_dict: Optional[dict[str, float]], SupportedMetricNames.E2E: e2e, SupportedMetricNames.REQUEST_QUEUE_TIME: request_queue_time } - if output_length > 1: + if output_length > 1 and not is_multiple_response: tpot = (req_perf_metrics_dict.get( RequestEventTiming.LAST_TOKEN_TIME, 0) - req_perf_metrics_dict.get( RequestEventTiming.FIRST_TOKEN_TIME, 0)) / (output_length - 1) diff --git a/tensorrt_llm/llmapi/llm_args.py b/tensorrt_llm/llmapi/llm_args.py index 258edbabeff0..3e42b7595996 100644 --- a/tensorrt_llm/llmapi/llm_args.py +++ b/tensorrt_llm/llmapi/llm_args.py @@ -1311,9 +1311,9 @@ class BaseLlmArgs(StrictBaseModel): status="deprecated", ) - return_perf_metrics: Optional[bool] = Field( + return_perf_metrics: bool = Field( default=False, - description="optionally returns perf metrics", + description="Return perf metrics", ) _parallel_config: Optional[object] = PrivateAttr(default=None) diff --git a/tensorrt_llm/metrics/collector.py b/tensorrt_llm/metrics/collector.py index abbc7a50a7d1..4a2ad5a7f470 100644 --- a/tensorrt_llm/metrics/collector.py +++ b/tensorrt_llm/metrics/collector.py @@ -96,3 +96,10 @@ def log_histogram(self, data: Optional[dict[str, float]]) -> None: self._log_histogram(self.histogram_queue_time_request, request_queue_time) self.last_log_time = time.time() + + def log_metrics_dict(self, metrics_dict: dict[str, float]) -> None: + if finish_reason := metrics_dict.get( + MetricsCollector.labelname_finish_reason): + self.log_request_success( + 1, {MetricsCollector.labelname_finish_reason: finish_reason}) + self.log_histogram(metrics_dict) diff --git a/tensorrt_llm/sampling_params.py b/tensorrt_llm/sampling_params.py index f37d357bbee4..94ea4dddad97 100644 --- a/tensorrt_llm/sampling_params.py +++ b/tensorrt_llm/sampling_params.py @@ -438,9 +438,7 @@ def _get_sampling_config(self) -> tllme.SamplingConfig: return tllme.SamplingConfig(**llmapi_to_rt_param_map) def _get_output_config( - self, - is_pytorch_backend: bool = False, - return_perf_metrics: Optional[bool] = False + self, is_pytorch_backend: bool = False, return_perf_metrics: Optional[bool] = False ) -> tllme.OutputConfig: sampling_param_fields = set(dir(SamplingParams)) fields = [ diff --git a/tensorrt_llm/utils/__init__.py b/tensorrt_llm/utils/__init__.py deleted file mode 100644 index e69de29bb2d1..000000000000 diff --git a/tensorrt_llm/utils/utils.py b/tensorrt_llm/utils/utils.py deleted file mode 100644 index 6d946014ad0a..000000000000 --- a/tensorrt_llm/utils/utils.py +++ /dev/null @@ -1,22 +0,0 @@ -import os -import tempfile - -from tensorrt_llm.logger import logger - - -def set_prometheus_multiproc_dir(): - # Set prometheus multiprocess directory - # sglang uses prometheus multiprocess mode - # we need to set this before importing prometheus_client - # https://prometheus.github.io/client_python/multiprocess/ - global prometheus_multiproc_dir - - if "PROMETHEUS_MULTIPROC_DIR" in os.environ: - logger.info("User set PROMETHEUS_MULTIPROC_DIR detected.") - prometheus_multiproc_dir = tempfile.TemporaryDirectory( - dir=os.environ["PROMETHEUS_MULTIPROC_DIR"]) - else: - prometheus_multiproc_dir = tempfile.TemporaryDirectory() - os.environ["PROMETHEUS_MULTIPROC_DIR"] = prometheus_multiproc_dir.name - logger.info( - f"PROMETHEUS_MULTIPROC_DIR: {os.environ['PROMETHEUS_MULTIPROC_DIR']}") From 2ae6a062cbca767cd34cce0cf1828ce1e71b7f52 Mon Sep 17 00:00:00 2001 From: Ye Zhang Date: Tue, 22 Jul 2025 20:58:20 +0800 Subject: [PATCH 04/12] style: pre-commit format Signed-off-by: Ye Zhang rebase in progress; onto 82d03ca97 --- requirements.txt | 2 + tensorrt_llm/_torch/pyexecutor/llm_request.py | 6 ++ tensorrt_llm/executor/proxy.py | 1 + tensorrt_llm/executor/result.py | 6 +- tensorrt_llm/executor/worker.py | 42 +++++++---- tensorrt_llm/metrics/collector.py | 1 + .../llmapi/apps/_test_openai_prometheus.py | 71 +++++++++++++++++++ 7 files changed, 112 insertions(+), 17 deletions(-) create mode 100644 tests/unittest/llmapi/apps/_test_openai_prometheus.py diff --git a/requirements.txt b/requirements.txt index f6b201b1f574..e2582f503856 100644 --- a/requirements.txt +++ b/requirements.txt @@ -29,6 +29,8 @@ nvidia-modelopt[torch]~=0.33.0 nvidia-nccl-cu12 nvidia-cuda-nvrtc-cu12 transformers==4.55.0 +prometheus_client +prometheus_fastapi_instrumentator pydantic>=2.9.1 pydantic-settings[yaml] omegaconf diff --git a/tensorrt_llm/_torch/pyexecutor/llm_request.py b/tensorrt_llm/_torch/pyexecutor/llm_request.py index a068327b6db8..db360f64a839 100644 --- a/tensorrt_llm/_torch/pyexecutor/llm_request.py +++ b/tensorrt_llm/_torch/pyexecutor/llm_request.py @@ -250,6 +250,12 @@ def deserialize(self): self._result = tensorrt_llm.bindings.executor.deserialize_result( self._result) + def get_result(self): + if tmp_res := tensorrt_llm.bindings.executor.deserialize_result( + self._result): + return tmp_res + return None + @dataclass class LlmResponse: diff --git a/tensorrt_llm/executor/proxy.py b/tensorrt_llm/executor/proxy.py index 9eeb7a4350dd..0b46ddddf1c3 100644 --- a/tensorrt_llm/executor/proxy.py +++ b/tensorrt_llm/executor/proxy.py @@ -23,6 +23,7 @@ from .executor import GenerationExecutor from .ipc import FusedIpcQueue, IpcQueue from .postproc_worker import PostprocWorker, PostprocWorkerConfig +from .request import CancellingRequest, GenerationRequest from .result import GenerationResult, IterationResult, ResponseWrapper from .utils import (ErrorResponse, IntraProcessQueue, WorkerCommIpcAddrs, create_mpi_comm_session, get_spawn_proxy_process_env, diff --git a/tensorrt_llm/executor/result.py b/tensorrt_llm/executor/result.py index 2a5efff3d32a..3388734957f6 100644 --- a/tensorrt_llm/executor/result.py +++ b/tensorrt_llm/executor/result.py @@ -53,7 +53,7 @@ class LogProbsResult(NamedTuple): class ResponseWrapper: """ 1. Wrapper of runtime response with optional outputs computed post runtime. - 2. a workaround to pass around RequestPerfMetrics. + 2. A workaround to pass around RequestPerfMetrics. """ def __init__(self, @@ -730,8 +730,8 @@ def _topk_logprobs(logits: torch.Tensor, top_k: int, def _process_req_perf_metrics( req_perf_metrics_dict: Optional[dict[str, float]], output_length: int, - is_multiple_response: Optional[bool] = False -) -> Optional[dict[str, float]]: + is_multiple_response: bool = False +) -> dict[SupportedMetricNames, float]: stat = {} if not req_perf_metrics_dict: return stat diff --git a/tensorrt_llm/executor/worker.py b/tensorrt_llm/executor/worker.py index 74affffad3b4..0e35551f3c08 100644 --- a/tensorrt_llm/executor/worker.py +++ b/tensorrt_llm/executor/worker.py @@ -901,12 +901,8 @@ def handle_for_worker(self, responses: List[tllm.Response]) -> None: assert response is not None queue = self.worker.return_queue(response.client_id) - logprobs_result = _get_logprobs(self.worker, response, + response = _maybe_wrap_response(self.worker, response, self.worker._is_pytorch_backend) - req_perf_metrics = _get_req_perf_metrics(response) - if logprobs_result or req_perf_metrics: - response = ResponseWrapper(response, logprobs_result, - req_perf_metrics) # For AsyncQueue.sync_q, we will batch the events to avoid too many # event notifications, thus put without wait here. @@ -944,12 +940,8 @@ def handle_for_ipc_batched(self, responses: List[tllm.Response]) -> None: response = ErrorResponse(response.client_id, response.error_msg, response.request_id) else: - logprobs_result = _get_logprobs(self.worker, response, + response = _maybe_wrap_response(self.worker, response, self.worker._is_pytorch_backend) - req_perf_metrics = _get_req_perf_metrics(response) - if logprobs_result or req_perf_metrics: - response = ResponseWrapper(response, logprobs_result, - req_perf_metrics) _send_rsp(self.worker, response, @@ -1059,11 +1051,16 @@ def _send_rsp( raise ValueError(f"Unknown response type: {response}") -def _get_req_perf_metrics( - response: tllm.Response) -> Optional[dict[str, float]]: +def _get_metrics_dict( + response: tllm.Response) -> dict[RequestEventTiming, float]: req_perf_metrics, metrics_dict = None, {} - if response.result: - req_perf_metrics = response.result.request_perf_metrics + res = response.result + if res: + if hasattr(res, '_result'): + if result := res.get_result(): + req_perf_metrics = result.request_perf_metrics + else: + req_perf_metrics = res.request_perf_metrics if req_perf_metrics and req_perf_metrics.timing_metrics: metrics_dict = { RequestEventTiming.ARRIVAL_TIME: @@ -1078,3 +1075,20 @@ def _get_req_perf_metrics( req_perf_metrics.timing_metrics.last_token_time.total_seconds() } return metrics_dict + + +def _maybe_wrap_response( + worker, + response: tllm.Response, + is_pytorch_backend=False) -> Union[tllm.Response, ResponseWrapper]: + + logprobs_result = _get_logprobs(worker, response, is_pytorch_backend) + req_perf_metrics = _get_metrics_dict(response) + if logprobs_result or req_perf_metrics: + if hasattr(response, 'result') and hasattr(response.result, "_result"): + response.result.deserialize() + response = tllm.Response(request_id=response.request_id, + result=response.result._result, + client_id=response.client_id) + response = ResponseWrapper(response, logprobs_result, req_perf_metrics) + return response diff --git a/tensorrt_llm/metrics/collector.py b/tensorrt_llm/metrics/collector.py index 4a2ad5a7f470..7d6a4558c69d 100644 --- a/tensorrt_llm/metrics/collector.py +++ b/tensorrt_llm/metrics/collector.py @@ -6,6 +6,7 @@ from .enums import SupportedMetricNames +# Adapted from https://github.com/vllm-project/vllm/blob/v0.10.0rc1/vllm/engine/metrics.py#L30 class MetricsCollector: labelname_finish_reason = "finished_reason" diff --git a/tests/unittest/llmapi/apps/_test_openai_prometheus.py b/tests/unittest/llmapi/apps/_test_openai_prometheus.py new file mode 100644 index 000000000000..d40d2d5a1288 --- /dev/null +++ b/tests/unittest/llmapi/apps/_test_openai_prometheus.py @@ -0,0 +1,71 @@ +import logging +import os +import tempfile +from urllib.request import urlopen + +import pytest +import torch +import yaml + +from ..test_llm import get_model_path +from .openai_server import RemoteOpenAIServer + +# Configure logging +logging.basicConfig(level=logging.INFO) +logger = logging.getLogger(__name__) + + +@pytest.fixture(scope="module", ids=["TinyLlama-1.1B-Chat"]) +def model_name(): + return "llama-models-v2/TinyLlama-1.1B-Chat-v1.0" + + +@pytest.fixture(scope="module") +def temp_extra_llm_api_options_file(request): + temp_dir = tempfile.gettempdir() + temp_file_path = os.path.join(temp_dir, "extra_llm_api_options.yaml") + try: + extra_llm_api_options_dict = {"return_perf_metrics": True} + + with open(temp_file_path, 'w') as f: + yaml.dump(extra_llm_api_options_dict, f) + + yield temp_file_path + finally: + if os.path.exists(temp_file_path): + os.remove(temp_file_path) + + +@pytest.fixture(scope="module") +def server(model_name: str, + temp_extra_llm_api_options_file: str) -> RemoteOpenAIServer: + model_path = get_model_path(model_name) + # Set parallelism based on GPU count + tp_size = torch.cuda.device_count() + assert tp_size > 0, "At least 1 GPU is required to run tests" + args = ["--backend", "pytorch", "--tp_size", str(tp_size)] + args.extend(["--extra_llm_api_options", temp_extra_llm_api_options_file]) + logger.info(f"Starting server, model: {model_name}, args: {args}") + with RemoteOpenAIServer(model_path, args) as remote_server: + yield remote_server + logger.info("Tests completed, shutting down server") + + +def test_metrics_endpoint(server: RemoteOpenAIServer): + + client = server.get_client() + client.completions.create( + model="Server", + prompt="Hello, my name is", + max_tokens=25, + stream=False, + ) + + response = urlopen(f'{server.url_root}/prometheus/metrics') + assert response.status is 200 + + data = response.read().decode("utf-8") + assert "request_success_total" in data + assert "e2e_request_latency_seconds" in data + assert "time_to_first_token_seconds" in data + assert "request_queue_time_seconds" in data From abe01680ebd24820497164ba5f2ba243829b3270 Mon Sep 17 00:00:00 2001 From: Ye Zhang Date: Tue, 29 Jul 2025 10:48:25 +0800 Subject: [PATCH 05/12] fix: pre-commit B105 false positive Signed-off-by: Ye Zhang --- tensorrt_llm/metrics/enums.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tensorrt_llm/metrics/enums.py b/tensorrt_llm/metrics/enums.py index 1bab966106cc..442f7bbbc54f 100644 --- a/tensorrt_llm/metrics/enums.py +++ b/tensorrt_llm/metrics/enums.py @@ -10,6 +10,6 @@ class SupportedMetricNames(Enum): class RequestEventTiming(Enum): ARRIVAL_TIME = "arrival_time" - FIRST_TOKEN_TIME = "first_token_time" + FIRST_TOKEN_TIME = "first_token_time" # nosec: B105 FIRST_SCHEDULED_TIME = "first_scheduled_time" - LAST_TOKEN_TIME = "last_token_time" + LAST_TOKEN_TIME = "last_token_time" # nosec: B105 From 6558f5e092c44c63d34b34bbc3aad7c7f4789dae Mon Sep 17 00:00:00 2001 From: Ye Zhang Date: Tue, 29 Jul 2025 16:23:47 +0800 Subject: [PATCH 06/12] test: add prometheus to integration test Signed-off-by: Ye Zhang --- tests/integration/defs/test_e2e.py | 7 +++++++ tests/integration/test_lists/test-db/l0_a10.yml | 1 + tests/unittest/llmapi/apps/_test_openai_prometheus.py | 6 +----- 3 files changed, 9 insertions(+), 5 deletions(-) diff --git a/tests/integration/defs/test_e2e.py b/tests/integration/defs/test_e2e.py index 33e3f55765ed..7c6a203d0b4e 100644 --- a/tests/integration/defs/test_e2e.py +++ b/tests/integration/defs/test_e2e.py @@ -1497,6 +1497,13 @@ def test_openai_chat_with_logit_bias(llm_root, llm_venv, sampler: str): ]) +def test_openai_prometheus(llm_root, llm_venv): + test_root = unittest_path() / "llmapi" / "apps" + llm_venv.run_cmd( + ["-m", "pytest", + str(test_root / "_test_openai_prometheus.py")]) + + def test_openai_lora(llm_root, llm_venv): test_root = unittest_path() / "llmapi" / "apps" llm_venv.run_cmd(["-m", "pytest", str(test_root / "_test_openai_lora.py")]) diff --git a/tests/integration/test_lists/test-db/l0_a10.yml b/tests/integration/test_lists/test-db/l0_a10.yml index 891649e5b9f2..ce285faa7994 100644 --- a/tests/integration/test_lists/test-db/l0_a10.yml +++ b/tests/integration/test_lists/test-db/l0_a10.yml @@ -25,6 +25,7 @@ l0_a10: - test_e2e.py::test_openai_chat_structural_tag_example - test_e2e.py::test_openai_chat_json_example - test_e2e.py::test_openai_chat_multimodal_example + - test_e2e.py::test_openai_prometheus - test_e2e.py::test_openai_lora - test_e2e.py::test_trtllm_serve_multimodal_example - test_e2e.py::test_trtllm_serve_lora_example diff --git a/tests/unittest/llmapi/apps/_test_openai_prometheus.py b/tests/unittest/llmapi/apps/_test_openai_prometheus.py index d40d2d5a1288..8a360668fd5d 100644 --- a/tests/unittest/llmapi/apps/_test_openai_prometheus.py +++ b/tests/unittest/llmapi/apps/_test_openai_prometheus.py @@ -4,7 +4,6 @@ from urllib.request import urlopen import pytest -import torch import yaml from ..test_llm import get_model_path @@ -40,10 +39,7 @@ def temp_extra_llm_api_options_file(request): def server(model_name: str, temp_extra_llm_api_options_file: str) -> RemoteOpenAIServer: model_path = get_model_path(model_name) - # Set parallelism based on GPU count - tp_size = torch.cuda.device_count() - assert tp_size > 0, "At least 1 GPU is required to run tests" - args = ["--backend", "pytorch", "--tp_size", str(tp_size)] + args = ["--backend", "pytorch", "--tp_size", "1"] args.extend(["--extra_llm_api_options", temp_extra_llm_api_options_file]) logger.info(f"Starting server, model: {model_name}, args: {args}") with RemoteOpenAIServer(model_path, args) as remote_server: From 267ed726bce07f9064a158606865a132625b8c01 Mon Sep 17 00:00:00 2001 From: Ye Zhang Date: Tue, 29 Jul 2025 21:54:42 +0800 Subject: [PATCH 07/12] fix: optimize pytorch response handling in _handle_response Signed-off-by: Ye Zhang --- tensorrt_llm/executor/result.py | 3 ++- tensorrt_llm/executor/worker.py | 5 ----- 2 files changed, 2 insertions(+), 6 deletions(-) diff --git a/tensorrt_llm/executor/result.py b/tensorrt_llm/executor/result.py index 3388734957f6..4bf60d86274a 100644 --- a/tensorrt_llm/executor/result.py +++ b/tensorrt_llm/executor/result.py @@ -322,7 +322,8 @@ def _handle_response(self, handler(response.error_msg) response_result = response.result - if hasattr(response_result, "_result"): + if hasattr(response_result, "_result") and isinstance( + response_result._result, bytes): response_result.deserialize() self._done = response_result.is_final diff --git a/tensorrt_llm/executor/worker.py b/tensorrt_llm/executor/worker.py index 0e35551f3c08..55e2ce41272a 100644 --- a/tensorrt_llm/executor/worker.py +++ b/tensorrt_llm/executor/worker.py @@ -1085,10 +1085,5 @@ def _maybe_wrap_response( logprobs_result = _get_logprobs(worker, response, is_pytorch_backend) req_perf_metrics = _get_metrics_dict(response) if logprobs_result or req_perf_metrics: - if hasattr(response, 'result') and hasattr(response.result, "_result"): - response.result.deserialize() - response = tllm.Response(request_id=response.request_id, - result=response.result._result, - client_id=response.client_id) response = ResponseWrapper(response, logprobs_result, req_perf_metrics) return response From 9e9bd0a31bdc37e2a774e4ce543db24a87124c72 Mon Sep 17 00:00:00 2001 From: Ye Zhang Date: Mon, 4 Aug 2025 11:13:00 +0800 Subject: [PATCH 08/12] fix: move return_perf_metrics arg into sampling_params to keep executor api unchanged Signed-off-by: Ye Zhang --- tensorrt_llm/executor/result.py | 13 ++++++------- tensorrt_llm/executor/worker.py | 3 +-- tensorrt_llm/llmapi/llm.py | 2 +- tensorrt_llm/metrics/__init__.py | 2 +- tensorrt_llm/metrics/collector.py | 11 +++++------ tensorrt_llm/metrics/enums.py | 6 +++--- tensorrt_llm/sampling_params.py | 5 +---- tests/unittest/api_stability/references/llm.yaml | 4 ++++ 8 files changed, 22 insertions(+), 24 deletions(-) diff --git a/tensorrt_llm/executor/result.py b/tensorrt_llm/executor/result.py index 4bf60d86274a..1f7adf77e1c8 100644 --- a/tensorrt_llm/executor/result.py +++ b/tensorrt_llm/executor/result.py @@ -15,7 +15,7 @@ from ..disaggregated_params import DisaggregatedParams from ..llmapi.tracer import global_tracer from ..llmapi.utils import AsyncQueue -from ..metrics import MetricsCollector, RequestEventTiming, SupportedMetricNames +from ..metrics import MetricNames, MetricsCollector, RequestEventTiming from ..sampling_params import LogprobParams, SamplingParams from .utils import ErrorResponse, has_event_loop, is_llm_response @@ -731,8 +731,7 @@ def _topk_logprobs(logits: torch.Tensor, top_k: int, def _process_req_perf_metrics( req_perf_metrics_dict: Optional[dict[str, float]], output_length: int, - is_multiple_response: bool = False -) -> dict[SupportedMetricNames, float]: + is_multiple_response: bool = False) -> dict[MetricNames, float]: stat = {} if not req_perf_metrics_dict: return stat @@ -743,14 +742,14 @@ def _process_req_perf_metrics( request_queue_time = req_perf_metrics_dict.get(RequestEventTiming.FIRST_SCHEDULED_TIME, 0) - \ req_perf_metrics_dict.get(RequestEventTiming.ARRIVAL_TIME, 0) stat = { - SupportedMetricNames.TTFT: ttft, - SupportedMetricNames.E2E: e2e, - SupportedMetricNames.REQUEST_QUEUE_TIME: request_queue_time + MetricNames.TTFT: ttft, + MetricNames.E2E: e2e, + MetricNames.REQUEST_QUEUE_TIME: request_queue_time } if output_length > 1 and not is_multiple_response: tpot = (req_perf_metrics_dict.get( RequestEventTiming.LAST_TOKEN_TIME, 0) - req_perf_metrics_dict.get( RequestEventTiming.FIRST_TOKEN_TIME, 0)) / (output_length - 1) - stat.update({SupportedMetricNames.TPOT: tpot}) + stat.update({MetricNames.TPOT: tpot}) stat = dict(filter(lambda item: item[1] > 0, stat.items())) return stat diff --git a/tensorrt_llm/executor/worker.py b/tensorrt_llm/executor/worker.py index 55e2ce41272a..c3a827bb00b1 100644 --- a/tensorrt_llm/executor/worker.py +++ b/tensorrt_llm/executor/worker.py @@ -473,8 +473,7 @@ def _deduce_max_tokens(request: GenerationRequest, request.sampling_params.end_id, pad_id=request.sampling_params.pad_id, output_config=request.sampling_params._get_output_config( - is_pytorch_backend=self._is_pytorch_backend, - return_perf_metrics=request.return_perf_metrics), + is_pytorch_backend=self._is_pytorch_backend), # Beam search enforces return_all_generated_tokens=True regardless of the passed value return_all_generated_tokens=False, # convert python config into pybind config diff --git a/tensorrt_llm/llmapi/llm.py b/tensorrt_llm/llmapi/llm.py index 12bb079eaf5a..b578ba072116 100644 --- a/tensorrt_llm/llmapi/llm.py +++ b/tensorrt_llm/llmapi/llm.py @@ -548,7 +548,7 @@ def _prepare_sampling_params( if sampling_params._stream_interval is None: sampling_params._stream_interval = getattr(self.args, "stream_interval", 1) - + sampling_params.return_perf_metrics = sampling_params.return_perf_metrics or self.args.return_perf_metrics return sampling_params def _check_arguments(self, prompt_len: int, query_len: int, diff --git a/tensorrt_llm/metrics/__init__.py b/tensorrt_llm/metrics/__init__.py index 339982c2d61d..f68d9f698ace 100644 --- a/tensorrt_llm/metrics/__init__.py +++ b/tensorrt_llm/metrics/__init__.py @@ -1,4 +1,4 @@ from .collector import * from .enums import * -__all__ = ["SupportedMetricNames", "MetricsCollector", "RequestEventTiming"] +__all__ = ["MetricsCollector", "MetricNames", "RequestEventTiming"] diff --git a/tensorrt_llm/metrics/collector.py b/tensorrt_llm/metrics/collector.py index 7d6a4558c69d..952529393c64 100644 --- a/tensorrt_llm/metrics/collector.py +++ b/tensorrt_llm/metrics/collector.py @@ -3,7 +3,7 @@ import time from typing import Dict, Optional, Union -from .enums import SupportedMetricNames +from .enums import MetricNames # Adapted from https://github.com/vllm-project/vllm/blob/v0.10.0rc1/vllm/engine/metrics.py#L30 @@ -86,14 +86,13 @@ def log_request_success(self, data: Union[int, float], self.last_log_time = time.time() def log_histogram(self, data: Optional[dict[str, float]]) -> None: - if e2e := data.get(SupportedMetricNames.E2E, 0): + if e2e := data.get(MetricNames.E2E, 0): self._log_histogram(self.histogram_e2e_time_request, e2e) - if ttft := data.get(SupportedMetricNames.TTFT, 0): + if ttft := data.get(MetricNames.TTFT, 0): self._log_histogram(self.histogram_time_to_first_token, ttft) - if tpot := data.get(SupportedMetricNames.TPOT, 0): + if tpot := data.get(MetricNames.TPOT, 0): self._log_histogram(self.histogram_time_per_output_token, tpot) - if request_queue_time := data.get( - SupportedMetricNames.REQUEST_QUEUE_TIME, 0): + if request_queue_time := data.get(MetricNames.REQUEST_QUEUE_TIME, 0): self._log_histogram(self.histogram_queue_time_request, request_queue_time) self.last_log_time = time.time() diff --git a/tensorrt_llm/metrics/enums.py b/tensorrt_llm/metrics/enums.py index 442f7bbbc54f..fb34b9b8ccc4 100644 --- a/tensorrt_llm/metrics/enums.py +++ b/tensorrt_llm/metrics/enums.py @@ -1,14 +1,14 @@ -from enum import Enum +from enum import StrEnum -class SupportedMetricNames(Enum): +class MetricNames(StrEnum): TTFT = "ttft" TPOT = "tpot" E2E = "e2e" REQUEST_QUEUE_TIME = "request_queue_time" -class RequestEventTiming(Enum): +class RequestEventTiming(StrEnum): ARRIVAL_TIME = "arrival_time" FIRST_TOKEN_TIME = "first_token_time" # nosec: B105 FIRST_SCHEDULED_TIME = "first_scheduled_time" diff --git a/tensorrt_llm/sampling_params.py b/tensorrt_llm/sampling_params.py index 94ea4dddad97..361c0fc0c0f4 100644 --- a/tensorrt_llm/sampling_params.py +++ b/tensorrt_llm/sampling_params.py @@ -437,9 +437,7 @@ def _get_sampling_config(self) -> tllme.SamplingConfig: return tllme.SamplingConfig(**llmapi_to_rt_param_map) - def _get_output_config( - self, is_pytorch_backend: bool = False, return_perf_metrics: Optional[bool] = False - ) -> tllme.OutputConfig: + def _get_output_config(self, is_pytorch_backend: bool = False) -> tllme.OutputConfig: sampling_param_fields = set(dir(SamplingParams)) fields = [ f @@ -453,7 +451,6 @@ def _get_output_config( config_kwargs["return_log_probs"] = bool(self.logprobs) else: config_kwargs["return_log_probs"] = self._return_log_probs - config_kwargs["return_perf_metrics"] = return_perf_metrics return tllme.OutputConfig(**config_kwargs) diff --git a/tests/unittest/api_stability/references/llm.yaml b/tests/unittest/api_stability/references/llm.yaml index 5ac588dc9116..5a846dd7869f 100644 --- a/tests/unittest/api_stability/references/llm.yaml +++ b/tests/unittest/api_stability/references/llm.yaml @@ -27,6 +27,10 @@ methods: annotation: Optional[int] default: null status: prototype + return_perf_metrics: + annotation: bool + default: False + status: prototype # Bindings and mirrored configs peft_cache_config: annotation: Optional[tensorrt_llm.llmapi.llm_args.PeftCacheConfig] From 65c6bcc8883a2aaf80b2df9a41c047acd256b9ee Mon Sep 17 00:00:00 2001 From: Ye Zhang Date: Mon, 4 Aug 2025 15:20:35 +0800 Subject: [PATCH 09/12] chore: move prometheus dep. to trtllm-serve Signed-off-by: Ye Zhang --- tensorrt_llm/_utils.py | 5 +---- tensorrt_llm/executor/proxy.py | 24 +++-------------------- tensorrt_llm/executor/result.py | 2 ++ tensorrt_llm/serve/openai_server.py | 23 +++++++++++++++++++--- tests/unittest/llmapi/test_llm_pytorch.py | 22 +++++++++++++++++++++ 5 files changed, 48 insertions(+), 28 deletions(-) diff --git a/tensorrt_llm/_utils.py b/tensorrt_llm/_utils.py index 7c6179148435..c68777b96a77 100644 --- a/tensorrt_llm/_utils.py +++ b/tensorrt_llm/_utils.py @@ -1116,10 +1116,7 @@ def is_multi_device_enable(): def set_prometheus_multiproc_dir() -> object: - # Set prometheus multiprocess directory - # sglang uses prometheus multiprocess mode - # we need to set this before importing prometheus_client - # https://prometheus.github.io/client_python/multiprocess/ + # Adapted from: https://github.com/sgl-project/sglang/blob/v0.4.10/python/sglang/srt/utils.py#L1266 global prometheus_multiproc_dir if "PROMETHEUS_MULTIPROC_DIR" in os.environ: logger.info("User set PROMETHEUS_MULTIPROC_DIR detected.") diff --git a/tensorrt_llm/executor/proxy.py b/tensorrt_llm/executor/proxy.py index 0b46ddddf1c3..1cb86dfdff71 100644 --- a/tensorrt_llm/executor/proxy.py +++ b/tensorrt_llm/executor/proxy.py @@ -10,10 +10,8 @@ import zmq.asyncio from tensorrt_llm.logger import logger -from tensorrt_llm.metrics.collector import MetricsCollector -from .._utils import (customized_gc_thresholds, mpi_rank, nvtx_range_debug, - set_prometheus_multiproc_dir) +from .._utils import customized_gc_thresholds, mpi_rank, nvtx_range_debug from ..llmapi.mpi_session import (MpiCommSession, MpiPoolSession, MpiSession, RemoteMpiCommSessionClient) from ..llmapi.tracer import enable_llm_tracer, get_tracer, global_tracer @@ -22,9 +20,9 @@ print_colored_debug) from .executor import GenerationExecutor from .ipc import FusedIpcQueue, IpcQueue -from .postproc_worker import PostprocWorker, PostprocWorkerConfig +from .postproc_worker import PostprocWorkerConfig from .request import CancellingRequest, GenerationRequest -from .result import GenerationResult, IterationResult, ResponseWrapper +from .result import GenerationResult, IterationResult from .utils import (ErrorResponse, IntraProcessQueue, WorkerCommIpcAddrs, create_mpi_comm_session, get_spawn_proxy_process_env, is_llm_response, print_alive_threads) @@ -62,12 +60,6 @@ def __init__( self.workers_started = False self.worker_cls = worker_cls - set_prometheus_multiproc_dir() - self.metrics_collector = MetricsCollector({ - "model_name": "undefined", - "engine_type": "undefined" - }) - mpi_process_pre_spawned: bool = get_spawn_proxy_process_env() if mpi_session is None: @@ -186,16 +178,6 @@ def process_res(res): else: queue.put(res) - # log metrics to prometheus when response is finished and return_perf_metrics is on - if isinstance(res, ResponseWrapper): - if res._response and res._response.result and res._response.result.is_final: - self._results[client_id]._handle_response(res) - if metrics_dict := self._results[client_id].metrics_dict: - self.metrics_collector.log_metrics_dict(metrics_dict) - if isinstance(res, PostprocWorker.Output): - if metrics_dict := res.metrics: - self.metrics_collector.log_metrics_dict(metrics_dict) - if (is_llm_response(res) and res.result.is_final) or isinstance( res, ErrorResponse): self._results.pop(client_id) diff --git a/tensorrt_llm/executor/result.py b/tensorrt_llm/executor/result.py index 1f7adf77e1c8..1b5561455311 100644 --- a/tensorrt_llm/executor/result.py +++ b/tensorrt_llm/executor/result.py @@ -310,6 +310,8 @@ def _handle_response(self, self._outputs[0] = response.res else: self._outputs[0]._postprocess_result = response.res + if response.metrics: + self.metrics_dict = response.metrics if response.error: if self._background_error_handler is not None and ( diff --git a/tensorrt_llm/serve/openai_server.py b/tensorrt_llm/serve/openai_server.py index f4001c0c2802..1ceeb77431d9 100644 --- a/tensorrt_llm/serve/openai_server.py +++ b/tensorrt_llm/serve/openai_server.py @@ -27,6 +27,7 @@ from tensorrt_llm.llmapi.disagg_utils import MetadataServerConfig, ServerRole from tensorrt_llm.llmapi.llm import RequestOutput from tensorrt_llm.logger import logger +from tensorrt_llm.metrics.collector import MetricsCollector from tensorrt_llm.serve.chat_utils import (check_multiple_response, parse_chat_messages_coroutines) from tensorrt_llm.serve.metadata_server import create_metadata_server @@ -44,7 +45,7 @@ completion_stream_post_processor) from tensorrt_llm.version import __version__ as VERSION -from .._utils import nvtx_mark +from .._utils import nvtx_mark, set_prometheus_multiproc_dir # yapf: enale TIMEOUT_KEEP_ALIVE = 5 # seconds. @@ -81,6 +82,13 @@ def __init__(self, else: self.model = model + if self.llm.args.return_perf_metrics: + set_prometheus_multiproc_dir() + self.metrics_collector = MetricsCollector({ + "model_name": "undefined", + "engine_type": "undefined" + }) + @asynccontextmanager async def lifespan(app: FastAPI): if self.metadata_server is not None: @@ -153,8 +161,9 @@ def register_routes(self): self.app.add_api_route("/v1/chat/completions", self.openai_chat, methods=["POST"]) - # register /prometheus/metrics - self.mount_metrics() + if self.llm.args.return_perf_metrics: + # register /prometheus/metrics + self.mount_metrics() def mount_metrics(self): # Lazy import for prometheus multiprocessing. @@ -255,6 +264,8 @@ async def chat_stream_generator( post_processor, args = postproc_params.post_processor, postproc_params.postproc_args async for res in promise: pp_results = res.outputs[0]._postprocess_result if self.postproc_worker_enabled else post_processor(res, args) + if res.finished and self.metrics_collector: + self.metrics_collector.log_metrics_dict(res.metrics_dict) for pp_res in pp_results: yield pp_res yield "data: [DONE]\n\n" @@ -272,6 +283,8 @@ async def create_chat_response( # Add prompt_tokens_ids to the response if disaggregated_params and disaggregated_params.request_type and disaggregated_params.request_type == "context_only": chat_response.prompt_token_ids = promise.prompt_token_ids + if promise.finished and self.metrics_collector: + self.metrics_collector.log_metrics_dict(promise.metrics_dict) return chat_response try: @@ -364,6 +377,8 @@ async def completion_response(promise: RequestOutput, if disaggregated_params and disaggregated_params.request_type and disaggregated_params.request_type == "context_only": # Include prompt token ids for context-only requests pp_result.prompt_token_ids = response.prompt_token_ids + if response.finished and self.metrics_collector: + self.metrics_collector.log_metrics_dict(response.metrics_dict) return pp_result def merge_completion_responses(responses: List[CompletionResponse]) -> CompletionResponse: @@ -399,6 +414,8 @@ async def completion_generator(promise: RequestOutput, params: Optional[Postproc pp_result = post_processor(output, args) else: pp_result = output.outputs[0]._postprocess_result + if output.finished and self.metrics_collector: + self.metrics_collector.log_metrics_dict(output.metrics_dict) for pp_res in pp_result: yield pp_res diff --git a/tests/unittest/llmapi/test_llm_pytorch.py b/tests/unittest/llmapi/test_llm_pytorch.py index bb5d028b5564..541965b588f0 100644 --- a/tests/unittest/llmapi/test_llm_pytorch.py +++ b/tests/unittest/llmapi/test_llm_pytorch.py @@ -6,6 +6,7 @@ from tensorrt_llm.llmapi import KvCacheConfig from tensorrt_llm.llmapi.llm_args import PeftCacheConfig from tensorrt_llm.llmapi.tokenizer import TransformersTokenizer +from tensorrt_llm.metrics import MetricNames from tensorrt_llm.sampling_params import SamplingParams # isort: off @@ -195,6 +196,27 @@ def test_llm_perf_metrics(): assert perf_metrics.last_iter == perf_metrics.iter +def test_llm_prometheus(): + test_prompts = [ + "Hello, my name is", + "The president of the United States is", + "The capital of France is", + "The future of AI is", + ] + sampling_params = SamplingParams(max_tokens=10, temperature=0.8, top_p=0.95) + llm = LLM(model=llama_model_path, + return_perf_metrics=True, + kv_cache_config=global_kvcache_config) + for test_prompt in test_prompts: + request_output = llm.generate(test_prompt, sampling_params) + assert request_output.metrics_dict is not None + assert MetricNames.REQUEST_QUEUE_TIME in request_output.metrics_dict + assert MetricNames.TPOT in request_output.metrics_dict + assert MetricNames.TTFT in request_output.metrics_dict + assert MetricNames.E2E in request_output.metrics_dict + assert request_output.outputs is not None + + @pytest.mark.parametrize("streaming", [True, False]) def test_llm_with_postprocess_parallel_and_result_handler(streaming): run_llm_with_postprocess_parallel_and_result_handler(streaming, From 097d4ece6d096715292abbfe3ab2c71caa828708 Mon Sep 17 00:00:00 2001 From: Ye Zhang Date: Thu, 7 Aug 2025 10:51:21 +0800 Subject: [PATCH 10/12] fix: CI test Signed-off-by: Ye Zhang --- tensorrt_llm/metrics/enums.py | 6 +++--- tensorrt_llm/serve/openai_server.py | 2 +- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/tensorrt_llm/metrics/enums.py b/tensorrt_llm/metrics/enums.py index fb34b9b8ccc4..5ce982281bcd 100644 --- a/tensorrt_llm/metrics/enums.py +++ b/tensorrt_llm/metrics/enums.py @@ -1,14 +1,14 @@ -from enum import StrEnum +from enum import Enum -class MetricNames(StrEnum): +class MetricNames(Enum): TTFT = "ttft" TPOT = "tpot" E2E = "e2e" REQUEST_QUEUE_TIME = "request_queue_time" -class RequestEventTiming(StrEnum): +class RequestEventTiming(Enum): ARRIVAL_TIME = "arrival_time" FIRST_TOKEN_TIME = "first_token_time" # nosec: B105 FIRST_SCHEDULED_TIME = "first_scheduled_time" diff --git a/tensorrt_llm/serve/openai_server.py b/tensorrt_llm/serve/openai_server.py index 1ceeb77431d9..1b1e15ec6259 100644 --- a/tensorrt_llm/serve/openai_server.py +++ b/tensorrt_llm/serve/openai_server.py @@ -81,7 +81,7 @@ def __init__(self, self.model = model_dir.name else: self.model = model - + self.metrics_collector = None if self.llm.args.return_perf_metrics: set_prometheus_multiproc_dir() self.metrics_collector = MetricsCollector({ From 8ebae0ece6ad3e3a882b443472a46b8b26de8463 Mon Sep 17 00:00:00 2001 From: Ye Zhang Date: Fri, 8 Aug 2025 10:35:00 +0800 Subject: [PATCH 11/12] fix: add status field to fix CI tests Signed-off-by: Ye Zhang --- tensorrt_llm/llmapi/llm_args.py | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/tensorrt_llm/llmapi/llm_args.py b/tensorrt_llm/llmapi/llm_args.py index 3e42b7595996..279d26999b27 100644 --- a/tensorrt_llm/llmapi/llm_args.py +++ b/tensorrt_llm/llmapi/llm_args.py @@ -1311,10 +1311,9 @@ class BaseLlmArgs(StrictBaseModel): status="deprecated", ) - return_perf_metrics: bool = Field( - default=False, - description="Return perf metrics", - ) + return_perf_metrics: bool = Field(default=False, + description="Return perf metrics.", + status="prototype") _parallel_config: Optional[object] = PrivateAttr(default=None) _model_format: Optional[_ModelFormatKind] = PrivateAttr(default=None) From d0f4acf123362e4b475d517a4cec4191e90ea258 Mon Sep 17 00:00:00 2001 From: Shunkang <182541032+Shunkangz@users.noreply.github.co> Date: Fri, 8 Aug 2025 07:31:25 +0000 Subject: [PATCH 12/12] Fix request output reference Signed-off-by: Shunkang <182541032+Shunkangz@users.noreply.github.co> --- tensorrt_llm/executor/result.py | 8 +++++++- .../api_stability/references/request_output.yaml | 9 +++++++++ 2 files changed, 16 insertions(+), 1 deletion(-) diff --git a/tensorrt_llm/executor/result.py b/tensorrt_llm/executor/result.py index 1b5561455311..2566a699aa4d 100644 --- a/tensorrt_llm/executor/result.py +++ b/tensorrt_llm/executor/result.py @@ -367,7 +367,13 @@ def _handle_response(self, def record_stats(self, output: CompletionOutput, - stats: Optional[dict[str, float]] = None): + stats: Optional[dict[str, float]] = None) -> None: + """Record the stats of the generation result. + + Args: + output (CompletionOutput): The output of the generation result. + stats (Optional[dict[str, float]]): The stats of the generation result. Defaults to None. + """ if not stats: return metrics_stats = {} diff --git a/tests/unittest/api_stability/references/request_output.yaml b/tests/unittest/api_stability/references/request_output.yaml index 52e499dd1476..7e3054cd5ef6 100644 --- a/tests/unittest/api_stability/references/request_output.yaml +++ b/tests/unittest/api_stability/references/request_output.yaml @@ -11,4 +11,13 @@ methods: clear_logprob_params: parameters: {} return_annotation: None + record_stats: + parameters: + output: + annotation: tensorrt_llm.executor.result.CompletionOutput + default: inspect._empty + stats: + annotation: Optional[dict[str, float]] + default: None + return_annotation: None properties: {}