diff --git a/airflow-core/docs/administration-and-deployment/logging-monitoring/traces.rst b/airflow-core/docs/administration-and-deployment/logging-monitoring/traces.rst index 1c9258c6be24e..513b4dfa23a16 100644 --- a/airflow-core/docs/administration-and-deployment/logging-monitoring/traces.rst +++ b/airflow-core/docs/administration-and-deployment/logging-monitoring/traces.rst @@ -62,6 +62,24 @@ Add the following lines to your configuration file e.g. ``airflow.cfg`` See the OpenTelemetry `exporter protocol specification `_ and `SDK environment variable documentation `_ for more information. +Context propagation +------------------- + +Airflow propagates trace context across its own boundaries — Dag run to task, API server to worker, +supervisor to task runner — using the globally configured OpenTelemetry propagators rather than a +fixed W3C trace context propagator. Select them with the standard ``OTEL_PROPAGATORS`` environment +variable: + +.. code-block:: bash + + export OTEL_PROPAGATORS=tracecontext,baggage,my-vendor-propagator + +The default is ``tracecontext,baggage``, so ``baggage`` now propagates alongside the trace context, +and a vendor propagator registered under the ``opentelemetry_propagator`` entry point group is +honored wherever Airflow injects or extracts context. Carrier keys are whatever the configured +propagators write, so a deployment that keeps the default emits the same ``traceparent`` / +``tracestate`` pair as before. + Adding Custom Spans in Tasks ----------------------------- @@ -106,17 +124,28 @@ run — from the UI/API "Trigger with config" dialog, the ``airflow dags trigger ``airflow/dagrun_parent_trace_context`` Embed this run in an **external** trace instead of it being a root trace. Supply a W3C - ``traceparent`` string (optionally a mapping with ``traceparent`` and ``tracestate``) captured - from the system that triggered the run — an upstream orchestrator, event pipeline, CI job, or - another Airflow deployment. The whole run (the ``dag_run`` span and all task/worker spans) then - lives inside that trace, and parent-based samplers inherit the external sampling decision. When - the key is absent (default), the run is a root trace. A missing or malformed value is - ignored (the run stays a root) rather than failing run creation. + ``traceparent`` string, or a carrier mapping, captured from the system that triggered the run — + an upstream orchestrator, event pipeline, CI job, or another Airflow deployment. The whole run + (the ``dag_run`` span and all task/worker spans) then lives inside that trace, and parent-based + samplers inherit the external sampling decision. When the key is absent (default), the run is a + root trace. A missing or malformed value is ignored (the run stays a root) rather than failing + run creation. .. code-block:: json {"airflow/dagrun_parent_trace_context": "00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01"} + The two forms are read differently, on purpose. The bare-string shorthand is *defined* as a W3C + ``traceparent`` — its value syntax, not just its key name — so it is always parsed as one, + regardless of ``OTEL_PROPAGATORS``. A carrier mapping is read by the configured propagators + instead, so under the default ``OTEL_PROPAGATORS=tracecontext,baggage`` it accepts + ``traceparent``, ``tracestate`` and ``baggage``, and a deployment running a vendor propagator can + supply that propagator's own keys. Use the mapping form when you are not propagating W3C. + + .. code-block:: json + + {"airflow/dagrun_parent_trace_context": {"traceparent": "00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01", "baggage": "tenant=acme"}} + Enable Https ----------------- diff --git a/airflow-core/newsfragments/70840.improvement.rst b/airflow-core/newsfragments/70840.improvement.rst new file mode 100644 index 0000000000000..4ec6bb4402c78 --- /dev/null +++ b/airflow-core/newsfragments/70840.improvement.rst @@ -0,0 +1 @@ +Trace context is now propagated across Airflow's boundaries — Dag run to task, API server to worker, supervisor to task runner — using the globally configured OpenTelemetry propagators (``OTEL_PROPAGATORS``) instead of a hardcoded W3C trace context propagator. Deployments running a vendor propagator now have it honored wherever Airflow injects or extracts context, and ``baggage`` propagates alongside the trace context under the default ``tracecontext,baggage`` setting. Deployments keeping the default emit the same ``traceparent``/``tracestate`` carriers as before. diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py index a49d8fe30e546..2c33361c93f19 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py @@ -29,9 +29,8 @@ import structlog from cadwyn import VersionedAPIRouter from fastapi import Body, HTTPException, Query, Response, Security, status -from opentelemetry import trace +from opentelemetry import propagate, trace from opentelemetry.trace import StatusCode -from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator from pydantic import JsonValue, ValidationError from sqlalchemy import and_, func, or_, tuple_, update from sqlalchemy.engine import CursorResult @@ -572,7 +571,7 @@ def _emit_task_span(ti, state): return if not isinstance(ti.context_carrier, dict): return - dr_ctx = TraceContextTextMapPropagator().extract(ti.dag_run.context_carrier) + dr_ctx = propagate.extract(ti.dag_run.context_carrier) # Skip if the run was head-sampled out, so every span in the run agrees with the # carrier's decision. A parent-based sampler would already drop this child span, @@ -584,7 +583,7 @@ def _emit_task_span(ti, state): if dr_span_context.is_valid and not dr_span_context.trace_flags.sampled: return - ti_ctx = TraceContextTextMapPropagator().extract(ti.context_carrier) + ti_ctx = propagate.extract(ti.context_carrier) ti_span = trace.get_current_span(context=ti_ctx) span_context = ti_span.get_span_context() start_time_candidates = (x for x in (ti.queued_dttm, ti.start_date, timezone.utcnow()) if x) diff --git a/airflow-core/src/airflow/jobs/triggerer_job_runner.py b/airflow-core/src/airflow/jobs/triggerer_job_runner.py index c6cfe55f4b85c..7c2e5a24f4c40 100644 --- a/airflow-core/src/airflow/jobs/triggerer_job_runner.py +++ b/airflow-core/src/airflow/jobs/triggerer_job_runner.py @@ -39,9 +39,8 @@ import attrs import greenback import structlog -from opentelemetry import trace +from opentelemetry import propagate, trace from opentelemetry.trace import Status, StatusCode -from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator from pydantic import BaseModel, Field, TypeAdapter from sqlalchemy import func, select from structlog.contextvars import bind_contextvars as bind_log_contextvars @@ -157,9 +156,7 @@ def _make_trigger_span( ti: TaskInstanceDTO | None, trigger_id: int, name: str ) -> _AgnosticContextManager[trace.Span]: - parent_context = ( - TraceContextTextMapPropagator().extract(ti.context_carrier) if ti and ti.context_carrier else None - ) + parent_context = propagate.extract(ti.context_carrier) if ti and ti.context_carrier else None attributes: dict[str, str | int] = { "airflow.trigger.name": name, } diff --git a/airflow-core/src/airflow/models/dagrun.py b/airflow-core/src/airflow/models/dagrun.py index 725bf09c77d07..bad5af86628cc 100644 --- a/airflow-core/src/airflow/models/dagrun.py +++ b/airflow-core/src/airflow/models/dagrun.py @@ -28,7 +28,7 @@ from uuid import UUID import structlog -from opentelemetry import trace +from opentelemetry import propagate, trace from opentelemetry.context import context from opentelemetry.trace import StatusCode from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator @@ -210,25 +210,31 @@ def parent_trace_context(conf) -> context.Context | None: Lets a run be embedded in a trace owned by an external system (an upstream orchestrator, event pipeline, CI job, another Airflow) instead of being a root - trace. The value is a W3C ``traceparent`` string, or a carrier dict carrying - ``traceparent`` (and optionally ``tracestate``); anything else -- including a - malformed traceparent -- is ignored so it can neither silently mis-parent the - run nor fail run creation. Returns an OpenTelemetry ``Context`` when a valid - parent is present, else None (root trace, the default). + trace. The value is a W3C ``traceparent`` string, or a carrier dict read by the + configured propagators (``OTEL_PROPAGATORS``), so a deployment running a vendor + propagator can supply that propagator's own carrier keys. Anything else -- + including a malformed traceparent -- is ignored so it can neither silently + mis-parent the run nor fail run creation. Returns an OpenTelemetry ``Context`` + when a valid parent is present, else None (root trace, the default). """ if not conf: return None match conf.get(DAGRUN_PARENT_TRACE_CONTEXT_KEY): case str() as traceparent: + # The shorthand's *value* syntax is W3C, not just its key name, so it keeps + # being parsed as W3C whatever OTEL_PROPAGATORS says -- handing it to a + # propagator expecting another format would silently drop the parent. + extract = TraceContextTextMapPropagator().extract carrier = {"traceparent": traceparent} - case {"traceparent": str()} as raw: + case dict() as raw: # Keep only str members: a non-str tracestate reaches TraceState.from_header # unvalidated and raises TypeError, which would drop the otherwise-valid parent. - carrier = {k: raw[k] for k in ("traceparent", "tracestate") if isinstance(raw.get(k), str)} + extract = propagate.extract + carrier = {k: v for k, v in raw.items() if isinstance(v, str)} case _: return None try: - ctx = TraceContextTextMapPropagator().extract(carrier) + ctx = extract(carrier) except Exception: # Never let a malformed conf value fail run creation; fall back to a root trace. return None @@ -1196,7 +1202,7 @@ def _emit_dagrun_span(self, state: DagRunState): if not isinstance(self.context_carrier, dict): return - ctx = TraceContextTextMapPropagator().extract(self.context_carrier) + ctx = propagate.extract(self.context_carrier) span = trace.get_current_span(context=ctx) span_context = span.get_span_context() diff --git a/airflow-core/tests/unit/models/test_dagrun.py b/airflow-core/tests/unit/models/test_dagrun.py index 0a3250c3c62e9..3de8f2ea7789f 100644 --- a/airflow-core/tests/unit/models/test_dagrun.py +++ b/airflow-core/tests/unit/models/test_dagrun.py @@ -4302,15 +4302,82 @@ def test_parent_trace_context_tracestate(tracestate, expected): assert span_ctx.trace_state.get("foo") == expected -@mock.patch("airflow.models.dagrun.TraceContextTextMapPropagator.extract", side_effect=ValueError("boom")) -def test_parent_trace_context_swallows_propagator_error(mock_extract): - """A propagator failure degrades to a root trace instead of failing run creation.""" +def test_parent_trace_context_carrier_read_by_configured_propagators(): + """A carrier dict is handed to the configured propagators, not filtered down to W3C keys.""" + from opentelemetry import baggage + + from airflow.models.dagrun import parent_trace_context + + conf = { + DAGRUN_PARENT_TRACE_CONTEXT_KEY: { + "traceparent": f"00-{_EXTERNAL_TRACE_ID}-{_EXTERNAL_SPAN_ID}-01", + "baggage": "tenant=acme", + } + } + ctx = parent_trace_context(conf) + assert format(otel_trace.get_current_span(ctx).get_span_context().trace_id, "032x") == _EXTERNAL_TRACE_ID + assert baggage.get_all(ctx) == {"tenant": "acme"} + + +@pytest.fixture +def global_propagator_without_tracecontext(): + """Stand in for a deployment whose OTEL_PROPAGATORS does not include ``tracecontext``.""" + from opentelemetry import propagate + from opentelemetry.baggage.propagation import W3CBaggagePropagator + from opentelemetry.propagators.composite import CompositePropagator + + original = propagate.get_global_textmap() + propagate.set_global_textmap(CompositePropagator([W3CBaggagePropagator()])) + yield + propagate.set_global_textmap(original) + + +def test_parent_trace_context_string_shorthand_ignores_configured_propagators( + global_propagator_without_tracecontext, +): + """The bare-string shorthand is a W3C traceparent by definition, so OTEL_PROPAGATORS cannot break it.""" from airflow.models.dagrun import parent_trace_context conf = {DAGRUN_PARENT_TRACE_CONTEXT_KEY: f"00-{_EXTERNAL_TRACE_ID}-{_EXTERNAL_SPAN_ID}-01"} + span_ctx = otel_trace.get_current_span(parent_trace_context(conf)).get_span_context() + assert format(span_ctx.trace_id, "032x") == _EXTERNAL_TRACE_ID + + +def test_parent_trace_context_carrier_mapping_follows_configured_propagators( + global_propagator_without_tracecontext, +): + """The mapping form is the propagator-agnostic path, so a W3C-less propagator set drops it.""" + from airflow.models.dagrun import parent_trace_context + + conf = { + DAGRUN_PARENT_TRACE_CONTEXT_KEY: {"traceparent": f"00-{_EXTERNAL_TRACE_ID}-{_EXTERNAL_SPAN_ID}-01"} + } assert parent_trace_context(conf) is None +@pytest.mark.parametrize( + ("conf_value", "patch_target"), + [ + pytest.param( + f"00-{_EXTERNAL_TRACE_ID}-{_EXTERNAL_SPAN_ID}-01", + "airflow.models.dagrun.TraceContextTextMapPropagator.extract", + id="string-shorthand", + ), + pytest.param( + {"traceparent": f"00-{_EXTERNAL_TRACE_ID}-{_EXTERNAL_SPAN_ID}-01"}, + "airflow.models.dagrun.propagate.extract", + id="carrier-mapping", + ), + ], +) +def test_parent_trace_context_swallows_propagator_error(conf_value, patch_target): + """A propagator failure degrades to a root trace instead of failing run creation.""" + from airflow.models.dagrun import parent_trace_context + + with mock.patch(patch_target, side_effect=ValueError("boom")): + assert parent_trace_context({DAGRUN_PARENT_TRACE_CONTEXT_KEY: conf_value}) is None + + class TestDagRunTracing: """Tests for DagRun OpenTelemetry span behavior.""" diff --git a/shared/observability/src/airflow_shared/observability/traces/__init__.py b/shared/observability/src/airflow_shared/observability/traces/__init__.py index ea866997db902..b2891c8b501aa 100644 --- a/shared/observability/src/airflow_shared/observability/traces/__init__.py +++ b/shared/observability/src/airflow_shared/observability/traces/__init__.py @@ -23,14 +23,13 @@ from importlib.metadata import entry_points from typing import TYPE_CHECKING -from opentelemetry import context, trace +from opentelemetry import context, propagate, trace from opentelemetry.sdk.resources import Resource from opentelemetry.sdk.trace import TracerProvider from opentelemetry.sdk.trace.export import BatchSpanProcessor, SpanExporter from opentelemetry.sdk.trace.id_generator import RandomIdGenerator from opentelemetry.sdk.trace.sampling import Decision from opentelemetry.trace import NonRecordingSpan, Span, SpanContext, TraceFlags, TraceState -from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator if TYPE_CHECKING: from configparser import ConfigParser @@ -149,18 +148,16 @@ def new_dagrun_trace_carrier( ) ctx = trace.set_span_in_context(NonRecordingSpan(span_ctx)) carrier: dict[str, str] = {} - TraceContextTextMapPropagator().inject(carrier, context=ctx) + propagate.inject(carrier, context=ctx) return carrier def new_task_run_carrier(dag_run_context_carrier): - parent_context = ( - TraceContextTextMapPropagator().extract(dag_run_context_carrier) if dag_run_context_carrier else None - ) + parent_context = propagate.extract(dag_run_context_carrier) if dag_run_context_carrier else None span = tracer.start_span("notused", context=parent_context) # intentionally never closed new_ctx = trace.set_span_in_context(span) carrier: dict[str, str] = {} - TraceContextTextMapPropagator().inject(carrier, context=new_ctx) + propagate.inject(carrier, context=new_ctx) return carrier diff --git a/shared/observability/tests/observability/test_traces.py b/shared/observability/tests/observability/test_traces.py index eb0f499a250e1..1c20b5aa8d86e 100644 --- a/shared/observability/tests/observability/test_traces.py +++ b/shared/observability/tests/observability/test_traces.py @@ -18,7 +18,7 @@ from __future__ import annotations import pytest -from opentelemetry import context, trace +from opentelemetry import baggage, context, propagate, trace from opentelemetry.sdk.trace import TracerProvider from opentelemetry.sdk.trace.sampling import ( ALWAYS_OFF, @@ -35,6 +35,7 @@ build_trace_state_entries, get_task_span_detail_level, new_dagrun_trace_carrier, + new_task_run_carrier, ) @@ -360,3 +361,30 @@ def test_roundtrip_via_carrier(self): span = trace.get_current_span(ctx) assert get_task_span_detail_level(span) == 3 + + +class TestCarriersUseConfiguredPropagators: + """ + Carriers are written by the globally configured propagators, not a fixed W3C one. + + Under the default ``OTEL_PROPAGATORS=tracecontext,baggage`` that means baggage rides along, + which is what lets a vendor propagate implementation-specific details across the boundary. + """ + + @pytest.fixture + def baggage_attached(self): + token = context.attach(baggage.set_baggage("tenant", "acme")) + yield + context.detach(token) + + def test_dagrun_carrier_carries_baggage(self, baggage_attached): + carrier = new_dagrun_trace_carrier() + assert baggage.get_all(propagate.extract(carrier)) == {"tenant": "acme"} + + def test_task_run_carrier_carries_baggage(self, baggage_attached): + carrier = new_task_run_carrier(new_dagrun_trace_carrier()) + assert baggage.get_all(propagate.extract(carrier)) == {"tenant": "acme"} + + def test_no_baggage_leaves_carrier_untouched(self): + """Deployments not using baggage keep emitting exactly traceparent/tracestate.""" + assert set(new_dagrun_trace_carrier()) <= {"traceparent", "tracestate"} diff --git a/task-sdk/src/airflow/sdk/api/client.py b/task-sdk/src/airflow/sdk/api/client.py index ea220c68fde71..3589c40f7358c 100644 --- a/task-sdk/src/airflow/sdk/api/client.py +++ b/task-sdk/src/airflow/sdk/api/client.py @@ -31,8 +31,7 @@ import httpx import msgspec import structlog -from opentelemetry import trace -from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator +from opentelemetry import propagate, trace from pydantic import BaseModel, JsonValue from tenacity import ( before_log, @@ -177,7 +176,6 @@ def getuser() -> str: log = structlog.get_logger(logger_name=__name__) -_trace_propagator = TraceContextTextMapPropagator() _log_retry_warning = before_log(log, logging.WARNING) __all__ = [ @@ -230,7 +228,7 @@ def add_correlation_id(request: httpx.Request): def inject_trace_context(request: httpx.Request) -> None: - _trace_propagator.inject(request.headers) + propagate.inject(request.headers) def _log_and_trace_retry(retry_state) -> None: diff --git a/task-sdk/src/airflow/sdk/execution_time/comms.py b/task-sdk/src/airflow/sdk/execution_time/comms.py index 4a26f56e297dc..1aefabe021222 100644 --- a/task-sdk/src/airflow/sdk/execution_time/comms.py +++ b/task-sdk/src/airflow/sdk/execution_time/comms.py @@ -108,10 +108,7 @@ # Available on Unix and Windows (so "everywhere") but lets be safe recv_fds = None # type: ignore[assignment] -from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator - -_trace_propagator = TraceContextTextMapPropagator() - +from opentelemetry import propagate as otel_propagate if TYPE_CHECKING: from structlog.typing import FilteringBoundLogger as Logger @@ -194,7 +191,7 @@ class _RequestFrame(_FrameMixin, msgspec.Struct, array_like=True, frozen=True, o """ body: dict[str, Any] | None context_carrier: dict[str, str] | None = None - """W3C trace context carrier (traceparent + tracestate) of the task runner's active span. + """Trace context carrier of the task runner's active span, written by the configured propagators. The supervisor extracts this to restore the task runner's trace context before making outbound HTTP calls, so that server-side spans (e.g. POST /xcoms/…) appear as children of the correct task span @@ -237,7 +234,7 @@ class CommsDecoder(Generic[ReceiveMsgType, SendMsgType]): def _make_frame(self, msg: SendMsgType) -> _RequestFrame: carrier: dict[str, str] = {} - _trace_propagator.inject(carrier) + otel_propagate.inject(carrier) return _RequestFrame(id=next(self.id_counter), body=msg.model_dump(), context_carrier=carrier or None) @property diff --git a/task-sdk/src/airflow/sdk/execution_time/supervisor.py b/task-sdk/src/airflow/sdk/execution_time/supervisor.py index 0a4808512856f..641748d343ffb 100644 --- a/task-sdk/src/airflow/sdk/execution_time/supervisor.py +++ b/task-sdk/src/airflow/sdk/execution_time/supervisor.py @@ -161,10 +161,7 @@ except ImportError: send_fds = None # type: ignore[assignment] -from opentelemetry import context as otel_context, trace -from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator - -_trace_propagator = TraceContextTextMapPropagator() +from opentelemetry import context as otel_context, propagate as otel_propagate, trace if TYPE_CHECKING: from structlog.typing import FilteringBoundLogger, WrappedLogger @@ -957,7 +954,7 @@ def handle_requests(self, log: FilteringBoundLogger) -> Generator[None, _Request token = None try: if request.context_carrier: - ctx = _trace_propagator.extract(request.context_carrier) + ctx = otel_propagate.extract(request.context_carrier) token = otel_context.attach(ctx) self._handle_request(msg, log, request.id) except ServerResponseError as e: diff --git a/task-sdk/src/airflow/sdk/execution_time/task_runner.py b/task-sdk/src/airflow/sdk/execution_time/task_runner.py index eeb2788062248..434d5b7cf70cf 100644 --- a/task-sdk/src/airflow/sdk/execution_time/task_runner.py +++ b/task-sdk/src/airflow/sdk/execution_time/task_runner.py @@ -36,9 +36,8 @@ import attrs import lazy_object_proxy import structlog -from opentelemetry import trace +from opentelemetry import propagate, trace from opentelemetry.trace import INVALID_SPAN, Status, StatusCode -from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator from pydantic import AwareDatetime, ConfigDict, Field, JsonValue, TypeAdapter from structlog.contextvars import bind_contextvars @@ -203,9 +202,7 @@ def wrapper(*inner_args, **inner_kwargs): @contextmanager def _make_task_span(msg: StartupDetails): - parent_context = ( - TraceContextTextMapPropagator().extract(msg.ti.context_carrier) if msg.ti.context_carrier else None - ) + parent_context = propagate.extract(msg.ti.context_carrier) if msg.ti.context_carrier else None ti = msg.ti span_name = f"worker.{ti.task_id}" if ti.map_index is not None and ti.map_index >= 0: