Dagster Pipeline Testing: Assets, Jobs, and Sensors
Dagster is a data orchestration platform built around Software-Defined Assets (SDAs) — a model where pipelines declare what data assets they produce rather than what tasks to run. This asset-centric model makes testing more natural: you test that assets produce the right data, not just that code runs. This guide covers Dagster's testing utilities and patterns for assets, jobs, sensors, and resources.
Dagster Testing Philosophy
Dagster's testing approach centers on:
- Unit testing ops and assets — test the computation in isolation
- Integration testing with
materialize— run asset materialization with test resources - Job execution testing — run full jobs with
execute_in_process - Sensor/schedule testing — validate trigger conditions
Setting Up
pip install dagster dagster-webserver pytest# conftest.py
import pytest
from dagster import build_asset_context, materialize, DagsterInstance
@pytest.fixture(scope="session")
def dagster_instance():
"""Use an ephemeral instance for tests."""
with DagsterInstance.ephemeral() as instance:
yield instanceTesting Software-Defined Assets
# assets.py
from dagster import asset, AssetExecutionContext
import pandas as pd
@asset
def raw_orders(context: AssetExecutionContext) -> pd.DataFrame:
"""Fetch raw orders from source."""
# In production: read from database or API
# For testability, the actual data source is injected via resource
db = context.resources.database
return db.query("SELECT * FROM orders")
@asset
def cleaned_orders(raw_orders: pd.DataFrame) -> pd.DataFrame:
"""Clean and validate order data."""
df = raw_orders.copy()
# Remove duplicates
df = df.drop_duplicates(subset=['order_id'])
# Filter invalid amounts
df = df[df['amount'] > 0]
# Standardize status
df['status'] = df['status'].str.lower().str.strip()
return df
@asset
def order_metrics(cleaned_orders: pd.DataFrame) -> dict:
"""Compute aggregate order metrics."""
return {
'total_orders': len(cleaned_orders),
'total_revenue': float(cleaned_orders['amount'].sum()),
'avg_order_value': float(cleaned_orders['amount'].mean()),
'status_breakdown': cleaned_orders['status'].value_counts().to_dict(),
}# test_assets.py
import pytest
import pandas as pd
from dagster import materialize, build_asset_context
from mypackage.assets import cleaned_orders, order_metrics
class TestCleanedOrders:
def test_removes_duplicates(self):
raw = pd.DataFrame([
{'order_id': 'o1', 'amount': 100, 'status': 'pending'},
{'order_id': 'o1', 'amount': 100, 'status': 'pending'}, # duplicate
{'order_id': 'o2', 'amount': 50, 'status': 'completed'},
])
result = cleaned_orders(raw)
assert len(result) == 2
assert result['order_id'].nunique() == 2
def test_filters_zero_amount_orders(self):
raw = pd.DataFrame([
{'order_id': 'o1', 'amount': 0, 'status': 'pending'},
{'order_id': 'o2', 'amount': 100, 'status': 'completed'},
])
result = cleaned_orders(raw)
assert len(result) == 1
assert result.iloc[0]['order_id'] == 'o2'
def test_normalizes_status(self):
raw = pd.DataFrame([
{'order_id': 'o1', 'amount': 100, 'status': ' PENDING '},
{'order_id': 'o2', 'amount': 50, 'status': 'Completed'},
])
result = cleaned_orders(raw)
assert result.iloc[0]['status'] == 'pending'
assert result.iloc[1]['status'] == 'completed'
def test_handles_empty_dataframe(self):
raw = pd.DataFrame(columns=['order_id', 'amount', 'status'])
result = cleaned_orders(raw)
assert len(result) == 0
assert list(result.columns) == ['order_id', 'amount', 'status']
class TestOrderMetrics:
def test_computes_correct_metrics(self):
orders = pd.DataFrame([
{'order_id': 'o1', 'amount': 100.0, 'status': 'completed'},
{'order_id': 'o2', 'amount': 200.0, 'status': 'completed'},
{'order_id': 'o3', 'amount': 50.0, 'status': 'pending'},
])
metrics = order_metrics(orders)
assert metrics['total_orders'] == 3
assert metrics['total_revenue'] == 350.0
assert metrics['avg_order_value'] == pytest.approx(116.67, rel=0.01)
assert metrics['status_breakdown']['completed'] == 2
assert metrics['status_breakdown']['pending'] == 1Testing Assets with Resources
When assets depend on resources (databases, APIs), inject test resources:
# resources.py
from dagster import ConfigurableResource, resource
import pandas as pd
class DatabaseResource(ConfigurableResource):
connection_string: str
def query(self, sql: str) -> pd.DataFrame:
import sqlalchemy as sa
engine = sa.create_engine(self.connection_string)
with engine.connect() as conn:
return pd.read_sql(sql, conn)
# Test resource
class MockDatabaseResource(ConfigurableResource):
"""In-memory mock database for testing."""
def query(self, sql: str) -> pd.DataFrame:
# Return test data based on query
if 'orders' in sql.lower():
return pd.DataFrame([
{'order_id': 'test-1', 'amount': 100.0, 'status': 'completed'},
{'order_id': 'test-2', 'amount': 200.0, 'status': 'pending'},
])
return pd.DataFrame()# test_assets_with_resources.py
from dagster import materialize
from mypackage.assets import raw_orders, cleaned_orders, order_metrics
from mypackage.resources import MockDatabaseResource
def test_full_asset_pipeline():
result = materialize(
assets=[raw_orders, cleaned_orders, order_metrics],
resources={
'database': MockDatabaseResource(),
},
)
assert result.success
# Check materialized asset values
raw_df = result.output_for_node('raw_orders')
assert len(raw_df) == 2
metrics = result.output_for_node('order_metrics')
assert metrics['total_orders'] == 2
assert metrics['total_revenue'] == 300.0
def test_asset_metadata_logged():
result = materialize(
assets=[raw_orders, cleaned_orders],
resources={'database': MockDatabaseResource()},
)
assert result.success
# Check that metadata was logged
materialization = result.get_asset_materialization_event('cleaned_orders')
assert materialization is not NoneTesting Jobs
# jobs.py
from dagster import define_asset_job
daily_orders_job = define_asset_job(
name="daily_orders_job",
selection=["raw_orders", "cleaned_orders", "order_metrics"],
)# test_jobs.py
from dagster import materialize
from mypackage.jobs import daily_orders_job
from mypackage.resources import MockDatabaseResource
def test_daily_orders_job_succeeds():
result = daily_orders_job.execute_in_process(
resources={'database': MockDatabaseResource()},
)
assert result.success
assert not result.get_failed_step_keys()
def test_job_handles_empty_data():
class EmptyDatabaseResource(MockDatabaseResource):
def query(self, sql: str):
import pandas as pd
return pd.DataFrame(columns=['order_id', 'amount', 'status'])
result = daily_orders_job.execute_in_process(
resources={'database': EmptyDatabaseResource()},
)
# Job should succeed even with empty data
assert result.success
metrics = result.output_for_node('order_metrics')
assert metrics['total_orders'] == 0
assert metrics['total_revenue'] == 0.0Testing Sensors
# sensors.py
from dagster import sensor, RunRequest, SensorEvaluationContext, SkipReason
@sensor(job=daily_orders_job)
def new_orders_sensor(context: SensorEvaluationContext):
"""Trigger job when new orders arrive in S3."""
s3 = context.resources.s3
new_files = s3.list_new_files(
bucket='orders-bucket',
prefix='incoming/',
after_timestamp=context.cursor,
)
if not new_files:
return SkipReason("No new order files found")
# Update cursor to latest file timestamp
latest_timestamp = max(f.timestamp for f in new_files)
context.update_cursor(str(latest_timestamp))
return [
RunRequest(
run_key=f.name,
run_config={'ops': {'raw_orders': {'config': {'file': f.name}}}},
)
for f in new_files
]# test_sensors.py
from dagster import build_sensor_context
from unittest.mock import MagicMock
from mypackage.sensors import new_orders_sensor
def test_sensor_triggers_on_new_files():
mock_s3 = MagicMock()
mock_s3.list_new_files.return_value = [
MagicMock(name='orders-2024-01-15.csv', timestamp=1705276800),
MagicMock(name='orders-2024-01-15-late.csv', timestamp=1705280400),
]
context = build_sensor_context(
cursor="0",
resources={'s3': mock_s3},
)
result = new_orders_sensor.evaluate_tick(context)
assert len(result.run_requests) == 2
assert result.cursor == "1705280400" # Updated to latest
def test_sensor_skips_when_no_files():
mock_s3 = MagicMock()
mock_s3.list_new_files.return_value = []
context = build_sensor_context(
cursor="1705276800",
resources={'s3': mock_s3},
)
result = new_orders_sensor.evaluate_tick(context)
assert len(result.run_requests) == 0
assert result.skip_message == "No new order files found"
assert result.cursor == "1705276800" # Cursor unchanged
def test_sensor_deduplicates_on_run_key():
"""Same file should not trigger duplicate runs."""
mock_s3 = MagicMock()
mock_s3.list_new_files.return_value = [
MagicMock(name='orders-2024-01-15.csv', timestamp=1705276800),
]
context = build_sensor_context(
cursor="0",
resources={'s3': mock_s3},
)
# Evaluate twice — second evaluation should not re-run
result = new_orders_sensor.evaluate_tick(context)
assert result.run_requests[0].run_key == 'orders-2024-01-15.csv'
# Dagster uses run_key to deduplicate — same key = no duplicate runTesting Partitioned Assets
# partitioned_assets.py
from dagster import asset, DailyPartitionsDefinition, AssetExecutionContext
import pandas as pd
daily_partitions = DailyPartitionsDefinition(start_date="2024-01-01")
@asset(partitions_def=daily_partitions)
def daily_revenue(context: AssetExecutionContext) -> pd.DataFrame:
"""Revenue aggregated per day."""
partition_date = context.partition_key # e.g., "2024-01-15"
db = context.resources.database
return db.query(f"""
SELECT date, SUM(amount) as revenue, COUNT(*) as orders
FROM orders
WHERE DATE(created_at) = '{partition_date}'
GROUP BY date
""")# test_partitioned.py
from dagster import materialize, build_asset_context
from mypackage.assets import daily_revenue
from mypackage.resources import MockDatabaseResource
def test_daily_revenue_for_partition():
result = materialize(
assets=[daily_revenue],
resources={'database': MockDatabaseResource()},
partition_key="2024-01-15",
)
assert result.success
df = result.output_for_node('daily_revenue')
assert 'revenue' in df.columns
assert 'orders' in df.columns
def test_backfill_multiple_partitions():
"""Test materializing multiple partitions in sequence."""
for date in ['2024-01-01', '2024-01-02', '2024-01-03']:
result = materialize(
assets=[daily_revenue],
resources={'database': MockDatabaseResource()},
partition_key=date,
)
assert result.success, f"Failed for partition {date}"Integration Testing with Real Database
# conftest.py
import pytest
import sqlalchemy as sa
import pandas as pd
from dagster import DagsterInstance
@pytest.fixture(scope="session")
def test_database():
"""Create and seed a test SQLite database."""
engine = sa.create_engine("sqlite:///test.db")
with engine.connect() as conn:
conn.execute(sa.text("""
CREATE TABLE IF NOT EXISTS orders (
order_id TEXT PRIMARY KEY,
amount REAL,
status TEXT,
created_at TEXT
)
"""))
conn.execute(sa.text("""
INSERT OR REPLACE INTO orders VALUES
('o1', 100.0, 'completed', '2024-01-15 10:00:00'),
('o2', 200.0, 'pending', '2024-01-15 11:00:00'),
('o3', 50.0, 'completed', '2024-01-15 12:00:00')
"""))
conn.commit()
yield engine
# Cleanup
import os
os.unlink("test.db")
def test_full_pipeline_integration(test_database):
from mypackage.resources import DatabaseResource
result = materialize(
assets=[raw_orders, cleaned_orders, order_metrics],
resources={
'database': DatabaseResource(
connection_string="sqlite:///test.db"
),
},
)
assert result.success
metrics = result.output_for_node('order_metrics')
assert metrics['total_orders'] == 3
assert metrics['total_revenue'] == pytest.approx(350.0)Dagster's testing utilities make it straightforward to test assets in isolation (just call the function with test data) and in integration (use materialize with mock resources). The asset model's explicit inputs and outputs make tests readable — you can see exactly what data flows between stages.