Skip to main content

dagster-otel

PyPI Python versions Release CI CodeQL k8s e2e License: MIT

OpenTelemetry tracing for Dagster ops and assets -- with trace/span IDs correlated into your own log lines -- without giving up ownership of your op/asset definitions to a third-party decorator, and without monkeypatching Dagster internals.

Need tracing without adding @traced() calls -- pipelines you don't own the source of? See opentelemetry-instrumentation-dagster, a separate, explicitly-riskier monkeypatch-based companion package built for that one use case. This project's own core stays decorator-based; see Why this exists for why.

A Jaeger trace showing jaffle_shop_dbt_assets nested into per-model and per-test spans A real trace from the included example -- traced_dbt() turns one opaque @dbt_assets step into a real step → asset → check tree.

Table of Contents

Installation

pip install dagster-otel

Usage

from dagster import asset, job, op

from dagster_otel import traced

@op(...)                # Dagster's own @op still owns op-ness; @traced() is a thin
@traced()                # layer underneath. No @resource/required_resource_keys, no
def upstream_op(context) -> int:  # manual "root" step -- the first @traced() step to
    ...                            # run in a given run just becomes the root.

@op(...)
@traced()
def downstream_op(context, x: int) -> int:
    ...

@asset(...)
@traced()  # same decorator, works for assets too
def downstream_asset(context) -> None:
    ...

@op(...)
@traced  # bare works too, like @op/@asset themselves -- same as @traced()
def another_op(context) -> None:
    ...

@job(...)
def my_job():
    downstream_op(upstream_op())

Works across Dagster's multiprocess and k8s_job_executor executors: each step usually runs in its own process, sometimes on its own node, so trace context is propagated via Dagster's own run storage rather than in-process memory. See docs/design.md for how, and what's verified vs. still assumed.

The actual trace shape for a multi-root, fan-in graph (root_a/root_b independent, merge_op depending on both) -- verified against real Dagster + Jaeger, not just drawn for illustration:

graph TD
    root_a[root_a] --> child_a[child_a]
    root_b[root_b] --> child_b[child_b]
    root_a --> merge_op[merge_op]
    root_b -.->|Link| merge_op

merge_op gets a real parent (root_a, deterministic) plus a Link to the other dependency it can't have as a second parent -- both relationships stay visible on the span, not just whichever upstream happened to be found first.

For @dbt_assets, dagster_otel.dbt.traced_dbt() is a drop-in replacement for @traced() that additionally opens a child span per dbt node (model/seed/test), keyed by the real Dagster asset_key/check_name -- no changes needed to the function body:

from dagster_otel.dbt import traced_dbt

@dbt_assets(manifest=...)
@traced_dbt()
def my_dbt_assets(context, dbt: DbtCliResource):
    yield from dbt.cli(["build"], context=context).stream()

In Jaeger, that's a real step → asset → check tree, not one opaque span for the whole dbt build -- see the screenshot at the top of this README (the real jaffle_shop example: 16 spans, 1 step + 3 assets + 12 checks), each with the accurate duration dbt itself measured.

For @sensor/@schedule tick evaluation -- what decides whether a run happens, before any run_id exists, so @traced() itself doesn't apply -- use traced_sensor()/traced_schedule():

from dagster_otel import traced_sensor, traced_schedule

@sensor(job=my_job)
@traced_sensor()
def my_sensor(context: SensorEvaluationContext):
    ...
    return RunRequest(...)

@schedule(cron_schedule="0 * * * *", job=my_job)
@traced_schedule()
def my_schedule(context: ScheduleEvaluationContext):
    ...
    return RunRequest(...)

Each tick gets its own span (a fresh root -- a tick has no run_id or upstream step to attach to, unlike @traced()). Any RunRequest the tick returns/yields gets tagged so the run it launches (if any) nests under that tick's span in the trace backend -- so "why did/didn't this run fire" is answerable from the trace directly. See docs/design.md for the full design and real-Jaeger verification.

