dagster trained

versions practiced: dagster 1.13.24, dagster-dbt 0.29.24, dagster-duckdb 0.29.24, dagster-dg-cli 1.13.24, dagster-webserver 1.13.24, dbt-core 1.12.5, dbt-duckdb 1.11.0, duckdb 1.5.5, pytest 9.1.1, python 3.12. pipeline design principles (idempotency, late data, gates) live in data-pipeline; this page covers dagster itself.

mental model

definitions
Definitions ─┬─ assets (@asset, @dbt_assets, @multi_asset) <- what should exist: a table, file, model
├─ asset_checks (@asset_check, blocking=True) <- is it correct; blocking = gate downstream
├─ partitions_def (Daily/Static/Multi) <- the unit of work & backfill; context.partition_key
├─ jobs (define_asset_job) <- a selection of assets to run together
├─ schedules / sensors / automation_condition <- when to run: time, external change, declarative
├─ resources (ConfigurableResource, DuckDBResource, DbtCliResource) <- injected clients + config
└─ executor (multiprocess default | in_process) <- how steps run
run = one execution of a job over one partition (or a range when every asset has BackfillPolicy.single_run)
  • assets, not tasks. you declare what should exist and its deps, and dagster infers the graph and can answer "is this up to date?" (sda blog). ops and op jobs are for side-effect workflows with no persistent output (ops).
  • returning a value invokes an io manager. the default PickledObjectFilesystemIOManager pickles it to $DAGSTER_HOME/storage (observed). return MaterializeResult when the asset writes its own storage (io managers).
  • checks are steps. a blocking check that fails stops downstream steps from starting in the same run.
  • partition mappings describe cross-partition deps. TimeWindowPartitionMapping(start_offset=-1) means "d reads d-1" (dependencies between partitioned assets).
  • two front doors in 1.13: pythonic Definitions (dagster ... -f defs.py), and the dg cli with components (create-dagster project, dg scaffold defs, dg check defs, dg launch) (projects).

examples

the dagster-crypto-assets example is verified end to end: raw parquet (binance zips) -> stg_klines (DuckDBResource) -> blocking gap check -> dbt daily_bars via @dbt_assets -> bar_features. it includes a schedule, a revised-file sensor, a Landing ConfigurableResource, and 6 pytest tests.

blocking check (the gate):

@dg.asset_check(asset=stg_klines, blocking=True)
def stg_klines_no_gaps(context: dg.AssetCheckExecutionContext, duckdb: DuckDBResource, landing: Landing):
day = context.run.tags["dagster/partition"]
...
return dg.AssetCheckResult(passed=not short, metadata={"short": json.dumps(short)})

partitioned dbt that also survives a range backfill:

@dbt_assets(manifest=dbt_project.manifest_path, partitions_def=daily)
def dbt_models(context: dg.AssetExecutionContext, dbt: DbtCliResource):
w = context.partition_time_window # NOT context.partition_key: a range run has none
yield from dbt.cli(["build", "--vars", json.dumps({"start": f"{w.start:%Y-%m-%d}", "end": f"{w.end:%Y-%m-%d}"})],
context=context).stream()

sensor, one run per changed day, with a sha cursor:

for day, shas in sorted(changed.items()):
yield dg.RunRequest(run_key=f"{day}:{'-'.join(shas)}", partition_key=day)
context.update_cursor(json.dumps(seen))

unit test without the network or a ui:

res = dg.materialize([raw_klines, stg_klines, stg_klines_no_gaps, dbt_models, bar_features],
partition_key="2025-01-02", resources=resources(str(tmp_path), ["BTCUSDT"], offline=True),
raise_on_error=False)

