Dagster
Dagster is an asset-oriented orchestration platform for data engineering. Instead of focusing on task execution like Airflow or Prefect, Dagster centers on software-defined assets: data artifacts that are produced and consumed by your pipelines.
Core Concepts
Software-Defined Assets (SDAs)
An asset is a data object (table, file, ML model) that your pipeline produces. Assets are defined in Python with metadata about their dependencies, upstream sources, and downstream consumers. Dagster automatically tracks the asset graph and knows what to recompute when upstream data changes.
from dagster import asset, AssetIn, Output
import pandas as pd
@asset
def raw_orders() -> pd.DataFrame:
"""Extract raw orders from the source database."""
return pd.read_sql("SELECT * FROM orders", ...)
@asset
def cleaned_orders(raw_orders: pd.DataFrame) -> pd.DataFrame:
"""Clean and validate orders data."""
df = raw_orders.dropna(subset=["order_id", "customer_id"])
df["total"] = df["quantity"] * df["unit_price"]
return df
@asset
def daily_order_summary(cleaned_orders: pd.DataFrame) -> pd.DataFrame:
"""Aggregate orders by day."""
return cleaned_orders.groupby("order_date").agg(
order_count=("order_id", "count"),
total_revenue=("total", "sum"),
).reset_index()
@asset
def top_products(cleaned_orders: pd.DataFrame) -> pd.DataFrame:
"""Identify top-selling products."""
return (
cleaned_orders.groupby("product_id")
.agg(total_revenue=("total", "sum"))
.sort_values("total_revenue", ascending=False)
.head(20)
.reset_index()
)Ops and Graphs
Before SDAs, Dagster's primitives were ops(computational units) and graphs (compositions of ops). SDAs are now the recommended approach for data pipelines, but ops are still useful for non-asset-oriented workflows.
Jobs and Schedules
A job materializes a selection of assets. Aschedule runs a job on a cadence. Asensor triggers a job based on external events.
from dagster import define_asset_job, ScheduleDefinition, AssetSelection # Define a job that materializes all assets all_assets_job = define_asset_job( name="daily_materialize", selection=AssetSelection.all(), ) # Schedule it to run daily at 6 AM daily_schedule = ScheduleDefinition( job=all_assets_job, cron_schedule="0 6 * * *", ) # Sensor that triggers when a file lands in S3 from dagster_s3 import S3Sensor s3_sensor = S3Sensor( name="s3_file_sensor", bucket_name="landing-zone", prefix="incoming/", job=all_assets_job, )Data Quality and Asset Checks
Dagster natively supports data quality checks on assets using theAssetCheck framework. You can define checks that run after an asset is materialized.
from dagster import asset_check, AssetCheckResult
@asset_check(asset=cleaned_orders)
def orders_have_no_nulls(cleaned_orders: pd.DataFrame):
null_count = cleaned_orders["customer_id"].isna().sum()
return AssetCheckResult(
passed=null_count == 0,
metadata={"null_customer_ids": null_count},
)
@asset_check(asset=daily_order_summary)
def positive_revenue(daily_order_summary: pd.DataFrame):
negative = daily_order_summary[
daily_order_summary["total_revenue"] < 0
]
return AssetCheckResult(
passed=len(negative) == 0,
metadata={"negative_revenue_days": len(negative)},
)Dagster UI (Dagit)
Dagster's web UI provides a rich interface for exploring the asset graph, viewing lineage, launching runs, and inspecting logs. It shows the full dependency graph of your data assets, making it easy to understand the impact of changes.
Key Differentiators
- Asset-centric - Focus on data products rather than task execution
- Automatic lineage - Dagster tracks which assets produce and consume other assets
- Type system - Define schemas for asset inputs and outputs for validation
- Code locations - Deploy multiple code repositories under a single Dagster instance
- Cloud-native - Runs on Kubernetes with automatic scaling and isolation
Best Practices
- Define assets for every meaningful data product in your system
- Use asset checks to enforce data quality contracts
- Partition assets by date for efficient incremental processing
- Use freshness policies to alert on stale assets
- Separate code locations by team or domain boundary