Dagster: Modern Data Orchestration with Software-Defined Assets
Dagster is a data orchestrator that organizes pipelines around the data they produce rather than the tasks that run. This tutorial walks through the core concepts using a sales-data pipeline as a running example, from installation to deployment. By the end you will understand assets, resources, IO managers, partitions, schedules, sensors, and data-quality checks.
What Dagster Is
Dagster is an orchestration framework for building, running, and observing data pipelines. It was designed to make pipelines testable, observable, and maintainable as software. Instead of treating a pipeline as a graph of opaque tasks, Dagster encourages you to declare the assets your pipeline produces — tables, files, machine-learning models, dashboards — and lets the framework figure out execution order from the dependencies between them.
An asset is a persistent object in storage that captures some understanding of the world. A daily sales table in a warehouse is an asset. A cleaned Parquet file is an asset. A trained model is an asset. When you describe your pipeline in terms of these outputs, the orchestrator gains a richer model of your system, which improves observability and lineage.
How Dagster Differs from Airflow and Prefect
Airflow and Prefect are primarily task-centric. You define tasks (or flows) and wire up the order in which they execute. The orchestrator knows that task B runs after task A, but it does not inherently know that task A produces a rawsales table and task B consumes it. Lineage and data awareness must be added on top.
Dagster is asset-centric. You declare the assets, and the dependency graph between assets is derived from how they reference each other. This shift has practical consequences:
- Lineage is built in. The asset graph is the lineage graph. You can see which upstream asset feeds which downstream asset directly in the UI.
- Reasoning in terms of data. "Rematerialize the
dailysalessummaryasset" is a more natural operation than "rerun the transform task and then the load task." - Partitions are first-class. Each asset can be partitioned (for example, one partition per day), and Dagster tracks which partitions are materialized.
- Built-in data quality. Asset checks let you attach validations to assets and surface failures in the UI.
Dagster still supports task-style execution through ops and jobs (covered later), so you are not forced into the asset model for everything. But for most data pipelines, assets are the recommended starting point.
Installation
Dagster runs on Python 3.9 or later. Create a virtual environment and install the core package together with the web UI.
python -m venv .venv
source .venv/bin/activate
pip install dagster dagster-webserver
The dagster package provides the core library and command-line tools. The dagster-webserver package provides the local web UI you will use to inspect and run assets. Integration packages such as dagster-pandas, dagster-duckdb, and dagster-aws are installed separately as needed.
Verify the installation:
dagster --version
A clean project layout for the examples below looks like this:
salespipeline/
salespipeline/
init.py
assets.py
resources.py
pyproject.toml
The @asset Decorator and Software-Defined Assets
A software-defined asset is a Python function decorated with @asset. The function computes the asset's value, and its return value is handed to an IO manager for persistence. The function name becomes the asset's key.
from dagster import asset
import pandas as pd
@asset
def rawsales() -> pd.DataFrame:
"""Extract raw sales records from a source CSV."""
return pd.readcsv("https://example.com/sales.csv")
Dependencies between assets are expressed by adding a parameter whose name matches an upstream asset key. Dagster injects the upstream asset's value at runtime.
@asset
def cleanedsales(rawsales: pd.DataFrame) -> pd.DataFrame:
"""Transform: drop nulls and normalize column names."""
df = rawsales.dropna(subset=["orderid", "amount"])
df = df.rename(columns={"amt": "amount", "cust": "customerid"})
df["amount"] = df["amount"].astype(float)
return df
Here cleanedsales depends on rawsales simply because it declares a parameter with that name. Dagster reads this and builds the dependency edge automatically — there is no separate wiring step.
Building a Small ETL Pipeline as Assets
Let us build a coherent extract-transform-load pipeline for sales data. The three assets are rawsales (extract), cleanedsales (transform), and dailysalessummary (load/aggregate).
# salespipeline/assets.py
from dagster import asset
import pandas as pd
@asset
def raw
sales() -> pd.DataFrame:
"""Extract: read raw sales records."""
data = {
"orderid": [1, 2, 3, 4, None],
"customerid": ["c1", "c2", "c1", "c3", "c2"],
"amount": [120.0, 80.5, None, 45.0, 60.0],
"orderdate": [
"2026-05-01", "2026-05-01", "2026-05-02",
"2026-05-02", "2026-05-02",
],
}
return pd.DataFrame(data)
@asset
def cleanedsales(rawsales: pd.DataFrame) -> pd.DataFrame:
"""Transform: remove invalid rows and enforce types."""
df = rawsales.dropna(subset=["orderid", "amount"]).copy()
df["orderid"] = df["orderid"].astype(int)
df["amount"] = df["amount"].astype(float)
df["orderdate"] = pd.todatetime(df["orderdate"])
return df
@asset
def dailysalessummary(cleanedsales: pd.DataFrame) -> pd.DataFrame:
"""Load: aggregate revenue per day."""
summary = (
cleanedsales.groupby("orderdate")["amount"]
.agg(["sum", "count"])
.resetindex()
.rename(columns={"sum": "totalrevenue", "count": "ordercount"})
)
return summary
The dependency chain is rawsales -> cleanedsales -> dailysalessummary. Dagster infers it from the function parameters.
The Definitions Object
A Dagster project exposes its contents — assets, resources, schedules, sensors, and jobs — through a single Definitions object. This is the entry point Dagster loads.
# salespipeline/init.py
from dagster import Definitions, load
assetsfrommodules
from salespipeline import assets
allassets = loadassetsfrommodules([assets])
defs = Definitions(
assets=allassets,
)
loadassetsfrommodules scans a module and collects every @asset-decorated function, so you do not have to list them by hand. You can also pass an explicit list of assets if you prefer tighter control.
Running the UI with dagster dev
The local development server combines the web UI and the background daemon needed for schedules and sensors. Run it from the project root, pointing at the module that defines defs.
dagster dev -m salespipeline
By default the UI is served at http://localhost:3000. Open it and you will see the asset graph for the sales pipeline. From the Assets page you can click Materialize all to run the full chain, or select an individual asset to rematerialize only that one and its required upstreams.
Each materialization is recorded. The UI shows run history, logs, the lineage graph, and metadata attached to each asset.
Ops and Jobs (Briefly)
Before software-defined assets, Dagster's core abstractions were ops and jobs. An op is a unit of computation; a job is a graph of ops wired together. Ops remain useful for procedural work that does not produce a persistent asset — sending a notification, calling an external API for a side effect, or running an arbitrary script.
from dagster import op, job
@op
def fetchcount() -> int:
return 42
@op
def report(count: int) -> None:
print(f"Processed {count} records")
@job
def reportingjob():
report(fetchcount())
In modern projects you mostly work with assets, and Dagster can build an asset job that materializes a selection of assets. Reach for ops and jobs when the task genuinely has no asset output.
Resources and Configuration
Resources model connections and external systems — databases, APIs, cloud storage. Defining them separately keeps assets clean and makes them easy to swap in tests. A modern resource is a subclass of ConfigurableResource.
# salespipeline/resources.py
from dagster import ConfigurableResource
import duckdb
import pandas as pd
class DuckDBResource(ConfigurableResource):
databasepath: str
def query(self, sql: str) -> pd.DataFrame:
with duckdb.connect(self.databasepath) as conn:
return conn.execute(sql).fetchdf()
def writetable(self, table: str, df: pd.DataFrame) -> None:
with duckdb.connect(self.databasepath) as conn:
conn.register("tmpdf", df)
conn.execute(
f"CREATE OR REPLACE TABLE {table} AS SELECT FROM tmpdf"
)
An asset declares the resource it needs as a typed parameter, and Dagster injects the configured instance.
from dagster import asset
import pandas as pd
from salespipeline.resources import DuckDBResource
@asset
def dailysalessummary(
cleanedsales: pd.DataFrame, warehouse: DuckDBResource
) -> pd.DataFrame:
summary = (
cleanedsales.groupby("orderdate")["amount"]
.agg(["sum", "count"])
.resetindex()
.rename(columns={"sum": "totalrevenue", "count": "ordercount"})
)
warehouse.writetable("dailysalessummary", summary)
return summary
Wire the resource into Definitions and supply configuration. Reading the path from an environment variable keeps secrets and deployment-specific values out of code.
from dagster import Definitions, EnvVar, loadassetsfrommodules
from salespipeline import assets
from salespipeline.resources import DuckDBResource
defs = Definitions(
assets=loadassetsfrommodules([assets]),
resources={
"warehouse": DuckDBResource(databasepath=EnvVar("DUCKDBPATH")),
},
)
IO Managers
An IO manager controls how an asset's return value is stored and how it is loaded back when a downstream asset needs it. This separates what an asset computes from where its data lives. By default Dagster pickles values to the local filesystem, which is fine for development but rarely what you want in production.
You can attach an IO manager globally or per asset. A typical setup uses a warehouse-backed IO manager so that returning a DataFrame writes a table, and depending on that asset reads the table back.
from dagster import Definitions, FilesystemIOManager, loadassetsfrommodules
from salespipeline import assets
defs = Definitions(
assets=loadassetsfrommodules([assets]),
resources={
"iomanager": FilesystemIOManager(basedir="data/storage"),
},
)
The key idea: assets focus on transformation logic and return plain Python objects, while the IO manager owns persistence. Swapping storage (local files in dev, a warehouse in prod) becomes a configuration change rather than a code change. Integration packages provide ready-made IO managers, for example dagster-duckdb-pandas and dagster-snowflake-pandas.
Schedules and Sensors
Schedules and sensors decide when assets should be materialized.
Schedules
A schedule triggers a job on a cron-like cadence. Build a job from an asset selection, then attach a schedule.
from dagster import (
AssetSelection,
defineassetjob,
ScheduleDefinition,
)
salesjob = defineassetjob(
name="salesjob", selection=AssetSelection.all()
)
dailysalesschedule = ScheduleDefinition(
job=salesjob,
cronschedule="0 6 ", # every day at 06:00
)
Register both in Definitions:
defs = Definitions(
assets=allassets,
jobs=[salesjob],
schedules=[dailysalesschedule],
)
Sensors
A sensor polls for an external condition and launches a run when it is met — for example, a new file landing in a bucket. The sensor function yields a RunRequest when work should happen.
import os
from dagster import sensor, RunRequest, SkipReason
@sensor(job=salesjob)
def newfilesensor(context):
dropdir = "data/incoming"
files = os.listdir(dropdir) if os.path.isdir(dropdir) else []
if not files:
return SkipReason("No new files found")
for filename in files:
yield RunRequest(runkey=filename)
The runkey ensures each file triggers exactly one run, even if the sensor sees the same file on later ticks. Both schedules and sensors require the daemon, which dagster dev runs automatically.
Partitions
Partitions split an asset into independently materializable slices — most commonly one slice per day. This lets you backfill history, reprocess a single bad day, and track exactly which days have completed.
from dagster import asset, DailyPartitionsDefinition
import pandas as pd
dailypartitions = DailyPartitionsDefinition(startdate="2026-05-01")
@asset(partitionsdef=dailypartitions)
def partitionedsales(context) -> pd.DataFrame:
partitiondate = context.partitionkey
# Load only the data for this specific day.
df = readsalesfordate(partitiondate)
context.log.info(f"Loaded {len(df)} rows for {partitiondate}")
return df
context.partitionkey is the date string for the partition being materialized. Downstream partitioned assets map their partitions onto the upstream ones, so dailysalessummary for 2026-05-02 depends on partitionedsales for 2026-05-02. In the UI you get a partition grid showing which days are materialized, missing, or failed, and you can launch a backfill across a range of partitions in one action.
Asset Checks for Data Quality
Asset checks attach validations to an asset and report pass/fail status in the UI. They turn data-quality expectations into first-class, observable objects instead of scattered assertions.
from dagster import assetcheck, AssetCheckResult
import pandas as pd
@asset
check(asset="cleanedsales")
def no
negativeamounts(cleanedsales: pd.DataFrame) -> AssetCheckResult:
badrows = (cleanedsales["amount"] < 0).sum()
return AssetCheckResult(
passed=bool(badrows == 0),
metadata={"negativerows": int(badrows)},
)
@assetcheck(asset="cleanedsales")
def orderidisunique(cleanedsales: pd.DataFrame) -> AssetCheckResult:
duplicates = int(cleanedsales["orderid"].duplicated().sum())
return AssetCheckResult(
passed=duplicates == 0,
metadata={"duplicateorderids": duplicates},
)
Register checks alongside assets:
from dagster import Definitions, loadassetchecksfrommodules
from sales
pipeline import assets, checks
defs = Definitions(
assets=loadassetsfrommodules([assets]),
assetchecks=loadassetchecksfrommodules([checks]),
)
Checks run after the asset materializes (by default). A failing check is visible in the UI next to the asset, and you can configure runs to block downstream materialization when a check fails.
Testing Assets in Python
Because assets are plain Python functions, you can test them directly by calling them with ordinary inputs — no orchestrator required. This is one of the main practical benefits of the asset model.
# tests/testassets.py
import pandas as pd
from sales
pipeline.assets import cleanedsales, dailysalessummary
def test
cleanedsalesdropsnulls():
raw = pd.DataFrame(
{
"order
id": [1, None],
"customerid": ["c1", "c2"],
"amount": [10.0, None],
"orderdate": ["2026-05-01", "2026-05-01"],
}
)
result = cleanedsales(raw)
assert len(result) == 1
assert result["orderid"].iloc[0] == 1
def testdailysummaryaggregates():
cleaned = pd.DataFrame(
{
"orderid": [1, 2],
"customerid": ["c1", "c2"],
"amount": [10.0, 20.0],
"orderdate": pd.todatetime(["2026-05-01", "2026-05-01"]),
}
)
summary = dailysalessummary(cleaned)
assert summary["totalrevenue"].iloc[0] == 30.0
assert summary["ordercount"].iloc[0] == 2
For assets that need resources, Dagster provides materializetomemory, which executes a selection of assets in-process with test resources supplied directly.
from dagster import materializetomemory
from sales
pipeline.assets import rawsales, cleanedsales
def testpipelineinmemory():
result = materializetomemory([rawsales, cleanedsales])
assert result.success
Run the suite with pytest:
pytest tests/
Deployment Notes
A production Dagster deployment has two long-running processes:
dagster-webserverserves the UI and the GraphQL API that clients use to launch and inspect runs.dagster-daemonruns schedules, sensors, run queuing, and backfills. Without the daemon, time-based and event-based triggers will not fire.
dagster-webserver -h 0.0.0.0 -p 3000 -m salespipeline
dagster-daemon run -m salespipeline
Both processes read the same DAGSTERHOME directory, which holds instance configuration (dagster.yaml) and run/event storage. In dagster.yaml you configure where runs execute (for example, a Kubernetes run launcher), where logs and run metadata are stored (often Postgres), and concurrency limits.
For teams that prefer a managed option, Dagster+ (Dagster Cloud) hosts the control plane — the web UI, metadata storage, schedules, and daemon — while your code runs in your own environment through agents. It adds features such as branch deployments for testing pull requests, role-based access control, and alerting. Self-hosting with the open-source webserver and daemon remains fully supported; Dagster+ is an operational convenience rather than a requirement.
A common containerized deployment pattern:
# docker-compose.yaml (simplified)
services:
webserver:
image: my-org/sales-pipeline:latest
command: dagster-webserver -h 0.0.0.0 -p 3000 -m salespipeline
ports:
- "3000:3000"
environment:
DAGSTERHOME: /opt/dagster/home
DUCKDBPATH: /data/warehouse.duckdb
daemon:
image: my-org/sales-pipeline:latest
command: dagster-daemon run -m salespipeline
environment:
DAGSTERHOME: /opt/dagster/home
DUCKDBPATH: /data/warehouse.duckdb
Best Practices
- Model the data, not the tasks. Name assets after the tables and files they produce. The asset graph should read like a description of your data warehouse.
- Keep transformation logic pure. Let assets return plain objects and let IO managers handle persistence, so functions stay easy to test.
- Put connections in resources. Never hard-code credentials or connection strings in assets. Use
ConfigurableResourcewithEnvVarfor secrets. - Partition time-series assets. Daily (or hourly) partitions make backfills and reprocessing cheap and precise.
- Encode quality expectations as asset checks. A check that lives next to its asset is far more durable than an ad-hoc assertion buried in transform code.
- Write unit tests for asset functions. Call them directly with sample inputs; reserve
materializetomemoryfor resource-dependent flows. - Separate concerns by module. Keep assets, resources, schedules, and checks in their own files and assemble them in one
Definitions. - Run the daemon in every environment that uses schedules or sensors. A missing daemon is the most common reason triggers silently fail to fire.
Conclusion and Key Takeaways
Dagster's asset-centric model reframes orchestration around the data your pipelines produce. Compared with task-centric tools, this gives you lineage, observability, and partition tracking without bolting them on afterward.
Key takeaways:
- A software-defined asset is a decorated Python function; dependencies come from parameter names.
- The
Definitionsobject is the single entry point that ties together assets, resources, schedules, sensors, and checks. - Resources and IO managers separate business logic from connections and storage, which makes pipelines testable and portable.
- Partitions, schedules, and sensors control when and over what slice assets are materialized.
- Asset checks make data quality a visible, first-class concern.
- Assets are plain functions, so you can test most of your pipeline with
pytestand no orchestrator. - In production, run
dagster-webserveranddagster-daemon; Dagster+ offers a managed control plane if you prefer not to operate them yourself.
Start small with one or two assets, run dagster dev, and grow the graph from there. The asset model scales smoothly from a local prototype to a production warehouse pipeline.