Expert guidance for Dagster data orchestration including assets, resources, schedules, sensors, partitions, testing, and ETL patterns...
Think in Assets: Dagster is built around the asset abstractionβpersistent objects like tables, files, or models that your pipeline produces. Assets provide:
Assets over Ops: For most data pipelines, prefer assets over ops. Use ops only when the asset abstraction doesn't fit (non-data workflows, complex execution patterns).
Environment Separation: Use resources and EnvVar to maintain separate configurations for dev, staging, and production without code changes.
| If you're writing... | Check this section/reference |
|---|---|
@dg.asset |
Assets or references/assets.md |
ConfigurableResource |
Resources or references/resources.md |
AutomationCondition |
Declarative Automation or references/automation.md |
@dg.schedule or ScheduleDefinition |
Automation or references/automation.md |
@dg.sensor |
Sensors or references/automation.md |
PartitionsDefinition |
Partitions or references/automation.md |
Tests with dg.materialize() |
Testing or references/testing.md |
@asset_check |
references/testing.md#asset-checks |
@dlt_assets or @sling_assets |
references/etl-patterns.md |
@dbt_assets |
dbt Integration or dbt-development skill |
Definitions or code locations |
references/project-structure.md |
Components (defs.yaml) |
references/project-structure.md#components |
Asset: A persistent object (table, file, model) that your pipeline produces. Define with @dg.asset.
Resource: External services/tools (databases, APIs) shared across assets. Define with ConfigurableResource.
Job: A selection of assets to execute together. Create with dg.define_asset_job().
Schedule: Time-based automation for jobs. Create with dg.ScheduleDefinition.
Sensor: Event-driven automation that watches for changes. Define with @dg.sensor.
Partition: Logical divisions of data (by date, category). Define with PartitionsDefinition.
Definitions: The container for all Dagster objects in a code location.
Component: Reusable, declarative building blocks that generate Definitions from configuration (YAML). Use for standardized patterns.
Declarative Automation: Modern automation framework where you set conditions on assets rather than scheduling jobs.
import dagster as dg
@dg.asset
def my_asset() -> None:
"""Asset description appears in the UI."""
# Your computation logic here
pass
@dg.asset
def downstream_asset(upstream_asset) -> dict:
"""Depends on upstream_asset by naming it as a parameter."""
return {"processed": upstream_asset}
@dg.asset(
group_name="analytics",
key_prefix=["warehouse", "staging"],
description="Cleaned customer data",
owners=["team:data-engineering", "alice@example.com"],
tags={"priority": "high", "domain": "sales"},
code_version="1.2.0",
)
def customers() -> None:
pass
Best Practices:
customers, daily_revenue), not verbs (load_customers)from dagster import ConfigurableResource
class DatabaseResource(ConfigurableResource):
connection_string: str
def query(self, sql: str) -> list:
# Implementation here
pass
@dg.asset
def my_asset(database: DatabaseResource) -> None:
results = database.query("SELECT * FROM table")
dg.Definitions(
assets=[my_asset],
resources={"database": DatabaseResource(connection_string="...")},
)
import dagster as dg
from my_project.defs.jobs import my_job
my_schedule = dg.ScheduleDefinition(
job=my_job,
cron_schedule="0 0 * * *", # Daily at midnight
)
| Pattern | Meaning |
|---|---|
0 * * * * |
Every hour |
0 0 * * * |
Daily at midnight |
0 0 * * 1 |
Weekly on Monday |
0 0 1 * * |
Monthly on the 1st |
0 0 5 * * |
Monthly on the 5th |
Modern automation pattern: Set conditions on assets instead of scheduling jobs.
from dagster import AutomationCondition
# Update when upstream data changes
@dg.asset(
automation_condition=AutomationCondition.on_missing()
)
def my_asset() -> None:
pass
# Update daily at a specific time
@dg.asset(
automation_condition=AutomationCondition.on_cron("0 9 * * *")
)
def daily_report() -> None:
pass
# Combine conditions
@dg.asset(
automation_condition=(
AutomationCondition.on_missing()
| AutomationCondition.on_cron("0 0 * * *")
)
)
def flexible_asset() -> None:
pass
Benefits over Schedules:
When to Use:
@dg.sensor(job=my_job)
def my_sensor(context: dg.SensorEvaluationContext):
# 1. Read cursor (previous state)
previous_state = json.loads(context.cursor) if context.cursor else {}
current_state = {}
runs_to_request = []
# 2. Check for changes
for item in get_items_to_check():
current_state[item.id] = item.modified_at
if item.id not in previous_state or previous_state[item.id] != item.modified_at:
runs_to_request.append(dg.RunRequest(
run_key=f"run_{item.id}_{item.modified_at}",
run_config={...}
))
# 3. Return result with updated cursor
return dg.SensorResult(
run_requests=runs_to_request,
cursor=json.dumps(current_state)
)
Key: Use cursors to track state between sensor evaluations.
weekly_partition = dg.WeeklyPartitionsDefinition(start_date="2023-01-01")
@dg.asset(partitions_def=weekly_partition)
def weekly_data(context: dg.AssetExecutionContext) -> None:
partition_key = context.partition_key # e.g., "2023-01-01"
# Process data for this partition
region_partition = dg.StaticPartitionsDefinition(["us-east", "us-west", "eu"])
@dg.asset(partitions_def=region_partition)
def regional_data(context: dg.AssetExecutionContext) -> None:
region = context.partition_key
| Type | Use Case |
|---|---|
DailyPartitionsDefinition |
One partition per day |
WeeklyPartitionsDefinition |
One partition per week |
MonthlyPartitionsDefinition |
One partition per month |
HourlyPartitionsDefinition |
One partition per hour |
StaticPartitionsDefinition |
Fixed set of partitions |
DynamicPartitionsDefinition |
Partitions created at runtime |
MultiPartitionsDefinition |
Combine multiple partition dimensions |
Best Practice: Limit partitions to 100,000 or fewer per asset for optimal UI performance.
def test_my_asset():
result = my_asset()
assert result == expected_value
def test_asset_graph():
result = dg.materialize(
assets=[asset_a, asset_b],
resources={"database": mock_database},
)
assert result.success
assert result.output_for_node("asset_b") == expected
from unittest.mock import Mock
def test_with_mocked_resource():
mocked_resource = Mock()
mocked_resource.query.return_value = [{"id": 1}]
result = dg.materialize(
assets=[my_asset],
resources={"database": mocked_resource},
)
assert result.success
@dg.asset_check(asset=my_asset)
def validate_non_empty(my_asset):
return dg.AssetCheckResult(
passed=len(my_asset) > 0,
metadata={"row_count": len(my_asset)},
)
For dbt integration, prefer the component-based approach for standard dbt projects. Use Pythonic assets only when you need custom logic or fine-grained control.
Use DbtProjectComponent with remote Git repository:
# defs/transform/defs.yaml
type: dagster_dbt.DbtProjectComponent
attributes:
project:
repo_url: https://github.com/dagster-io/jaffle-platform.git
repo_relative_path: jdbt
dbt:
target: dev
When to use:
For private repositories:
attributes:
project:
repo_url: https://github.com/your-org/dbt-project.git
repo_relative_path: dbt
token: '{{ env.GIT_TOKEN }}'
dbt:
target: dev
For custom logic or local development:
from dagster_dbt import DbtCliResource, dbt_assets
from pathlib import Path
dbt_project_dir = Path(__file__).parent / "dbt_project"
@dbt_assets(manifest=dbt_project_dir / "target" / "manifest.json")
def my_dbt_assets(context: dg.AssetExecutionContext, dbt: DbtCliResource):
yield from dbt.cli(["build"], context=context).stream()
dg.Definitions(
assets=[my_dbt_assets],
resources={"dbt": DbtCliResource(project_dir=dbt_project_dir)},
)
When to use:
Full patterns: See Dagster dbt docs
references/assets.md when:references/resources.md when:ConfigurableResource classesreferences/automation.md when:references/testing.md when:dg.materialize() for integration testsreferences/etl-patterns.md when:references/project-structure.md when:Definitions and code locationsdg CLI for scaffoldingmy_project/
βββ pyproject.toml
βββ src/
β βββ my_project/
β βββ definitions.py # Main Definitions
β βββ defs/
β βββ assets/
β β βββ __init__.py
β β βββ my_assets.py
β βββ jobs.py
β βββ schedules.py
β βββ sensors.py
β βββ resources.py
βββ tests/
βββ test_assets.py
Auto-Discovery (Simplest):
# src/my_project/definitions.py
from dagster import Definitions
from dagster_dg import load_defs
# Automatically discovers all definitions in defs/ folder
defs = Definitions.merge(
load_defs()
)
Combining Components with Pythonic Assets:
# src/my_project/definitions.py
from dagster import Definitions
from dagster_dg import load_defs
from my_project.assets import custom_assets
# Load component definitions from defs/ folder
component_defs = load_defs()
# Define pythonic assets separately
pythonic_defs = Definitions(
assets=custom_assets,
resources={...}
)
# Merge them together
defs = Definitions.merge(component_defs, pythonic_defs)
Traditional (Explicit):
# src/my_project/definitions.py
from dagster import Definitions
from my_project.defs import assets, jobs, schedules, resources
defs = Definitions(
assets=assets,
jobs=jobs,
schedules=schedules,
resources=resources,
)
# Create new project
uvx create-dagster my_project
# Scaffold new asset file
dg scaffold defs dagster.asset assets/new_asset.py
# Scaffold schedule
dg scaffold defs dagster.schedule schedules.py
# Scaffold sensor
dg scaffold defs dagster.sensor sensors.py
# Validate definitions
dg check defs
trip_update_job = dg.define_asset_job(
name="trip_update_job",
selection=["taxi_trips", "taxi_zones"],
)
from dagster import Config
class MyAssetConfig(Config):
filename: str
limit: int = 100
@dg.asset
def configurable_asset(config: MyAssetConfig) -> None:
print(f"Processing {config.filename} with limit {config.limit}")
@dg.asset(deps=["external_table"])
def derived_asset() -> None:
"""Depends on external_table which isn't managed by Dagster."""
pass
| Anti-Pattern | Better Approach |
|---|---|
| Hardcoding credentials in assets | Use ConfigurableResource with env vars |
| Giant assets that do everything | Split into focused, composable assets |
| Ignoring asset return types | Use type annotations for clarity |
| Skipping tests for assets | Test assets like regular Python functions |
| Not using partitions for time-series | Use DailyPartitionsDefinition etc. |
| Putting all assets in one file | Organize by domain in separate modules |
# Development
dg dev # Start Dagster UI (port 3000)
dg check defs # Validate definitions load correctly
dg list defs # Show all loaded definitions
dg list components # Show available components
# Scaffolding
dg scaffold defs dagster.asset assets/file.py
dg scaffold defs dagster.schedule schedules.py
dg scaffold defs dagster.sensor sensors.py
dg scaffold defs dagster.resources resources.py
# Execution
dg launch --assets my_asset # Materialize specific asset
dg launch --assets "*" # Materialize all assets
# Use for non-dg projects or advanced scenarios
dagster dev # Start Dagster UI
dagster job execute -j my_job # Execute a job
dagster asset materialize -a my_asset # Materialize an asset
Use dg CLI for projects created with create-dagster. It provides auto-discovery, scaffolding, and modern workflow support.
references/assets.md - Detailed asset patternsreferences/resources.md - Resource configurationreferences/automation.md - Schedules, sensors, partitionsreferences/testing.md - Testing patterns and asset checksreferences/etl-patterns.md - dlt, Sling, file/API ingestionreferences/project-structure.md - Definitions, Components