Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,24 @@ Add the following lines to your configuration file e.g. ``airflow.cfg``
See the OpenTelemetry `exporter protocol specification <https://opentelemetry.io/docs/specs/otel/protocol/exporter/#configuration-options>`_ and
`SDK environment variable documentation <https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-variables/#periodic-exporting-metricreader>`_ 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
-----------------------------

Expand Down Expand Up @@ -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
-----------------

Expand Down
1 change: 1 addition & 0 deletions airflow-core/newsfragments/70840.improvement.rst
Original file line number Diff line number Diff line change
@@ -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.
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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)
Expand Down
7 changes: 2 additions & 5 deletions airflow-core/src/airflow/jobs/triggerer_job_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
}
Expand Down
26 changes: 16 additions & 10 deletions airflow-core/src/airflow/models/dagrun.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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()

Expand Down
73 changes: 70 additions & 3 deletions airflow-core/tests/unit/models/test_dagrun.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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


Expand Down
30 changes: 29 additions & 1 deletion shared/observability/tests/observability/test_traces.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -35,6 +35,7 @@
build_trace_state_entries,
get_task_span_detail_level,
new_dagrun_trace_carrier,
new_task_run_carrier,
)


Expand Down Expand Up @@ -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"}
6 changes: 2 additions & 4 deletions task-sdk/src/airflow/sdk/api/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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__ = [
Expand Down Expand Up @@ -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:
Expand Down
9 changes: 3 additions & 6 deletions task-sdk/src/airflow/sdk/execution_time/comms.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
Loading