dagster-openlineage
Asset-centric OpenLineage emission for Dagster. Emits schema, column-lineage, data-quality-assertion, and partition-nominal-time facets alongside existing run/step events. Two emission mechanisms are available — pick exactly one per deployment.
Features
- Asset-centric emission — materializations, observations, check evaluations, and synthesized failures
- Schema + column lineage facets from
dagster/column_schemaanddagster/column_lineagemetadata - Data quality assertions placed on
InputDataset(spec-conformant) - Partition → nominal time heuristic (ISO date or date-hour)
- Multi-tenant namespaces via string templates (
{namespace},{tag:KEY}) - Bounded emit — synchronous, 2s default timeout, retries disabled, failures swallowed
- Pipeline / step events preserved (v0.1 surface unchanged)
Requirements
- Python 3.10+
- Dagster
>=1.11.6
Installation
pip install dagster-openlineage
Compatibility matrix
| Environment | Mechanism A (storage wrapper) | Mechanism B (sensor) |
|---|---|---|
| OSS Dagster (self-hosted) | ✅ | ✅ |
OSS Dagster, dagster dev |
✅ | ✅ |
| Dagster+ Hybrid | ❌ (operator doesn't control instance.yaml) |
✅ |
| Dagster+ Serverless | ❌ | ✅ |
| Dagster+ Branch Deployments | ❌ | ✅ |
Which mechanism do I want?
- You control
instance.yaml→ Mechanism A (push). Every event hits OpenLineage as it is persisted. No daemon dependency. Process-local failure synthesis (see Known limitations). - You run on Dagster+ Hybrid / Serverless / Branch Deployments → Mechanism B (polled). The sensor tails the event log and converts asset events into OpenLineage emissions.
- Both at once → don't. See Pick exactly one mechanism below.
Mechanism A — storage wrapper (push)
Configure OpenLineageEventLogStorage as your event log storage. It composes any inner EventLogStorage class and intercepts store_event / store_event_batch to emit OpenLineage events.
# instance.yaml
event_log_storage:
module: dagster_openlineage
class: OpenLineageEventLogStorage
config:
wrapped:
module: dagster_postgres.event_log
class: PostgresEventLogStorage
config:
postgres_url:
env: DAGSTER_PG_URL
namespace: my-company
# Optional:
# namespace_template: "{namespace}/{tag:tenant}"
# timeout: 2.0
# strict_assertion_mapping: false
Set OPENLINEAGE_URL (and optionally OPENLINEAGE_API_KEY) in the environment of any process that writes Dagster events — typically the run worker and the daemon.
Mechanism B — sensor (polled)
Add openlineage_sensor(include_asset_events=True) to your Definitions. v0.2 keeps include_asset_events=False as the default (v0.1 parity); v0.3 will flip it.
from dagster import Definitions
from dagster_openlineage import openlineage_sensor
defs = Definitions(
assets=[...],
sensors=[openlineage_sensor(include_asset_events=True)],
)
Environment variables go on the process that runs the Dagster daemon:
OPENLINEAGE_URL(required)OPENLINEAGE_API_KEY(optional)OPENLINEAGE_NAMESPACE(optional, defaultdagster)
Namespace templates
Provide namespace_template (wrapper) or let the adapter/sensor derive the default. v0.2 supports two token families:
{namespace}— the configured default namespace{tag:KEY}— the run tag namedKEY, empty if unset
Adjacent slashes collapse; trailing slashes strip. Unknown tokens raise NamespaceTemplateError at construction. {code_location} and {repository} are deferred to v0.3 — they do not reliably reach store_event time.
Example:
# Template
"{namespace}/{tag:tenant}"
# Run tags -> resolved namespace
{"tenant": "acme"} -> "dagster/acme"
{} -> "dagster" # tag unresolved, trailing slash stripped
Emit path
The adapter emits synchronously and swallows Exception; BaseException (shutdown signals, SystemExit, KeyboardInterrupt, MemoryError) propagates. The default transport disables retries to keep the per-event wall clock bounded by timeout (default 2s). Configure retry/async behavior via OpenLineage's own openlineage.yml or OPENLINEAGE_CONFIG if you need it.
Pick exactly one mechanism
Configuring Mechanism A and Mechanism B simultaneously produces duplicate OpenLineage events. OpenLineage has no client-side idempotency primitive and backend dedup is not spec-defined (Marquez dedupes on (runId, eventType, eventTime); DataHub and OpenMetadata do not). This is a deployment-time contract — the library does not enforce it at runtime. Mechanism B logs a WARN at sensor construction when opted in, reminding operators of this.
Migration from v0.1
- Pin
dagster>=1.11.6. - Namespace default is flat (
dagster). v0.1 usedrepository_nameas the namespace when reachable — existing OL-backend lineage keyed under that old namespace may need a one-time rename. SetOPENLINEAGE_NAMESPACEor configurenamespace_templateto match your prior layout. - If you relied on the
openlineage_sensorin v0.1, it still works unchanged — asset events are now supported via the newinclude_asset_events=Trueflag (default remainsFalse). OpenLineageEventListeneris removed; it was a dead stub with no call sites.
Sensor Configuration
| Option | Default | Description |
|---|---|---|
minimum_interval_seconds |
300 | Minimum seconds between sensor evaluations |
record_filter_limit |
30 | Max number of event logs to process per evaluation |
after_storage_id |
0 | Starting storage ID for event processing |
include_asset_events |
False (v0.2) |
Opt in to asset-level emission |
Not in v0.2
- Dagster+ Insights bridge / Catalog UI mirroring
- IO manager enrichment (
LOADED_INPUT/HANDLED_OUTPUT) - Spark / engine-specific facets
- SQL parsing or column-graph inference (reads
dagster/column_lineagemetadata if present) - Python-callable
naming_fn(string templating only; callable form deferred to v0.3) {code_location}and{repository}namespace tokens (deferred to v0.3)- Shipped JSON schema file for the custom
dagster_asset_checkrun facet (deferred to v0.2.1) - Wrapper-side synthesis reconciliation across process restarts (deferred to v0.2.1; use Mechanism B if you need crash-tolerant failure reporting)
Version Compatibility
This library supports Dagster >=1.11.6. End users get the constraint from pyproject.toml; CI and reproducible builds use the committed uv.lock.
Development
uv sync
make test
make ruff
make check
Release files for dagster-openlineage 0.2.1
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| dagster_openlineage-0.2.1.tar.gz | 26.0 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| dagster_openlineage-0.2.1-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 54.1 kB
Release files / dagster_openlineage-0.2.1.tar.gz
| Download URL | dagster_openlineage-0.2.1.tar.gz |
|---|---|
| Size | 26.0 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
e1589b018eb859b34697da199e7850592e291a865e19024ee092529b2fbd8856
|
|
BLAKE2b-256 checksum How to use checksums |
2e334b0c0e58fcd8092de016d9d2bf67c852b2222924bdbd4ae12221c1583133
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
uv/0.11.16 {"installer":{"name":"uv","version":"0.11.16","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}
|
Release files / dagster_openlineage-0.2.1-py3-none-any.whl
| Download URL | dagster_openlineage-0.2.1-py3-none-any.whl |
|---|---|
| Size | 28.1 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
5b53ef44a0410c5699580475d69e81d9005250d1a6e62cd4d68e22c0f3c58703
|
|
BLAKE2b-256 checksum How to use checksums |
9109b55f471876e194df8773c3bea101df1b9c0c1b6a9c498558a00614aaec57
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
uv/0.11.16 {"installer":{"name":"uv","version":"0.11.16","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}
|