feat(datadog): emit latency metrics as distributions and drop empty HOSTNAME tag

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
Devin AI 2026-09-03 04:46:33 +00:00
parent ff17e8b987
commit 39bda42458
3 changed files with 210 additions and 45 deletions

View file

@ -2,8 +2,9 @@ import asyncio
import gzip
import os
import time
import zlib
from datetime import datetime
from typing import Final
from typing import Final, Literal
from litellm._logging import verbose_logger
from litellm.integrations.custom_batch_logger import CustomBatchLogger
@ -21,6 +22,8 @@ from litellm.llms.custom_httpx.http_handler import (
)
from litellm.types.integrations.base_health_check import IntegrationHealthCheckStatus
from litellm.types.integrations.datadog_metrics import (
DatadogDistributionPayload,
DatadogDistributionSeries,
DatadogMetricPoint,
DatadogMetricSeries,
DatadogMetricsPayload,
@ -38,6 +41,7 @@ class DatadogMetricsLogger(CustomBatchLogger):
verbose_logger.warning("Datadog Metrics: DD_API_KEY is required. Integration will not work.")
self.upload_url = f"https://api.{self.dd_site}/api/v2/series"
self.distribution_upload_url = f"https://api.{self.dd_site}/api/v1/distribution_points"
self.async_client = get_async_httpx_client(llm_provider=httpxSpecialProvider.LoggingCallback)
@ -71,10 +75,12 @@ class DatadogMetricsLogger(CustomBatchLogger):
f"env:{get_datadog_env()}",
f"service:{get_datadog_service()}",
f"version:{os.getenv('DD_VERSION', 'unknown')}",
f"HOSTNAME:{get_datadog_hostname()}",
f"POD_NAME:{get_datadog_pod_name()}",
]
if hostname := get_datadog_hostname():
tags.append(f"HOSTNAME:{hostname}")
# Add metric-specific tags
if provider := log.get("custom_llm_provider"):
tags.append(f"provider:{provider}")
@ -102,6 +108,25 @@ class DatadogMetricsLogger(CustomBatchLogger):
return tags
def _add_latency_metric(self, metric: str, seconds: float, timestamp: int, tags: list[str]) -> None:
"""
Queues a latency sample as a gauge (legacy metric name) and as a distribution
(`<metric>.distribution`) so Datadog computes percentiles over every request
"""
gauge: Final[DatadogMetricSeries] = {
"metric": metric,
"type": 3, # gauge
"points": [{"timestamp": timestamp, "value": seconds}],
"tags": tags,
}
distribution: Final[DatadogDistributionSeries] = {
"metric": f"{metric}.distribution",
"type": "distribution",
"points": ((timestamp, (seconds,)),),
"tags": tuple(tags),
}
self.log_queue.extend((gauge, distribution))
def _add_metrics_from_log(
self,
log: StandardLoggingPayload,
@ -121,43 +146,20 @@ class DatadogMetricsLogger(CustomBatchLogger):
start_time_dt: Final = kwargs.get("start_time")
if start_time_dt and end_time_dt:
total_duration: Final = (end_time_dt - start_time_dt).total_seconds()
series_total_latency: Final[DatadogMetricSeries] = {
"metric": "litellm.request.total_latency",
"type": 3, # gauge
"points": [{"timestamp": timestamp, "value": total_duration}],
"tags": tags,
}
self.log_queue.append(series_total_latency)
self._add_latency_metric("litellm.request.total_latency", total_duration, timestamp, tags)
# 2. LLM API Latency Metric (Provider alone)
api_call_start_time: Final = kwargs.get("api_call_start_time")
if api_call_start_time and end_time_dt:
llm_api_duration: Final = (end_time_dt - api_call_start_time).total_seconds()
series_llm_latency: Final[DatadogMetricSeries] = {
"metric": "litellm.llm_api.latency",
"type": 3, # gauge
"points": [{"timestamp": timestamp, "value": llm_api_duration}],
"tags": tags,
}
self.log_queue.append(series_llm_latency)
self._add_latency_metric("litellm.llm_api.latency", llm_api_duration, timestamp, tags)
# 3. LiteLLM Overhead Latency Metric (total - llm_api time)
hidden_params: Final = log.get("hidden_params", {}) or {}
litellm_overhead_time_ms: Final = hidden_params.get("litellm_overhead_time_ms")
if litellm_overhead_time_ms is not None:
overhead_tags: Final = self._extract_tags(log) # no status_code on latency metric
series_overhead: Final[DatadogMetricSeries] = {
"metric": "litellm.overhead.latency",
"type": 3, # gauge
"points": [
{
"timestamp": timestamp,
"value": litellm_overhead_time_ms / 1000, # convert ms → seconds
}
],
"tags": overhead_tags,
}
self.log_queue.append(series_overhead)
self._add_latency_metric("litellm.overhead.latency", litellm_overhead_time_ms / 1000, timestamp, overhead_tags)
# 4. Request Count / Status Code
series_count: Final[DatadogMetricSeries] = {
@ -210,42 +212,49 @@ class DatadogMetricsLogger(CustomBatchLogger):
if not self.log_queue:
return
batch: Final = self.log_queue.copy()
payload_data: Final[DatadogMetricsPayload] = {"series": batch}
batch: Final = tuple(self.log_queue)
series: Final[list[DatadogMetricSeries]] = [s for s in batch if s["type"] != "distribution"]
distributions: Final[tuple[DatadogDistributionSeries, ...]] = tuple(
s for s in batch if s["type"] == "distribution"
)
try:
await self._upload_to_datadog(payload_data)
if series:
await self._upload_to_datadog({"series": series})
if distributions:
await self._upload_distributions_to_datadog({"series": distributions})
except Exception as e:
verbose_logger.exception("Datadog Metrics: Error in async_send_batch: %s", e)
raise
async def _upload_to_datadog(self, payload: DatadogMetricsPayload):
await self._post_compressed(self.upload_url, safe_dumps(payload), encoding="gzip")
async def _upload_distributions_to_datadog(self, payload: DatadogDistributionPayload):
# /api/v1/distribution_points only accepts deflate-compressed bodies
await self._post_compressed(self.distribution_upload_url, safe_dumps(payload), encoding="deflate")
async def _post_compressed(self, url: str, json_data: str, encoding: Literal["gzip", "deflate"]) -> None:
if not self.dd_api_key:
return
headers: Final = {
"Content-Type": "application/json",
"Content-Encoding": encoding,
"DD-API-KEY": self.dd_api_key,
}
if self.dd_app_key:
headers["DD-APPLICATION-KEY"] = self.dd_app_key
json_data: Final = safe_dumps(payload)
compressed_data: Final = gzip.compress(json_data.encode("utf-8"))
headers["Content-Encoding"] = "gzip"
raw: Final = json_data.encode("utf-8")
compressed_data: Final = gzip.compress(raw) if encoding == "gzip" else zlib.compress(raw)
response: Final = await self.async_client.post(
self.upload_url,
content=compressed_data,
headers=headers,
)
response: Final = await self.async_client.post(url, content=compressed_data, headers=headers)
response.raise_for_status()
verbose_logger.debug(
"Datadog Metrics: Uploaded %s metric points. Status: %s", len(payload["series"]), response.status_code
)
verbose_logger.debug("Datadog Metrics: Uploaded metrics to %s. Status: %s", url, response.status_code)
async def async_health_check(self) -> IntegrationHealthCheckStatus:
"""

View file

@ -1,4 +1,6 @@
from typing_extensions import TypedDict
from typing import Literal
from typing_extensions import ReadOnly, TypedDict
class DatadogMetricPoint(TypedDict):
@ -16,3 +18,17 @@ class DatadogMetricSeries(TypedDict, total=False):
class DatadogMetricsPayload(TypedDict):
series: list[DatadogMetricSeries]
DatadogDistributionPoint = tuple[int, tuple[float, ...]]
class DatadogDistributionSeries(TypedDict):
metric: ReadOnly[str]
type: ReadOnly[Literal["distribution"]]
points: ReadOnly[tuple[DatadogDistributionPoint, ...]]
tags: ReadOnly[tuple[str, ...]]
class DatadogDistributionPayload(TypedDict):
series: ReadOnly[tuple[DatadogDistributionSeries, ...]]

View file

@ -125,9 +125,9 @@ async def test_add_metrics_from_log(clean_env):
logger._add_metrics_from_log(log=payload, kwargs=kwargs, status_code="200")
# Should have 3 series: total_latency, llm_api_latency, request_count
# total_latency and llm_api_latency (each as gauge + distribution) plus request_count
# (no overhead metric because payload has no hidden_params litellm_overhead_time_ms)
assert len(logger.log_queue) == 3
assert len(logger.log_queue) == 5
metrics = {s["metric"]: s for s in logger.log_queue}
@ -148,6 +148,65 @@ async def test_add_metrics_from_log(clean_env):
assert "status_code:200" in count["tags"]
@pytest.mark.asyncio
async def test_latency_metrics_also_emitted_as_distributions(clean_env):
"""Each latency sample is queued as a distribution point so Datadog can compute per-request percentiles."""
logger = DatadogMetricsLogger(batch_size=100, start_periodic_flush=False)
now = datetime.now()
payload = StandardLoggingPayload(
custom_llm_provider="openai",
model="gpt-4o",
hidden_params={"litellm_overhead_time_ms": 250},
)
kwargs = {
"start_time": now - timedelta(seconds=2),
"api_call_start_time": now - timedelta(seconds=1),
"end_time": now,
}
logger._add_metrics_from_log(log=payload, kwargs=kwargs, status_code="200")
distributions = {s["metric"]: s for s in logger.log_queue if s["type"] == "distribution"}
assert set(distributions) == {
"litellm.request.total_latency.distribution",
"litellm.llm_api.latency.distribution",
"litellm.overhead.latency.distribution",
}
expected_seconds = {
"litellm.request.total_latency.distribution": 2.0,
"litellm.llm_api.latency.distribution": 1.0,
"litellm.overhead.latency.distribution": 0.25,
}
for metric, seconds in expected_seconds.items():
((timestamp, values),) = distributions[metric]["points"]
assert timestamp == int(now.timestamp())
assert len(values) == 1
assert abs(values[0] - seconds) < 0.1
assert "provider:openai" in distributions[metric]["tags"]
assert "model_name:gpt-4o" in distributions[metric]["tags"]
assert "status_code:200" in distributions["litellm.request.total_latency.distribution"]["tags"]
assert not any(
tag.startswith("status_code:") for tag in distributions["litellm.overhead.latency.distribution"]["tags"]
)
@pytest.mark.asyncio
async def test_extract_tags_omits_hostname_when_unset(clean_env, monkeypatch: pytest.MonkeyPatch):
"""An unset HOSTNAME must not produce an empty `HOSTNAME:` tag."""
monkeypatch.delenv("HOSTNAME", raising=False)
logger = DatadogMetricsLogger(start_periodic_flush=False)
tags = logger._extract_tags(log=StandardLoggingPayload(model="gpt-4o"))
assert not any(tag.startswith("HOSTNAME") for tag in tags)
monkeypatch.setenv("HOSTNAME", "pod-abc")
assert "HOSTNAME:pod-abc" in logger._extract_tags(log=StandardLoggingPayload(model="gpt-4o"))
@pytest.mark.asyncio
async def test_overhead_latency_metric_emitted(clean_env):
"""Test that litellm.overhead.latency is emitted when hidden_params contains litellm_overhead_time_ms."""
@ -351,6 +410,87 @@ async def test_async_send_batch(clean_env):
assert payload["series"][0]["metric"] == "litellm.request.total_latency"
@pytest.mark.asyncio
async def test_async_send_batch_routes_distributions_to_v1_endpoint(clean_env):
"""Distribution series go to /api/v1/distribution_points (deflate), gauges/counts stay on /api/v2/series."""
import gzip
import json
import zlib
logger = DatadogMetricsLogger(start_periodic_flush=False)
logger.async_client = AsyncMock()
logger.async_client.post.return_value = Response(
202, json={"status": "ok"}, request=Request("POST", "https://api.test.datadoghq.com")
)
timestamp = int(time.time())
logger.log_queue = [
{
"metric": "litellm.request.total_latency",
"type": 3,
"points": [{"timestamp": timestamp, "value": 1.5}],
"tags": ["env:test"],
},
{
"metric": "litellm.request.total_latency.distribution",
"type": "distribution",
"points": ((timestamp, (1.5,)),),
"tags": ("env:test",),
},
]
await logger.async_send_batch()
calls = {call.args[0]: call.kwargs for call in logger.async_client.post.call_args_list}
assert set(calls) == {
"https://api.test.datadoghq.com/api/v2/series",
"https://api.test.datadoghq.com/api/v1/distribution_points",
}
series_call = calls["https://api.test.datadoghq.com/api/v2/series"]
assert series_call["headers"]["Content-Encoding"] == "gzip"
series_payload = json.loads(gzip.decompress(series_call["content"]))
assert [s["metric"] for s in series_payload["series"]] == ["litellm.request.total_latency"]
distribution_call = calls["https://api.test.datadoghq.com/api/v1/distribution_points"]
assert distribution_call["headers"]["Content-Encoding"] == "deflate"
assert distribution_call["headers"]["DD-API-KEY"] == "test_api_key"
distribution_payload = json.loads(zlib.decompress(distribution_call["content"]))
assert distribution_payload == {
"series": [
{
"metric": "litellm.request.total_latency.distribution",
"type": "distribution",
"points": [[timestamp, [1.5]]],
"tags": ["env:test"],
}
]
}
@pytest.mark.asyncio
async def test_async_send_batch_skips_v2_when_only_distributions_queued(clean_env):
logger = DatadogMetricsLogger(start_periodic_flush=False)
logger.async_client = AsyncMock()
logger.async_client.post.return_value = Response(
202, json={"status": "ok"}, request=Request("POST", "https://api.test.datadoghq.com")
)
logger.log_queue = [
{
"metric": "litellm.llm_api.latency.distribution",
"type": "distribution",
"points": ((int(time.time()), (0.4,)),),
"tags": ("env:test",),
}
]
await logger.async_send_batch()
assert [call.args[0] for call in logger.async_client.post.call_args_list] == [
"https://api.test.datadoghq.com/api/v1/distribution_points"
]
@pytest.mark.asyncio
async def test_async_send_batch_empty_queue(clean_env):
"""Test that async_send_batch does nothing when queue is empty."""