From b2b2a8054fb73e530f8127b9fb714973354ebc6d Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Wed, 7 Oct 2026 21:01:16 +0000 Subject: [PATCH] feat(otel)!: compress span exports with gzip by default The default OTLP exporter now gzips span batches unless configured otherwise. Resolution order: the otel_compression argument, LANGFUSE_OTEL_COMPRESSION, OTEL_EXPORTER_OTLP_TRACES_COMPRESSION, OTEL_EXPORTER_OTLP_COMPRESSION, then gzip. Values are case-insensitive and an invalid value falls through to the next setting, matching the JS SDK. Set any of them to "none" to send uncompressed. Co-authored-by: Hassieb Pakzad --- langfuse/_client/client.py | 2 +- langfuse/_client/environment_variables.py | 10 ++--- langfuse/_client/span_processor.py | 53 ++++++++++++++++++----- tests/unit/test_span_processor.py | 51 ++++++++++++++++++++-- 4 files changed, 94 insertions(+), 22 deletions(-) diff --git a/langfuse/_client/client.py b/langfuse/_client/client.py index fe3b09426..a588b0e99 100644 --- a/langfuse/_client/client.py +++ b/langfuse/_client/client.py @@ -251,7 +251,7 @@ def mask_otel_spans( 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, `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. + otel_compression (Optional[Literal["gzip", "none"]]): Compression for span batches sent by the default OTLP span exporter: "gzip" (default) or "none" to send uncompressed. If unset, the LANGFUSE_OTEL_COMPRESSION, OTEL_EXPORTER_OTLP_TRACES_COMPRESSION and OTEL_EXPORTER_OTLP_COMPRESSION environment variables apply, in that order, then gzip. Values are case-insensitive; an invalid value falls through to the next setting. Example: ```python diff --git a/langfuse/_client/environment_variables.py b/langfuse/_client/environment_variables.py index b340c9ac4..3d2c84a92 100644 --- a/langfuse/_client/environment_variables.py +++ b/langfuse/_client/environment_variables.py @@ -77,13 +77,13 @@ """ .. 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. +Compression for span batches sent by the default OTLP exporter: ``gzip`` or ``none`` +(case-insensitive). The ``otel_compression`` client argument takes precedence. If this +variable is unset or invalid, ``OTEL_EXPORTER_OTLP_TRACES_COMPRESSION`` and then +``OTEL_EXPORTER_OTLP_COMPRESSION`` apply, and otherwise span exports are gzip-compressed. Custom span exporters are not affected. -**Default value:** unset +**Default value:** unset (span exports use gzip) """ LANGFUSE_DEBUG = "LANGFUSE_DEBUG" diff --git a/langfuse/_client/span_processor.py b/langfuse/_client/span_processor.py index ba2f5b2de..718d854e7 100644 --- a/langfuse/_client/span_processor.py +++ b/langfuse/_client/span_processor.py @@ -21,6 +21,10 @@ 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.environment_variables import ( + OTEL_EXPORTER_OTLP_COMPRESSION, + OTEL_EXPORTER_OTLP_TRACES_COMPRESSION, +) from opentelemetry.sdk.trace import ReadableSpan, Span from opentelemetry.sdk.trace.export import BatchSpanProcessor, SpanExporter from opentelemetry.trace import format_span_id, format_trace_id @@ -71,23 +75,18 @@ def _resolve_max_batch_size_bytes() -> Optional[int]: _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() +def _parse_compression( + raw_value: Optional[str], *, setting: str, warn: bool +) -> Optional[Compression]: + value = (raw_value or "").strip().lower() if not value: return None compression = _COMPRESSION_BY_NAME.get(value) - if compression is None: + if compression is None and warn: langfuse_logger.warning( - "Invalid %s=%r. Expected 'gzip' or 'none'. Falling back to the " - "OTEL_EXPORTER_OTLP_*COMPRESSION environment variables.", + "Invalid %s=%r. Expected 'gzip' or 'none'. Falling back to the next " + "compression setting, then gzip.", setting, raw_value, ) @@ -95,6 +94,36 @@ def _resolve_compression(otel_compression: Optional[str]) -> Optional[Compressio return compression +def _resolve_compression(otel_compression: Optional[str]) -> Compression: + """Resolve compression for the default exporter; gzip unless configured otherwise. + + Order: the ``otel_compression`` argument, ``LANGFUSE_OTEL_COMPRESSION``, + ``OTEL_EXPORTER_OTLP_TRACES_COMPRESSION``, ``OTEL_EXPORTER_OTLP_COMPRESSION``, + then gzip. Values are case-insensitive, and an invalid value falls through to + the next setting, matching the JS SDK. + """ + candidates = ( + ("otel_compression", otel_compression, True), + (LANGFUSE_OTEL_COMPRESSION, os.environ.get(LANGFUSE_OTEL_COMPRESSION), True), + ( + OTEL_EXPORTER_OTLP_TRACES_COMPRESSION, + os.environ.get(OTEL_EXPORTER_OTLP_TRACES_COMPRESSION), + False, + ), + ( + OTEL_EXPORTER_OTLP_COMPRESSION, + os.environ.get(OTEL_EXPORTER_OTLP_COMPRESSION), + False, + ), + ) + for setting, raw_value, warn in candidates: + compression = _parse_compression(raw_value, setting=setting, warn=warn) + if compression is not None: + return compression + + return Compression.Gzip + + class LangfuseSpanProcessor(BatchSpanProcessor): """OpenTelemetry span processor that exports spans to the Langfuse API. diff --git a/tests/unit/test_span_processor.py b/tests/unit/test_span_processor.py index 9591929be..d7c5366be 100644 --- a/tests/unit/test_span_processor.py +++ b/tests/unit/test_span_processor.py @@ -134,6 +134,11 @@ def _serialized_request_size(spans: List[ReadableSpan]) -> int: return len(encode_spans(spans).SerializePartialToString()) +def _uncompressed_size(request: _RecordedRequest) -> int: + body = gzip.decompress(request.body) if request.content_encoding else request.body + return len(body) + + def _default_exporter_processor(base_url: str) -> LangfuseSpanProcessor: return LangfuseSpanProcessor( public_key="pk-test", @@ -170,7 +175,9 @@ def test_default_exporter_enforces_max_batch_size_bytes_at_boundary( expected_requests = ( [] if expected_result == SpanExportResult.FAILURE else [request_size] ) - assert [len(request.body) for request in received_requests] == expected_requests + assert [_uncompressed_size(request) for request in received_requests] == ( + expected_requests + ) def test_oversized_batch_is_dropped_on_flush_without_blocking_later_batches( @@ -199,7 +206,7 @@ def test_oversized_batch_is_dropped_on_flush_without_blocking_later_batches( finally: processor.shutdown() - assert [len(request.body) for request in received_requests] == [ + assert [_uncompressed_size(request) for request in received_requests] == [ _serialized_request_size(small_spans) ] @@ -293,10 +300,46 @@ def test_client_otel_compression_sends_gzip_request(compression_env, otlp_http_s @pytest.mark.parametrize( ("otel_compression", "env", "expected_encoding"), [ - (None, {LANGFUSE_OTEL_COMPRESSION: " GZIP "}, "gzip"), + # Nothing configured: gzip by default. + (None, {}, "gzip"), + # The argument wins over every environment variable. ("gzip", {LANGFUSE_OTEL_COMPRESSION: "none"}, "gzip"), ("none", {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "gzip"}, None), - (None, {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "gzip"}, "gzip"), + # LANGFUSE_OTEL_COMPRESSION wins over the OTEL variables (case-insensitive). + (None, {LANGFUSE_OTEL_COMPRESSION: " GZIP "}, "gzip"), + ( + None, + {LANGFUSE_OTEL_COMPRESSION: "None", OTEL_EXPORTER_OTLP_COMPRESSION: "gzip"}, + None, + ), + # The traces-specific OTEL variable wins over the generic one. + ( + None, + { + OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "gzip", + OTEL_EXPORTER_OTLP_COMPRESSION: "none", + }, + "gzip", + ), + (None, {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "NONE"}, None), + (None, {OTEL_EXPORTER_OTLP_COMPRESSION: "none"}, None), + # Invalid values fall through to the next setting, then gzip. + ("brotli", {LANGFUSE_OTEL_COMPRESSION: "none"}, None), + ("brotli", {}, "gzip"), + ( + None, + {LANGFUSE_OTEL_COMPRESSION: "zstd", OTEL_EXPORTER_OTLP_COMPRESSION: "none"}, + None, + ), + ( + None, + { + OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "deflate", + OTEL_EXPORTER_OTLP_COMPRESSION: "none", + }, + None, + ), + (None, {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "bogus"}, "gzip"), ], ) def test_default_exporter_compression_precedence(