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 ─┬─ 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 runrun = 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
PickledObjectFilesystemIOManagerpickles it to$DAGSTER_HOME/storage(observed). returnMaterializeResultwhen 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 thedgcli 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 noneyield 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
- 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).
- 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. - resources for every external thing (paths, db, dbt cli, http). config lives in
ConfigurableResourcefields, and tests swap them (offline=True) (resources). - gates as
blocking=Trueasset checks on the upstream asset. dbt tests become non-blocking checks (#25112). - declare lookback with
AssetDep(..., partition_mapping=TimeWindowPartitionMapping(start_offset=-1)), so the ui and automation see d+1 as stale after d changes. - schedules:
build_schedule_from_partitioned_jobruns the last complete partition (verified: tick 2025-01-05 02:00 targets2025-01-04). - sensors: cursor +
run_key, grouped by partition. the run_key makes a re-evaluation idempotent, and the cursor keeps it cheap. - automation:
AutomationCondition(eager / on_cron / on_missing), notAutoMaterializePolicy. freshness:FreshnessPolicy.time_window(fail_window=...), which went ga in 1.12 and supersedesbuild_*_freshness_checks. - test with
dg.materialize()+ fixtures,build_sensor_context,build_schedule_context. resolve unresolved schedules throughdefs.resolve_schedule_def(name). - validate in ci:
dagster definitions validate -f defs.pyordg 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 devis optional. - official agent skill and
dgscaffolding 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+ newFreshnessPolicy,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
DbtCliResourcefails validation whendbtis not on path ("the dbt executable 'dbt' does not exist"), even with the venv's python. activate the venv or passdbt_executable.@dbt_assetsasset keys are the dbt model names (daily_bars), not the function name.--select dbt_modelsraisedDagsterInvalidSubsetError.@dbt_assetsdefaults to a single-run backfill policy, so a range run has nopartition_key(cannot access partition_key for a partitioned run with a range of partitions). usepartition_time_windowand range vars in dbt.- mixing policies warns "materializes assets with varying backfillpolicies ... using BackfillPolicy.multi_run(max_partitions_per_run=1)".
- removing
partitions_def=fromdefine_asset_job(as the deprecation says) leaves the partitioned schedule unresolved.schedule.evaluate_tickthen fails withAttributeError, so resolve it viaDefinitions. - blocked steps emit no
STEP_SKIPPED. downstream steps neverSTEP_START, so tests must assert on absence. - the check context has no
partition_key. readcontext.run.tags["dagster/partition"]in@asset_checkfor partitioned assets. - a sensor yielding one runrequest per file launched 6 runs for 3 days (2 symbols). group by partition key.
AutoMaterializePolicyis 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.- the default executor is multiprocess even for
dagster asset materialize. pinexecutor=dg.in_process_executorfor duckdb-backed projects.
known bugs
troubleshooting
| symptom | cause | fix |
|---|---|---|
| ValidationError ... dbt executable 'dbt' does not exist | venv not on path | export PATH=.venv/bin:$PATH or DbtCliResource(dbt_executable=...) |
| DagsterInvalidSubsetError: ['dbt_models'] ... no AssetsDefinition | selected the function name | select the dbt model key (daily_bars) or * |
| Cannot access partition_key for a partitioned run with a range | single-run backfill of dbt assets | context.partition_time_window -> dbt vars start/end |
| --partition-range CheckError: requires BackfillPolicy.single_run() | non-dbt assets default to multi-run | loop partitions, use the daemon/ui backfill, or set backfill_policy=dg.BackfillPolicy.single_run() and handle ranges |
| AttributeError: 'UnresolvedPartitionedAssetScheduleDefinition' ... evaluate_tick | schedule built from a job without partitions_def | defs.resolve_schedule_def(name) in tests |
| sensor launches duplicate runs for a day | one runrequest per file | aggregate changes per partition key; run_key = day + shas |
| Could not set lock on file ... .duckdb | parallel step processes | in_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.
- definitions load:
dagster definitions validate -f defs.py-> "all code locations passed validation".dagster asset listshowsbar_features, daily_bars, raw_klines, stg_klines.dagster dev -p 3999served/server_infowith http 200 and version 1.13.24, then stopped. - daily partitions end to end:
dagster asset materialize --select '*' --partitionfor 2025-01-01..03 gave 3/3 run_success. each run passedstg_klines_no_gapsand 2 dbtnot_nullchecks.daily_barsholds 6 rows with 1440 minutes each, andbar_features.ret_1dfor btc 01-02 is 0.02498 (01-01 is null: no d-1 in range). - range backfill:
--select daily_bars --partition-range 2025-01-01...2025-01-03failed first, then after thepartition_time_windowfix ran as one run, one dbt step in 7.14 s, with checks passing.--select '*'over a range is refused without single_run policies. - blocking gate: 1435-minute fixture ->
stg_klines_no_gapsfailed,raw_klinesandstg_klinesmaterialized, anddbt_modelsandbar_featuresnever started (pytesttest_gap_blocks_downstream). - sensor:
dagster sensor preview landing_sensorfirst 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. - schedule:
build_schedule_from_partitioned_job(hour_of_day=2)has cron0 2 * * *. a tick at 2025-01-05 02:00 utc targets partition2025-01-04. - pytest:
6 passed in 10.2 s(offline fixtures, a tmp duckdb, and a real dbt subprocess intest_full_graph_with_dbt). - dg cli:
uvx create-dagster@1.13.24 project dgproj --uv-sync, thendg scaffold defs dagster.asset assets/hello.py,dg check defs("all definitions loaded successfully"),dg list defs(1 asset), anddg launch --assets '*'(run_success).definitions.pyusesload_from_defs_folder. - 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. - 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
| need | pick | note |
|---|---|---|
| dbt in the graph | dagster-dbt @dbt_assets + DbtProject | map dbt sources to upstream assets with meta.dagster.asset_key (dbt) |
| duckdb | dagster-duckdb DuckDBResource | one writer; no read-only mode (#33802) |
| api ingestion | dagster-dlt | dlt itself in data-pipeline |
| external compute (spark, k8s, scripts) | pipes | covered by dagster's external-pipelines guide |
| scaffolding / project layout | create-dagster, dg, components | components |
| hosting | oss (webserver + daemon + dagster.yaml) or dagster+ | dagster's own deployment guides cover both shapes |
| alternatives | airflow 3 (assets now exist), prefect 3 | comparison in data-pipeline |