@traced() doesn't belong on a @graph_asset's own decorated function -- that function is a definition-time composition of other ops (it wires up which @op depends on which, called once at definition time), not a per-run compute_fn, so it never receives a runtime context for traced() to use in the first place. Put @traced() on the individual @ops the graph_asset composes instead -- those run at execution time with a real context, same as any other op, and tracing them already covers everything that actually executes under the graph asset.

Configuration

Standard OTel environment variables -- nothing bespoke:

Variable Purpose
OTEL_SERVICE_NAME Names your service in the trace backend.
OTEL_EXPORTER_OTLP_ENDPOINT (or ..._TRACES_ENDPOINT) Where to send spans (e.g. http://localhost:4317). Required -- without one of these set, no real exporter is attached at all (see below).
OTEL_EXPORTER_OTLP_TRACES_TIMEOUT / ..._TIMEOUT Per-export timeout. Set this yourself if the default (2s) doesn't fit -- see docs/design.md for why a default exists at all (an unreachable collector otherwise blocked every step for ~7s).
OTEL_EXPORTER_OTLP_TRACES_PROTOCOL / ..._PROTOCOL Transport to export over: grpc (default) or http/protobuf -- e.g. for a collector that only exposes HTTP ingest, or an environment that blocks gRPC egress. Any other value raises rather than silently keeping gRPC.
OTEL_SDK_DISABLED Set to true to force no export regardless of the endpoint vars above.

@traced() reads these itself (idempotently) the first time it runs in a process -- there's nothing else to wire up, no @resource/required_resource_keys needed. Call configure() yourself only if you want configuration to happen eagerly (e.g. at Definitions load time) rather than lazily on first use.

Without OTEL_EXPORTER_OTLP_ENDPOINT/..._TRACES_ENDPOINT set, no real OTLP exporter is created at all -- spans are still created (propagation and log correlation keep working), just never sent anywhere, so trying @traced() with zero setup never makes a surprise network call. See docs/design.md for the one deliberate tradeoff this makes.

Nesting a whole run's trace under an external caller's (a CI/CD pipeline, a scheduler, another OTel-instrumented system) is a run tag, not an env var -- set EXTERNAL_TRACE_CONTEXT_TAG_KEY (exported from dagster_otel) at launch time:

from dagster_otel import EXTERNAL_TRACE_CONTEXT_TAG_KEY

carrier: dict[str, str] = {}
TraceContextTextMapPropagator().inject(carrier)  # from your own active span
my_job.execute_in_process(tags={EXTERNAL_TRACE_CONTEXT_TAG_KEY: json.dumps(carrier)})

Every root step in the run (the ones that would otherwise seed a fresh trace) checks for this tag first. See docs/design.md for the full verification.

Trace backends checked

No backend-specific code exists here -- configure() constructs OTLPSpanExporter() with no endpoint=/headers=/credentials=, so anything speaking OTLP should work purely via the env vars above. What's actually been checked end-to-end, not just assumed to work by construction:

Backend License Status
Jaeger Apache-2.0 ✅ Verified -- multiprocess/k8s_job_executor/retry-from-failure/traced_dbt()/external trace context, see docs/design.md
Grafana Tempo AGPL-3.0 ✅ Verified -- real @dbt_assets jaffle_shop pipeline, through a real OTel Collector (not sent directly), see docs/design.md
SigNoz MIT ⬜ Not yet -- #43

Should work the same way against any other OTLP-compatible backend (Honeycomb, Datadog, New Relic, ...) -- just not individually checked off here yet. Open an issue if you hit something backend-specific.

Combined demo: this project + dagster-prometheus-exporter

docker compose up -d also starts a full traces-and-metrics demo: this project's traces and dagster-prometheus-exporter's metrics, both flowing through the same real OTel Collector into Grafana (http://localhost:3002, both Prometheus and Tempo datasources provisioned, plus a pre-built "dagster-otel combined demo" dashboard) -- see examples/README.md for how to run the example pipeline against it. The exporter needs zero changes; it's referenced as an external published image, not vendored here. See docs/design.md for the full verification writeup.

Compatibility

Built and verified against Dagster 1.13.22 and opentelemetry-sdk 1.44.0 (pyproject.toml/uv.lock) -- this is the combination every behavior described here has actually been checked against, including the multiprocess/k8s_job_executor/ retry-from-failure verification in docs/design.md. pyproject.toml declares a much wider floor (dagster >= 1.5) since nothing here relies on version-specific Dagster internals beyond what's documented as an accepted-risk private-API dependency there -- but that wide range isn't individually spot-checked the way it is for dagster-prometheus-exporter. If you hit an incompatibility on another version, please open an issue.

Requires Python 3.10+ (matches Dagster's own floor).

Why this exists

See docs/design.md for the full rationale, including a comparison against prior art (formenergy-observability, a monkeypatch-based prototype) and the design decisions (no monkeypatching, decorators stack under Dagster's own @op/@asset rather than replacing it, log correlation via a public logging.Filter on context.log).

Contributing

See CONTRIBUTING.md for how to set up your toolchain, run checks locally, and submit a pull request. Bug reports and feature requests go through GitHub issues; a security vulnerability goes to SECURITY.md instead. See CHANGELOG.md for what changed in each release.

Roadmap

  • @traced() for @op/@asset, bare or parameterized, with automatic root detection (no manual publish_trace_context() call needed)
  • Cross-process propagation via Dagster run tags -- verified across multiprocess and k8s_job_executor, and across retry-from-failure
  • Log/trace correlation via a public logging.Filter on context.log
  • traced_dbt() for per-dbt-node (model/seed/test) spans
  • traced_sensor() / traced_schedule() for tick evaluation
  • External trace context nesting (EXTERNAL_TRACE_CONTEXT_TAG_KEY)
  • Deterministic multi-root trace_id, so independent roots in the same run never land in separate traces
  • Configurable OTLP transport (grpc / http/protobuf)
  • Jaeger and Grafana Tempo verified end-to-end
  • SigNoz verification (#43)
  • User-supplied callback for custom span attributes (#39)
  • dagster.partition_key span attribute for partitioned assets (#36)
  • Automated k8s_job_executor e2e CI job (#34) -- scheduled + on-demand, not per-PR (see the workflow's own comment for why)

License

MIT -- see LICENSE.

Metadata

Release files for dagster-otel 0.4.1

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for dagster-otel 0.4.1
File Size Uploaded
dagster_otel-0.4.1.tar.gz 759.3 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for dagster-otel 0.4.1
File Interpreter ABI Platform
dagster_otel-0.4.1-py3-none-any.whl Python 3 none any Details

Total release size: 799.7 kB

Release files / dagster_otel-0.4.1.tar.gz

Download URL dagster_otel-0.4.1.tar.gz
Size 759.3 kB
Tags Source
SHA-256 checksum
How to use checksums
c7e6ef6f88c3051f2b72b5bbf76fae1f9549de36614e04eca25ade39bdfb7f62
BLAKE2b-256 checksum
How to use checksums
20cb1ab4721df827d3d264adaf908ca7175ecb9c01d6447e2881363e5d6bda2d
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 23, 2026.

Transparency log

Release files / dagster_otel-0.4.1-py3-none-any.whl

Download URL dagster_otel-0.4.1-py3-none-any.whl
Size 40.4 kB
Tags Python 3
SHA-256 checksum
How to use checksums
081c48631da2670a65af83308cbd3a1b68247893c298254c04548e6497b06e56
BLAKE2b-256 checksum
How to use checksums
7308350570379c5c249e982f23bd6fb14b68c6f393b380ceaf0909eace76b2ad
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 23, 2026.

Transparency log

Release history Release notifications | RSS feed

0.5.1

2 release files

0.5.0

2 release files

This release

0.4.1 This release

2 release files

0.4.0

2 release files

0.3.2

2 release files

0.3.1

2 release files

0.3.0

2 release files

0.2.1

2 release files

0.2.0

2 release files

0.1.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page