This lesson on Modern Orchestration with Dagster — Software-Defined Assets is hands-on and example-driven. You will build and orchestrate an end-to-end ETL data pipeline using Dagster's Software-Defined Assets (SDAs). You will chain asset dependencies implicitly through function signatures, schedule automated materializations, selectively re-execute failed assets from cached states, and secure API credentials using Dagster resources.
What You'll Be Able To Do
- Scaffold a standard Dagster project and launch the web UI using Dagster CLI commands.
- Define extract, transform, and load steps as Software-Defined Assets using the @asset decorator.
- Wire upstream and downstream asset dependencies automatically via function parameter names.
- Debug runtime failures and execute selective re-materializations using Dagster's cached outputs.
- Bundle asset graphs into executable production jobs and attach automated schedules.
- Extract API tokens and credentials from business logic into configurable @resource definitions.
Detailed Concept Walkthrough
1. Software-Defined Assets and Dependency Resolution
Software-Defined Assets (SDAs) declare the desired data artifact alongside the code needed to generate it. Dagster builds the pipeline DAG automatically by matching downstream function parameter names to upstream asset names.
- Mechanism: Decorating a Python function with
@assetturns it into a data asset definition where the return value represents the materialized artifact. - Dependency Wiring: Dagster resolves DAG edges implicitly; naming an asset function's input parameter after an upstream asset automatically establishes data flow between them.
- Under the Hood: IO managers handle intermediate serialization, persisting upstream return values and passing them into downstream asset functions as in-memory arguments.
- Execution Flow: During execution, Dagster topologically sorts the resolved asset graph, running extraction before transformation and aggregation.
from dagster import asset
import pandas as pd
@asset
def GitHubStargazers() -> pd.DataFrame:
# Extracts raw data from API
return pd.DataFrame([{"user": "alice", "date": "2023-01-01"}])
@asset
def GitHubStargazers_by_week(GitHubStargazers: pd.DataFrame) -> pd.DataFrame:
# Parameter name matches upstream asset name exactly
return GitHubStargazers.groupby("date").count()
Key Takeaway: Asset dependencies are declared implicitly through function argument names rather than manual pipeline wiring.
2. Runtime Context, Materialization, and Selective Recovery
Materialization executes asset code and updates data tracking in Dagster. When errors occur, intermediate data is cached, allowing selective downstream re-execution without recalculating upstream steps.
- Mechanism: Asset functions requiring runtime metadata, logging, or partition details must explicitly request
contextas their first parameter. - Under the Hood: Dagster caches successfully materialized asset outputs via the configured IO manager, decoupling upstream extraction from downstream reporting.
- Failure Recovery: If a downstream asset fails due to a bug, fixing the code allows developers to re-materialize only the failed asset while reusing cached upstream outputs.
- Best Practice: Always use
context.log.info()rather than standardprint()statements to ensure structured operational logs are captured in the UI.
from dagster import asset, AssetExecutionContext
@asset
def upload_report(context: AssetExecutionContext, GitHubStargazers_by_week):
# Access runtime logging through context
context.log.info("Uploading report to external destination...")
report_url = "https://gist.github.com/example_id"
context.log.info(f"Report published at {report_url}")
return report_url
Key Takeaway: Explicit context arguments enable structured logging, while asset caching prevents redundant re-computation of upstream stages.
3. Asset Jobs and Automated Scheduling
Asset definitions specify what data to build, while Jobs and Schedules determine when and how those assets are materialized in production.
- Mechanism:
define_asset_jobselects a subset or full graph of assets to materialize together as a cohesive batch job. - Scheduling:
ScheduleDefinitionbinds an asset job to a cron expression or standard interval, triggering automated execution runs. - Operational Control: Separating asset declarations from jobs allows the same pipeline graph to be triggered on different cadences (e.g., hourly vs. daily) without changing transformation code.
from dagster import define_asset_job, ScheduleDefinition
# Group all software-defined assets into an executable unit
daily_refresh_job = define_asset_job(name="daily_refresh", selection="*")
# Schedule the job to run every day at midnight
daily_schedule = ScheduleDefinition(
job=daily_refresh_job,
cron_schedule="0 0 * * *"
)
Key Takeaway: Jobs package asset graphs into runnable units, enabling scheduled automated execution without modifying asset logic.
4. Configurable Resources and Secret Isolation
Resources isolate infrastructure connections, client libraries, and secrets from data processing code, ensuring security and environmental portability.
- Mechanism: The
@resourcedecorator defines reusable external services (such as API clients) with strong configuration schemas. - Security: Sensitive credentials like API tokens are defined as resource configs rather than hardcoded inside asset transformation logic.
- Under the Hood: Dagster validates provided configuration at pipeline launch and supplies initialized resource objects directly to executing assets.
from dagster import resource, Field, String
@resource(config_schema={"access_token": Field(String, is_required=True)})
def github_api_client(init_context):
token = init_context.resource_config["access_token"]
# Return authenticated client or session
return {"Authorization": f"Bearer {token}"}
Key Takeaway: Resources decouple sensitive credentials and external system configurations from core asset transformation functions.
Topics Covered in Modern Orchestration with Dagster — Software-Defined Assets
- Project Scaffolding (0:15 - 1:23) — Initialize the Gitpod development environment and scaffold a standard Dagster project using the CLI.
- Environment & Dagit Launch (1:23 - 1:54) — Examine the default configuration files and start the Dagit web UI.
- Extraction Asset Definition (1:54 - 2:22) — Create the initial asset to fetch raw stargazer records from the GitHub API.
- Transformation Asset Chaining (2:22 - 3:41) — Build the weekly aggregation asset by referencing the upstream extraction asset as an input parameter.
- Visualization and Loading Assets (3:41 - 4:34) — Define assets that generate a Jupyter notebook report and upload the final output to GitHub Gist.
- Materialization and Error Recovery (4:34 - 5:38) — Materialize the asset graph, resolve a missing context parameter bug, and re-execute failed assets selectively.
- Job and Schedule Creation (5:38 - 6:36) — Wrap all asset definitions into a refresh job and configure a recurring daily schedule.
- Resource Configuration (6:36 - 7:59) — Refactor raw API tokens into a parameterized resource with a validated configuration schema.
Data Engineering Cheat Sheet
-
dagster project scaffold --name <proj>— Initializes a standard Dagster project directory structuredagster project scaffold --name stargazers_pipeline -
dagit— Launches the local web UI for pipeline managementdagit -
@asset— Decorates a function to define a Software-Defined Asset@asset def raw_data(): return [1, 2, 3] -
define_asset_job(name, selection)— Creates an executable job from selected assetsjob = define_asset_job(name="sync_job", selection="*") -
ScheduleDefinition(job, cron_schedule)— Attaches a recurring schedule to an asset jobsched = ScheduleDefinition(job=sync_job, cron_schedule="@daily") -
@resource(config_schema={...})— Defines a configurable external dependency or credential@resource(config_schema={"token": str}) def api_conn(): pass
Comparison Table
| Entity | Primary Responsibility | Defined By |
|---|---|---|
| Software-Defined Asset | Declares data artifact and transformation computation | @asset decorator |
| Asset Job | Groups assets for targeted runtime materialization | define_asset_job function |
| Resource | Manages external connections and credentials | @resource decorator |
Common Pitfalls
- Mistake: Hardcoding secrets like API tokens inside asset functions. Avoid: Isolate credentials in configurable resources using config schemas.
- Mistake: Calling context logging methods without declaring context in function signature. Avoid: Explicitly add context as the first argument in asset functions.
- Mistake: Giving downstream parameters names that differ from upstream assets. Avoid: Match input argument names exactly to upstream asset function names.
- Mistake: Re-executing all upstream assets when fixing a downstream bug. Avoid: Use selective re-materialization in the UI to run only failed assets.
FAQs
- How does Dagster determine execution order across assets? Dagster parses function parameter names in downstream assets and matches them to upstream asset names to construct the dependency graph.
- Do I need to manually serialize and load files between asset steps? No, Dagster's built-in IO managers automatically handle saving outputs and loading them as inputs for downstream functions.
- Why did my asset fail when calling context.log.info?
The
contextparameter is not globally available; it must be explicitly defined as an argument in your asset function signature. - What is the primary difference between Dagit and the Dagster CLI? The Dagster CLI scaffolds and executes project commands, whereas Dagit is the web-based graphical interface for visualizing and materializing assets.