Skip to main content
Pre-release

This release is a pre-release and may not be stable for production use.

Note
New integration maintainers are needed! Please open an issue to get started.

OpenLineage Dagster Integration

A library that integrates Dagster with OpenLineage for automatic metadata collection. It provides an OpenLineage sensor, a Dagster sensor that tails Dagster event logs for tracking metadata. On each sensor evaluation, the function processes a batch of event logs, converts Dagster events into OpenLineage events, and emits them to an OpenLineage backend.

Features

Metadata

  • Dagster job & op lifecycle

Requirements

Installation

$ python -m pip install openlineage-dagster

Usage

OpenLineage Sensor & Event Log Storage Requirements

Single OpenLineage sensor per Dagster instance
As the OpenLineage sensor processes all event logs for a given Dagster instance, define and enable only one sensor per instance. Running multiple sensors will result in duplicate OpenLineage job runs being emitted for Dagster steps with different OpenLineage run IDs. These are dynamically generated during sensor runs.

Non-sharded Event Log Storage
For the sensor to handle all event logs across runs, use non-sharded event log storage. If an event log storage sharded by run (i.e., the default SqliteEventLogStorage) is used, the cursor that tracks the last processed event storage ID may not update properly.

OpenLineage Sensor Setup

Get a OpenLineage sensor definition from the openlineage_sensor() factory function and add it to your Dagster repository.

from dagster import repository
from openlineage.dagster.sensor import openlineage_sensor


@repository
def my_repository():
    openlineage_sensor_def = openlineage_sensor()
    return other_defs + [openlineage_sensor_def]

Given that parallel sensor runs are not supported at the time of writing, some tuning may be necessary to avoid affecting other sensors' performance.

See Dagster's documentation on Evaluation Interval for more detail on minimum_interval_seconds, which defaults to 30 seconds. record_filter_limit is the maximum number of event logs to process on each sensor evaluation, and it defaults to 30 records per evaluation. Default values can be overridden:

@repository
def my_repository():
    openlineage_sensor_def = openlineage_sensor(
        minimum_interval_seconds=60,
        record_filter_limit=60,
    )
    return other_defs + [openlineage_sensor_def]

The OpenLineage sensor handles event logs in ascending order of storage ID and starts with the first log by default. Optionally, after_storage_id can be specified to customize the starting point. This is only applicable when the cursor is undefined or has been deleted.

@repository
def my_repository():
    openlineage_sensor_def = openlineage_sensor(
        after_storage_id=100
    )
    return other_defs + [openlineage_sensor_def]

OpenLineage Adapter & Client Configuration

The sensor uses an OpenLineage adapter and client to convert and push data to an OpenLineage backend. These depend on environment variables.

If using User Repository Deployments, add the below variables to the repository where the sensor is defined. Otherwise, add the variables to the Dagster Daemon.

  • OPENLINEAGE_URL - point to the service which will consume OpenLineage events.
  • OPENLINEAGE_API_KEY - set if the consumer of OpenLineage events requires a Bearer authentication key.
  • OPENLINEAGE_NAMESPACE - set if you are using something other than the default as the default namespace when a Dagster repository is undefined.

OpenLineage Namespace & Dagster Repository

For Dagster jobs organized in repositories, Dagster keeps track of the repository name for each pipeline run. When the repository name is present, it is always used as the OpenLineage namespace name. OPENLINEAGE_NAMESPACE option is a way to fall back and provide some other static default value.

Development

To install all dependencies for local development:

$ python -m pip install -e .[dev]  # or python -m pip install -e .\[dev\] in zsh 

To run the test suite:

$ pytest

SPDX-License-Identifier: Apache-2.0
Copyright 2018-2024 contributors to the OpenLineage project

Release files for openlineage-dagster 1.20.3.dev12141

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

Source distribution (sdist)

Source distribution for openlineage-dagster 1.20.3.dev12141
File Size Uploaded
openlineage-dagster-1.20.3.dev12141.tar.gz 13.0 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for openlineage-dagster 1.20.3.dev12141
File Interpreter ABI Platform
openlineage_dagster-1.20.3.dev12141-py3-none-any.whl Python 3 none any Details

Total release size: 22.1 kB

Release files / openlineage-dagster-1.20.3.dev12141.tar.gz

Download URL openlineage-dagster-1.20.3.dev12141.tar.gz
Size 13.0 kB
Tags Source
SHA-256 checksum
How to use checksums
71689eceddd35ee41f15dd9ffd4d611dec4eca8926a53af44707323fe1681d75
BLAKE2b-256 checksum
How to use checksums
b798bbcf1cb5d9da3334c94337297c8d979181a0819efea1c776d9829601a924
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/5.1.1 CPython/3.8.19

Release files / openlineage_dagster-1.20.3.dev12141-py3-none-any.whl

Download URL openlineage_dagster-1.20.3.dev12141-py3-none-any.whl
Size 9.2 kB
Tags Python 3
SHA-256 checksum
How to use checksums
d4f22b648d8e018c1bd59f2005f6caa945acfad5713a8aecb5bc19a4c55f20d8
BLAKE2b-256 checksum
How to use checksums
03ca379b719ae58ee1f75dbb271fdb172e786c067a515c5aac6bd75cecccaf61
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/5.1.1 CPython/3.8.19

Release history Release notifications | RSS feed

1.37.0

2 release files

1.36.0

2 release files

1.35.0

2 release files

1.34.0

2 release files

1.33.0

2 release files

1.32.0

2 release files

1.31.0

2 release files

1.30.1

2 release files

1.30.0

2 release files

1.29.0

2 release files

1.27.0

2 release files

1.26.0

2 release files

1.25.0

2 release files

1.21.1

2 release files

1.21.0

2 release files

1.20.6

2 release files

1.20.5

2 release files

1.20.4

2 release files

This release

1.19.0

2 release files

1.18.0

2 release files

1.17.1

2 release files

1.17.0

2 release files

1.16.0

2 release files

1.15.0

2 release files

1.13.1

2 release files

1.10.2

2 release files

1.10.1

2 release files

1.10.0

2 release files

1.9.1

2 release files

1.9.0

2 release files

1.8.0

2 release files

1.7.0

2 release files

1.6.2

2 release files

1.6.1

2 release files

1.6.0

2 release files

1.5.0

2 release files

1.4.1

2 release files

1.3.1

2 release files

1.3.0

2 release files

1.2.2

2 release files

1.2.1

2 release files

1.2.0

2 release files

1.1.0

2 release files

1.0.0

2 release files

0.30.1

2 release files

0.30.0

2 release files

0.29.2

2 release files

0.28.0

2 release files

0.26.0

2 release files

0.25.0

2 release files

0.23.0

2 release files

0.20.6

2 release files

0.17.0

2 release files

0.13.1

2 release files

0.13.0

2 release files

0.10.0

2 release files

0.9.0

2 release files

0.8.2

2 release files

0.8.1

2 release files

0.7.1

2 release files

0.7.0

2 release files

0.6.2

2 release files

0.6.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