best practices

  1. model data as assets with a time partition. loop over entities (symbols) inside the asset unless you truly need per-entity backfills. multi-partitions carry open backfill bugs (#34223).
  2. write storage explicitly and return MaterializeResult(metadata=...) when the layout is a contract (hive parquet, warehouse tables). use io managers for in-memory handoffs (dataframes between python assets). don't pickle production data by accident.
  3. resources for every external thing (paths, db, dbt cli, http). config lives in ConfigurableResource fields, and tests swap them (offline=True) (resources).
  4. gates as blocking=True asset checks on the upstream asset. dbt tests become non-blocking checks (#25112).
  5. declare lookback with AssetDep(..., partition_mapping=TimeWindowPartitionMapping(start_offset=-1)), so the ui and automation see d+1 as stale after d changes.
  6. schedules: build_schedule_from_partitioned_job runs the last complete partition (verified: tick 2025-01-05 02:00 targets 2025-01-04).
  7. sensors: cursor + run_key, grouped by partition. the run_key makes a re-evaluation idempotent, and the cursor keeps it cheap.
  8. automation: AutomationCondition (eager / on_cron / on_missing), not AutoMaterializePolicy. freshness: FreshnessPolicy.time_window(fail_window=...), which went ga in 1.12 and supersedes build_*_freshness_checks.
  9. test with dg.materialize() + fixtures, build_sensor_context, build_schedule_context. resolve unresolved schedules through defs.resolve_schedule_def(name).
  10. validate in ci: dagster definitions validate -f defs.py or dg check defs. both are headless and need no webserver.

strengths

  • the asset graph is the lineage. dbt sources map to upstream dagster assets with meta.dagster.asset_key, which produced one graph from raw parquet through dbt to the features asset.
  • partitions give backfill, reruns, and staleness per day for free. checks give gates with ui history.
  • everything runs headless: cli materialize, python materialize(), sensor preview, pytest. dagster dev is optional.
  • official agent skill and dg scaffolding keep projects conventional.

weaknesses / pain points

  • the default multiprocess executor spawns a subprocess per step (about 2–6 s each observed), which is slow for tiny steps and contends for single-writer stores such as duckdb.
  • range backfills from the cli need BackfillPolicy.single_run() on every selected asset. otherwise you need the daemon/ui, or a loop.
  • deprecation churn: AutoMaterializePolicy -> AutomationCondition, FreshnessPolicy -> LegacyFreshnessPolicy + new FreshnessPolicy, define_asset_job(partitions_def=) deprecated. its own warning still says "removed in 1.10.0" in 1.13.24.
  • dbt integration quirks: blocking dbt tests are not possible, test-name collisions drop checks (#34155), and dbt fusion is misdetected (#34215).

gotchas

  1. DbtCliResource fails validation when dbt is not on path ("the dbt executable 'dbt' does not exist"), even with the venv's python. activate the venv or pass dbt_executable.
  2. @dbt_assets asset keys are the dbt model names (daily_bars), not the function name. --select dbt_models raised DagsterInvalidSubsetError.
  3. @dbt_assets defaults to a single-run backfill policy, so a range run has no partition_key (cannot access partition_key for a partitioned run with a range of partitions). use partition_time_window and range vars in dbt.
  4. mixing policies warns "materializes assets with varying backfillpolicies ... using BackfillPolicy.multi_run(max_partitions_per_run=1)".
  5. removing partitions_def= from define_asset_job (as the deprecation says) leaves the partitioned schedule unresolved. schedule.evaluate_tick then fails with AttributeError, so resolve it via Definitions.
  6. blocked steps emit no STEP_SKIPPED. downstream steps never STEP_START, so tests must assert on absence.
  7. the check context has no partition_key. read context.run.tags["dagster/partition"] in @asset_check for partitioned assets.
  8. a sensor yielding one runrequest per file launched 6 runs for 3 days (2 symbols). group by partition key.
  9. AutoMaterializePolicy is still importable in 1.13.24 and only warns. don't assume a removal happened because the warning names a past version, and verify against the installed version.
  10. the default executor is multiprocess even for dagster asset materialize. pin executor=dg.in_process_executor for duckdb-backed projects.

known bugs

the most relevant upstream issues, each with a workaround below:

  • #25112: dbt checks cannot block.
  • #33802: no read-only DuckDBResource.
  • #13255: out-of-range backfill partitions not rejected.
  • #34223: multi-partition backfills collapse.

troubleshooting

symptomcausefix
ValidationError ... dbt executable 'dbt' does not existvenv not on pathexport PATH=.venv/bin:$PATH or DbtCliResource(dbt_executable=...)
DagsterInvalidSubsetError: ['dbt_models'] ... no AssetsDefinitionselected the function nameselect the dbt model key (daily_bars) or *
Cannot access partition_key for a partitioned run with a rangesingle-run backfill of dbt assetscontext.partition_time_window -> dbt vars start/end
--partition-range CheckError: requires BackfillPolicy.single_run()non-dbt assets default to multi-runloop partitions, use the daemon/ui backfill, or set backfill_policy=dg.BackfillPolicy.single_run() and handle ranges
AttributeError: 'UnresolvedPartitionedAssetScheduleDefinition' ... evaluate_tickschedule built from a job without partitions_defdefs.resolve_schedule_def(name) in tests
sensor launches duplicate runs for a dayone runrequest per fileaggregate changes per partition key; run_key = day + shas
Could not set lock on file ... .duckdbparallel step processesin_process_executor; readers via duckdb.connect(read_only=True)
dg list defs: "must be run inside a dagster project directory"plain defs.py, no [tool.dg]use dagster ... -f defs.py, or uvx create-dagster project for a dg project

practiced cases

all headless, on public binance data, btcusdt + ethusdt. example: dagster-crypto-assets.

  1. definitions load: dagster definitions validate -f defs.py -> "all code locations passed validation". dagster asset list shows bar_features, daily_bars, raw_klines, stg_klines. dagster dev -p 3999 served /server_info with http 200 and version 1.13.24, then stopped.
  2. daily partitions end to end: dagster asset materialize --select '*' --partition for 2025-01-01..03 gave 3/3 run_success. each run passed stg_klines_no_gaps and 2 dbt not_null checks. daily_bars holds 6 rows with 1440 minutes each, and bar_features.ret_1d for btc 01-02 is 0.02498 (01-01 is null: no d-1 in range).
  3. range backfill: --select daily_bars --partition-range 2025-01-01...2025-01-03 failed first, then after the partition_time_window fix ran as one run, one dbt step in 7.14 s, with checks passing. --select '*' over a range is refused without single_run policies.
  4. blocking gate: 1435-minute fixture -> stg_klines_no_gaps failed, raw_klines and stg_klines materialized, and dbt_models and bar_features never started (pytest test_gap_blocks_downstream).
  5. sensor: dagster sensor preview landing_sensor first returned 6 run requests for 3 days (the per-file bug above), and returned 3 after grouping. pytest proves: first tick requests the day, an unchanged file gives no request, a republished file requests it again, and two symbols give one request.
  6. schedule: build_schedule_from_partitioned_job(hour_of_day=2) has cron 0 2 * * *. a tick at 2025-01-05 02:00 utc targets partition 2025-01-04.
  7. pytest: 6 passed in 10.2 s (offline fixtures, a tmp duckdb, and a real dbt subprocess in test_full_graph_with_dbt).
  8. dg cli: uvx create-dagster@1.13.24 project dgproj --uv-sync, then dg scaffold defs dagster.asset assets/hello.py, dg check defs ("all definitions loaded successfully"), dg list defs (1 asset), and dg launch --assets '*' (run_success). definitions.py uses load_from_defs_folder.
  9. api verification (installed 1.13.24): AutomationCondition.{eager,on_cron,on_missing} exist. FreshnessPolicy.time_window(fail_window, warn_window=None) works on an asset. AutoMaterializePolicy.eager() still works and emits a deprecationwarning.
  10. executor: the multiprocess run of the crypto pipeline took about 80 s for one partition against about 13 s in-process (data-pipeline).

not practiced (unverified here): pipes, components yaml beyond scaffolding, the daemon-driven ui backfills, dagster+ (serverless/hybrid), run-status sensors, concurrency pools, dagster-dlt, and AutomationCondition evaluation by the daemon.

ecosystem

needpicknote
dbt in the graphdagster-dbt @dbt_assets + DbtProjectmap dbt sources to upstream assets with meta.dagster.asset_key (dbt)
duckdbdagster-duckdb DuckDBResourceone writer; no read-only mode (#33802)
api ingestiondagster-dltdlt itself in data-pipeline
external compute (spark, k8s, scripts)pipescovered by dagster's external-pipelines guide
scaffolding / project layoutcreate-dagster, dg, componentscomponents
hostingoss (webserver + daemon + dagster.yaml) or dagster+dagster's own deployment guides cover both shapes
alternativesairflow 3 (assets now exist), prefect 3comparison in data-pipeline

search pages

go to any page