Dagster Pipeline Testing: Assets, Jobs, and Sensors

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:

  1. Unit testing ops and assets — test the computation in isolation
  2. Integration testing with materialize — run asset materialization with test resources
  3. Job execution testing — run full jobs with execute_in_process
  4. 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 instance

Testing 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'] == 1

Testing 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 None

Testing 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.0

Testing 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 run

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

Read more

Start now free