From 81d61bff4a7492ef9c7e5673322ea7df27f673e9 Mon Sep 17 00:00:00 2001 From: David Traina <44659830+DavidTraina@users.noreply.github.com> Date: Fri, 2 Oct 2026 15:56:26 -0500 Subject: [PATCH 1/3] feat(otel): add opt-in gzip compression for the default OTLP exporter --- langfuse/_client/client.py | 5 +- langfuse/_client/environment_variables.py | 13 ++ langfuse/_client/get_client.py | 1 + langfuse/_client/resource_manager.py | 7 +- langfuse/_client/span_processor.py | 34 ++++- tests/unit/test_span_processor.py | 164 ++++++++++++++++++++-- 6 files changed, 210 insertions(+), 14 deletions(-) diff --git a/langfuse/_client/client.py b/langfuse/_client/client.py index a56188ab3..af06e69ef 100644 --- a/langfuse/_client/client.py +++ b/langfuse/_client/client.py @@ -268,7 +268,8 @@ def mask_otel_spans( additional_headers (Optional[Dict[str, str]]): Additional headers to include in all API requests and in the default OTLPSpanExporter requests. These headers will be merged with default headers. Note: If httpx_client is provided, additional_headers must be set directly on your custom httpx_client as well. If `span_exporter` is provided, these headers are not wired into that exporter and must be configured on the exporter instance directly. tracer_provider(Optional[TracerProvider]): OpenTelemetry TracerProvider to use for Langfuse. This can be useful to set to have disconnected tracing between Langfuse and other OpenTelemetry-span emitting libraries. Note: To track active spans, the context is still shared between TracerProviders. This may lead to broken trace trees. id_generator (Optional[IdGenerator]): OpenTelemetry ID generator to use when Langfuse creates its own TracerProvider. If omitted, the OpenTelemetry SDK default is used. If `tracer_provider` is provided, or an OpenTelemetry TracerProvider is already registered globally, configure the ID generator on that provider instead. - span_exporter (Optional[SpanExporter]): Custom OpenTelemetry span exporter for the Langfuse span processor. If omitted, Langfuse creates an OTLPSpanExporter pointed at the Langfuse OTLP endpoint. If provided, Langfuse does not wire `base_url`, exporter headers, exporter auth, exporter timeout, or the `LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES` request size limit into it. Configure endpoint, headers, and timeout on the exporter instance directly. If you are sending spans to Langfuse v4 or using Langfuse Cloud Fast Preview, include `x-langfuse-ingestion-version=4` on the exporter to enable real time processing of exported spans. + span_exporter (Optional[SpanExporter]): Custom OpenTelemetry span exporter for the Langfuse span processor. If omitted, Langfuse creates an OTLPSpanExporter pointed at the Langfuse OTLP endpoint. If provided, Langfuse does not wire `base_url`, exporter headers, exporter auth, exporter timeout, `otel_compression`, or the `LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES` request size limit into it. Configure endpoint, headers, timeout, and compression on the exporter instance directly. If you are sending spans to Langfuse v4 or using Langfuse Cloud Fast Preview, include `x-langfuse-ingestion-version=4` on the exporter to enable real time processing of exported spans. + otel_compression (Optional[Literal["gzip", "none"]]): Compression for span batches sent by the default OTLP span exporter. Use "gzip" to reduce network bytes or "none" to send uncompressed. Can also be set via LANGFUSE_OTEL_COMPRESSION environment variable. If unset, the standard OTEL_EXPORTER_OTLP_TRACES_COMPRESSION and OTEL_EXPORTER_OTLP_COMPRESSION environment variables apply (no compression by default). "gzip" requires Langfuse server v3.30.0 or later. Example: ```python @@ -335,6 +336,7 @@ def __init__( tracer_provider: Optional[TracerProvider] = None, id_generator: Optional[IdGenerator] = None, span_exporter: Optional[SpanExporter] = None, + otel_compression: Optional[Literal["gzip", "none"]] = None, ): self._base_url = ( base_url @@ -435,6 +437,7 @@ def __init__( tracer_provider=tracer_provider, id_generator=id_generator, span_exporter=span_exporter, + otel_compression=otel_compression, ) self._mask = self._resources.mask diff --git a/langfuse/_client/environment_variables.py b/langfuse/_client/environment_variables.py index 424ba53fd..b340c9ac4 100644 --- a/langfuse/_client/environment_variables.py +++ b/langfuse/_client/environment_variables.py @@ -73,6 +73,19 @@ **Default value:** ``67108864`` (64 MiB) """ +LANGFUSE_OTEL_COMPRESSION = "LANGFUSE_OTEL_COMPRESSION" +""" +.. envvar:: LANGFUSE_OTEL_COMPRESSION + +Compression for span batches sent by the default OTLP exporter: ``gzip`` or ``none``. +The ``otel_compression`` client argument takes precedence. If unset, the standard +``OTEL_EXPORTER_OTLP_TRACES_COMPRESSION`` and ``OTEL_EXPORTER_OTLP_COMPRESSION`` +environment variables apply. ``gzip`` requires Langfuse server v3.30.0 or later. +Custom span exporters are not affected. + +**Default value:** unset +""" + LANGFUSE_DEBUG = "LANGFUSE_DEBUG" """ .. envvar:: LANGFUSE_DEBUG diff --git a/langfuse/_client/get_client.py b/langfuse/_client/get_client.py index 0c4ccd321..ff06c7d29 100644 --- a/langfuse/_client/get_client.py +++ b/langfuse/_client/get_client.py @@ -57,6 +57,7 @@ def _create_client_from_instance( tracer_provider=instance.tracer_provider, id_generator=instance.id_generator, span_exporter=instance.span_exporter, + otel_compression=instance.otel_compression, httpx_client=instance.httpx_client, ) diff --git a/langfuse/_client/resource_manager.py b/langfuse/_client/resource_manager.py index a61101fe3..4395872db 100644 --- a/langfuse/_client/resource_manager.py +++ b/langfuse/_client/resource_manager.py @@ -21,7 +21,7 @@ import urllib.request import weakref from queue import Full, Queue -from typing import Any, Callable, Dict, List, Optional, cast +from typing import Any, Callable, Dict, List, Literal, Optional, cast import httpx from opentelemetry import trace as otel_trace_api @@ -133,6 +133,7 @@ def __new__( tracer_provider: Optional[TracerProvider] = None, id_generator: Optional[IdGenerator] = None, span_exporter: Optional[SpanExporter] = None, + otel_compression: Optional[Literal["gzip", "none"]] = None, ) -> "LangfuseResourceManager": if public_key in cls._instances: return cls._instances[public_key] @@ -171,6 +172,7 @@ def __new__( tracer_provider=tracer_provider, id_generator=id_generator, span_exporter=span_exporter, + otel_compression=otel_compression, ) cls._instances[public_key] = instance @@ -200,6 +202,7 @@ def _initialize_instance( tracer_provider: Optional[TracerProvider] = None, id_generator: Optional[IdGenerator] = None, span_exporter: Optional[SpanExporter] = None, + otel_compression: Optional[Literal["gzip", "none"]] = None, ) -> None: self.public_key = public_key self.secret_key = secret_key @@ -222,6 +225,7 @@ def _initialize_instance( self.additional_headers = additional_headers self.id_generator = id_generator self.span_exporter = span_exporter + self.otel_compression = otel_compression self.tracer_provider: Optional[TracerProvider] = None self._custom_httpx_client = httpx_client @@ -261,6 +265,7 @@ def _initialize_instance( span_exporter=span_exporter, media_manager=self._media_manager, mask_otel_spans=mask_otel_spans, + otel_compression=otel_compression, ) tracer_provider.add_span_processor(langfuse_processor) diff --git a/langfuse/_client/span_processor.py b/langfuse/_client/span_processor.py index 524b2996d..9c056490f 100644 --- a/langfuse/_client/span_processor.py +++ b/langfuse/_client/span_processor.py @@ -15,10 +15,11 @@ import logging import os import threading -from typing import Callable, Dict, List, Optional, cast +from typing import Callable, Dict, List, Literal, Optional, cast from opentelemetry import context as context_api from opentelemetry.context import Context +from opentelemetry.exporter.otlp.proto.http import Compression from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter from opentelemetry.sdk.trace import ReadableSpan, Span from opentelemetry.sdk.trace.export import BatchSpanProcessor, SpanExporter @@ -28,6 +29,7 @@ from langfuse._client.environment_variables import ( LANGFUSE_FLUSH_AT, LANGFUSE_FLUSH_INTERVAL, + LANGFUSE_OTEL_COMPRESSION, LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES, LANGFUSE_OTEL_TRACES_EXPORT_PATH, ) @@ -65,6 +67,34 @@ def _resolve_max_batch_size_bytes() -> Optional[int]: return None +# The Langfuse OTLP endpoint only decodes gzip, so deflate is not offered. +_COMPRESSION_BY_NAME = {"gzip": Compression.Gzip, "none": Compression.NoCompression} + + +def _resolve_compression(otel_compression: Optional[str]) -> Optional[Compression]: + """Return the configured compression, or None to defer to OTEL_EXPORTER_OTLP_*COMPRESSION.""" + setting = "otel_compression" + raw_value = otel_compression + if raw_value is None: + setting = LANGFUSE_OTEL_COMPRESSION + raw_value = os.environ.get(LANGFUSE_OTEL_COMPRESSION, "") + + value = raw_value.strip().lower() + if not value: + return None + + compression = _COMPRESSION_BY_NAME.get(value) + if compression is None: + langfuse_logger.warning( + "Invalid %s=%r. Expected 'gzip' or 'none'. Falling back to the " + "OTEL_EXPORTER_OTLP_*COMPRESSION environment variables.", + setting, + raw_value, + ) + + return compression + + class LangfuseSpanProcessor(BatchSpanProcessor): """OpenTelemetry span processor that exports spans to the Langfuse API. @@ -97,6 +127,7 @@ def __init__( span_exporter: Optional[SpanExporter] = None, media_manager: Optional[MediaManager] = None, mask_otel_spans: Optional[MaskOtelSpansFunction] = None, + otel_compression: Optional[Literal["gzip", "none"]] = None, ): self.public_key = public_key self.blocked_instrumentation_scopes = ( @@ -145,6 +176,7 @@ def __init__( endpoint=endpoint, headers=headers, timeout=timeout, + compression=_resolve_compression(otel_compression), max_request_size=_resolve_max_batch_size_bytes(), ) diff --git a/tests/unit/test_span_processor.py b/tests/unit/test_span_processor.py index a123fd108..99d7b1c1e 100644 --- a/tests/unit/test_span_processor.py +++ b/tests/unit/test_span_processor.py @@ -1,11 +1,19 @@ +import gzip import logging import threading from http.server import BaseHTTPRequestHandler, HTTPServer -from typing import List, Sequence +from typing import List, NamedTuple, Optional, Sequence from unittest.mock import patch import pytest from opentelemetry.exporter.otlp.proto.common.trace_encoder import encode_spans +from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ( + ExportTraceServiceRequest, +) +from opentelemetry.sdk.environment_variables import ( + OTEL_EXPORTER_OTLP_COMPRESSION, + OTEL_EXPORTER_OTLP_TRACES_COMPRESSION, +) from opentelemetry.sdk.trace import ReadableSpan, TracerProvider from opentelemetry.sdk.trace.export import ( SimpleSpanProcessor, @@ -17,9 +25,11 @@ ) import langfuse._client.span_processor as span_processor_module +from langfuse._client.client import Langfuse from langfuse._client.environment_variables import ( LANGFUSE_FLUSH_AT, LANGFUSE_FLUSH_INTERVAL, + LANGFUSE_OTEL_COMPRESSION, LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES, ) from langfuse._client.span_processor import LangfuseSpanProcessor @@ -71,12 +81,19 @@ def test_span_processor_uses_env_flush_settings_when_constructor_omits_them( processor.shutdown() +class _RecordedRequest(NamedTuple): + content_encoding: Optional[str] + body: bytes + + class _RecordingOTLPHandler(BaseHTTPRequestHandler): - received_body_sizes: List[int] + received_requests: List[_RecordedRequest] def do_POST(self): body = self.rfile.read(int(self.headers.get("Content-Length", 0))) - self.received_body_sizes.append(len(body)) + self.received_requests.append( + _RecordedRequest(self.headers.get("Content-Encoding"), body) + ) self.send_response(200) self.end_headers() @@ -86,17 +103,17 @@ def log_message(self, *args): @pytest.fixture def otlp_http_server(): - received_body_sizes: List[int] = [] + received_requests: List[_RecordedRequest] = [] handler = type( "Handler", (_RecordingOTLPHandler,), - {"received_body_sizes": received_body_sizes}, + {"received_requests": received_requests}, ) server = HTTPServer(("127.0.0.1", 0), handler) thread = threading.Thread(target=server.serve_forever, daemon=True) thread.start() - yield f"http://127.0.0.1:{server.server_port}", received_body_sizes + yield f"http://127.0.0.1:{server.server_port}", received_requests server.shutdown() server.server_close() @@ -136,7 +153,7 @@ def _default_exporter_processor(base_url: str) -> LangfuseSpanProcessor: def test_default_exporter_enforces_max_batch_size_bytes_at_boundary( monkeypatch, otlp_http_server, limit_offset, expected_result ): - base_url, received_body_sizes = otlp_http_server + base_url, received_requests = otlp_http_server spans = _finished_spans("x" * 1_000) request_size = _serialized_request_size(spans) monkeypatch.setenv( @@ -153,13 +170,13 @@ def test_default_exporter_enforces_max_batch_size_bytes_at_boundary( expected_requests = ( [] if expected_result == SpanExportResult.FAILURE else [request_size] ) - assert received_body_sizes == expected_requests + assert [len(request.body) for request in received_requests] == expected_requests def test_oversized_batch_is_dropped_on_flush_without_blocking_later_batches( monkeypatch, caplog, otlp_http_server ): - base_url, received_body_sizes = otlp_http_server + base_url, received_requests = otlp_http_server oversized_spans = _finished_spans("secret-payload" * 1_000) small_spans = _finished_spans("ok") monkeypatch.setenv( @@ -173,7 +190,7 @@ def test_oversized_batch_is_dropped_on_flush_without_blocking_later_batches( super(LangfuseSpanProcessor, processor).on_end(oversized_spans[0]) assert processor.force_flush() - assert received_body_sizes == [] + assert received_requests == [] assert "Dropping span batch" in caplog.text assert "secret-payload" not in caplog.text @@ -182,7 +199,9 @@ def test_oversized_batch_is_dropped_on_flush_without_blocking_later_batches( finally: processor.shutdown() - assert received_body_sizes == [_serialized_request_size(small_spans)] + assert [len(request.body) for request in received_requests] == [ + _serialized_request_size(small_spans) + ] def test_default_exporter_uses_64_mib_limit_when_env_unset(monkeypatch): @@ -213,6 +232,129 @@ def test_invalid_max_batch_size_bytes_falls_back_to_default_limit( processor.shutdown() +def _decode_request(request: _RecordedRequest) -> ExportTraceServiceRequest: + body = gzip.decompress(request.body) if request.content_encoding else request.body + return ExportTraceServiceRequest.FromString(body) + + +def _span_names(request: ExportTraceServiceRequest) -> List[str]: + return [ + span.name + for resource_spans in request.resource_spans + for scope_spans in resource_spans.scope_spans + for span in scope_spans.spans + ] + + +@pytest.fixture +def compression_env(monkeypatch): + for name in ( + LANGFUSE_OTEL_COMPRESSION, + OTEL_EXPORTER_OTLP_TRACES_COMPRESSION, + OTEL_EXPORTER_OTLP_COMPRESSION, + ): + monkeypatch.delenv(name, raising=False) + + return monkeypatch + + +def _export_one_span(processor: LangfuseSpanProcessor) -> None: + provider = TracerProvider() + provider.add_span_processor(processor) + with provider.get_tracer("test").start_as_current_span( + "llm-call", attributes={"gen_ai.system": "test"} + ): + pass + + try: + assert processor.force_flush() + finally: + provider.shutdown() + + +def test_client_otel_compression_sends_gzip_request(compression_env, otlp_http_server): + base_url, received_requests = otlp_http_server + langfuse = Langfuse( + public_key="pk-test", + secret_key="sk-test", + base_url=base_url, + tracer_provider=TracerProvider(), + otel_compression="gzip", + ) + + with langfuse.start_as_current_observation(name="compressed-span"): + pass + langfuse.flush() + + assert [request.content_encoding for request in received_requests] == ["gzip"] + assert _span_names(_decode_request(received_requests[0])) == ["compressed-span"] + + +@pytest.mark.parametrize( + ("otel_compression", "env", "expected_encoding"), + [ + (None, {LANGFUSE_OTEL_COMPRESSION: " GZIP "}, "gzip"), + ("gzip", {LANGFUSE_OTEL_COMPRESSION: "none"}, "gzip"), + ("none", {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "gzip"}, None), + ( + None, + {LANGFUSE_OTEL_COMPRESSION: "none", OTEL_EXPORTER_OTLP_COMPRESSION: "gzip"}, + None, + ), + (None, {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "gzip"}, "gzip"), + ], +) +def test_default_exporter_compression_precedence( + compression_env, otlp_http_server, otel_compression, env, expected_encoding +): + base_url, received_requests = otlp_http_server + for name, value in env.items(): + compression_env.setenv(name, value) + + _export_one_span( + LangfuseSpanProcessor( + public_key="pk-test", + secret_key="sk-test", + base_url=base_url, + otel_compression=otel_compression, + ) + ) + + assert [request.content_encoding for request in received_requests] == [ + expected_encoding + ] + assert _span_names(_decode_request(received_requests[0])) == ["llm-call"] + + +@pytest.mark.parametrize( + ("otel_compression", "env_value", "setting"), + [ + ("deflate", None, "otel_compression"), + (None, "deflate", LANGFUSE_OTEL_COMPRESSION), + ], +) +def test_invalid_compression_warns_and_falls_back_to_otel_env( + compression_env, caplog, otlp_http_server, otel_compression, env_value, setting +): + base_url, received_requests = otlp_http_server + compression_env.setenv(OTEL_EXPORTER_OTLP_TRACES_COMPRESSION, "gzip") + if env_value is not None: + compression_env.setenv(LANGFUSE_OTEL_COMPRESSION, env_value) + + with caplog.at_level(logging.WARNING, logger="langfuse"): + processor = LangfuseSpanProcessor( + public_key="pk-test", + secret_key="sk-test", + base_url=base_url, + otel_compression=otel_compression, + ) + + _export_one_span(processor) + + assert f"Invalid {setting}='deflate'" in caplog.text + assert [request.content_encoding for request in received_requests] == ["gzip"] + + @pytest.fixture def tracer_with_processor(): processor = LangfuseSpanProcessor( From 4f183f1fac4c95ec6dbbb6fb0c0b3ecda6895b56 Mon Sep 17 00:00:00 2001 From: David Traina <44659830+DavidTraina@users.noreply.github.com> Date: Fri, 2 Oct 2026 16:17:11 -0500 Subject: [PATCH 2/3] feat(otel): warn when compression is set with a custom exporter --- langfuse/_client/span_processor.py | 5 +++++ tests/unit/test_span_processor.py | 14 ++++++++++++++ 2 files changed, 19 insertions(+) diff --git a/langfuse/_client/span_processor.py b/langfuse/_client/span_processor.py index 9c056490f..b362b8da1 100644 --- a/langfuse/_client/span_processor.py +++ b/langfuse/_client/span_processor.py @@ -179,6 +179,11 @@ def __init__( compression=_resolve_compression(otel_compression), max_request_size=_resolve_max_batch_size_bytes(), ) + elif otel_compression is not None: + langfuse_logger.warning( + "otel_compression is ignored because a custom span_exporter was " + "provided. Configure compression on the exporter instead." + ) if media_manager is not None or mask_otel_spans is not None: span_exporter = LangfuseTransformingSpanExporter( diff --git a/tests/unit/test_span_processor.py b/tests/unit/test_span_processor.py index 99d7b1c1e..959558e8c 100644 --- a/tests/unit/test_span_processor.py +++ b/tests/unit/test_span_processor.py @@ -355,6 +355,20 @@ def test_invalid_compression_warns_and_falls_back_to_otel_env( assert [request.content_encoding for request in received_requests] == ["gzip"] +def test_otel_compression_with_custom_span_exporter_warns(caplog): + with caplog.at_level(logging.WARNING, logger="langfuse"): + processor = LangfuseSpanProcessor( + public_key="pk-test", + secret_key="sk-test", + base_url="http://localhost:3000", + span_exporter=InMemorySpanExporter(), + otel_compression="gzip", + ) + processor.shutdown() + + assert "otel_compression is ignored" in caplog.text + + @pytest.fixture def tracer_with_processor(): processor = LangfuseSpanProcessor( From 8bfa89e7cfc8e0b256ab18c484a61ffb39842499 Mon Sep 17 00:00:00 2001 From: David Traina <44659830+DavidTraina@users.noreply.github.com> Date: Fri, 2 Oct 2026 16:19:39 -0500 Subject: [PATCH 3/3] test(otel): drop low-value compression tests --- tests/unit/test_span_processor.py | 48 ------------------------------- 1 file changed, 48 deletions(-) diff --git a/tests/unit/test_span_processor.py b/tests/unit/test_span_processor.py index 959558e8c..9591929be 100644 --- a/tests/unit/test_span_processor.py +++ b/tests/unit/test_span_processor.py @@ -296,11 +296,6 @@ def test_client_otel_compression_sends_gzip_request(compression_env, otlp_http_s (None, {LANGFUSE_OTEL_COMPRESSION: " GZIP "}, "gzip"), ("gzip", {LANGFUSE_OTEL_COMPRESSION: "none"}, "gzip"), ("none", {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "gzip"}, None), - ( - None, - {LANGFUSE_OTEL_COMPRESSION: "none", OTEL_EXPORTER_OTLP_COMPRESSION: "gzip"}, - None, - ), (None, {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "gzip"}, "gzip"), ], ) @@ -326,49 +321,6 @@ def test_default_exporter_compression_precedence( assert _span_names(_decode_request(received_requests[0])) == ["llm-call"] -@pytest.mark.parametrize( - ("otel_compression", "env_value", "setting"), - [ - ("deflate", None, "otel_compression"), - (None, "deflate", LANGFUSE_OTEL_COMPRESSION), - ], -) -def test_invalid_compression_warns_and_falls_back_to_otel_env( - compression_env, caplog, otlp_http_server, otel_compression, env_value, setting -): - base_url, received_requests = otlp_http_server - compression_env.setenv(OTEL_EXPORTER_OTLP_TRACES_COMPRESSION, "gzip") - if env_value is not None: - compression_env.setenv(LANGFUSE_OTEL_COMPRESSION, env_value) - - with caplog.at_level(logging.WARNING, logger="langfuse"): - processor = LangfuseSpanProcessor( - public_key="pk-test", - secret_key="sk-test", - base_url=base_url, - otel_compression=otel_compression, - ) - - _export_one_span(processor) - - assert f"Invalid {setting}='deflate'" in caplog.text - assert [request.content_encoding for request in received_requests] == ["gzip"] - - -def test_otel_compression_with_custom_span_exporter_warns(caplog): - with caplog.at_level(logging.WARNING, logger="langfuse"): - processor = LangfuseSpanProcessor( - public_key="pk-test", - secret_key="sk-test", - base_url="http://localhost:3000", - span_exporter=InMemorySpanExporter(), - otel_compression="gzip", - ) - processor.shutdown() - - assert "otel_compression is ignored" in caplog.text - - @pytest.fixture def tracer_with_processor(): processor = LangfuseSpanProcessor(