Setup: the demo’s Orders schema, which every example on this page uses
import dagster as dg
import dataframely as dy
import polars as pl
import dagster_dataframely as dd
from dagster_dataframely_demo.schema import Orders
dd.asset passes partitions_def to dg.asset unchanged. Partitioning has no setting of its own. The writer writes the quarantine under the same partition key as the asset. Validation runs on each partition’s frame separately. dagster/row_count is the number of valid rows in that partition. If one partition’s frame fails the column-schema check, no other partition’s data changes.
A root asset reads its partition key from a declared context. It has no Dagster asset upstream, so it reads the partition’s file itself:
daily = dg.DailyPartitionsDefinition(start_date="2026-01-01")
@dd.asset(Orders, partitions_def=daily)
def orders(context: dg.AssetExecutionContext) -> pl.DataFrame:
day = context.partition_key # "2026-01-02"
raw = pl.read_parquet(f"raw/orders/{day}.parquet")
return raw
A partitioned @dg.asset reads its key the same way. In a direct invocation, only dg.build_asset_context(partition_key=...) can supply the key.
An asset downstream of it needs no partition key. An asset with the same partitions_def receives that partition’s rows as its parameter, because the IO manager loads only that partition:
@dd.asset(Orders, partitions_def=daily)
def priority_orders(orders: pl.DataFrame) -> pl.DataFrame:
# orders is one partition's rows, the same shape unpartitioned code sees
return orders.filter(pl.col("amount") > 100)
It is the same function you would write without partitioning.
An unpartitioned asset downstream of a partitioned one receives one frame per partition. The IO manager loads every partition and passes a dict keyed by partition key:
@dd.asset(Orders)
def rollup(orders: dict[str, pl.DataFrame]) -> pl.DataFrame:
# orders == {
# "2026-01-01": pl.DataFrame, # that day's rows
# "2026-01-02": pl.DataFrame,
# }
return pl.concat(orders.values())
Annotate the parameter as the dict, not as a frame. A pl.DataFrame annotation fails Dagster’s type check, and only after the IO manager has read every partition.
Annotate the values as pl.LazyFrame and the IO manager returns a scan for each partition:
@dd.asset(Orders)
def rollup(orders: dict[str, pl.LazyFrame]) -> pl.LazyFrame:
# orders == {
# "2026-01-01": pl.LazyFrame, # a scan of that day's file
# "2026-01-02": pl.LazyFrame,
# }
return pl.concat(orders.values())
pl.concat then builds one plan over every partition. Polars reads no rows until validation collects the plan. With a hundred partitions, the eager form holds all hundred in memory at once. The lazy form holds only the plan’s output. The dict’s keys and the validation are the same in both forms.
A MultiPartitionsDefinition needs nothing special either. The quarantine uses the same partitions as the asset:
grid = dg.MultiPartitionsDefinition({
"day": dg.DailyPartitionsDefinition(start_date="2026-01-01"),
"region": dg.StaticPartitionsDefinition(["eu", "us"]),
})
@dd.asset(Orders, partitions_def=grid)
def orders(context: dg.AssetExecutionContext) -> pl.DataFrame:
cell = context.partition_key.keys_by_dimension # {"day": ..., "region": ...}
raw = pl.read_parquet(f"raw/orders/{cell['region']}/{cell['day']}.parquet")
return raw
@dd.asset(Orders, partitions_def=grid, quarantine=True)
def priority_orders(
context: dg.AssetExecutionContext, orders: pl.DataFrame
) -> pl.DataFrame:
# orders is one cell's rows: one day, one region
return orders.filter(pl.col("amount") > 100)
As with one dimension, the downstream asset receives one frame, for the partition the run materializes, not a nested dict.
context.partition_key is a dg.MultiPartitionKey, a str subclass whose string form is 2026-01-01|eu. Read a dimension from keys_by_dimension instead of parsing that string. The string and the paths below order the dimensions by name, not in the order you declared them. So renaming a dimension can change the order.
A UPathIOManager writes one file per partition, with one path segment per dimension:
orders/2026-01-01/eu.parquet
orders/2026-01-01/us.parquet
priority_orders/2026-01-01/eu.parquet
priority_orders_quarantine/2026-01-01/eu.parquet
An unpartitioned asset downstream of a multi-partitioned one receives a flat dict, not a nested one. It has one entry per partition, keyed by the multi-partition key:
@dd.asset(Orders)
def eu_orders(orders: dict[dg.MultiPartitionKey, pl.LazyFrame]) -> pl.LazyFrame:
# orders == {
# "2026-01-01|eu": pl.LazyFrame, # one cell, one scan
# "2026-01-01|us": pl.LazyFrame,
# "2026-01-02|eu": pl.LazyFrame,
# "2026-01-02|us": pl.LazyFrame,
# }
return pl.concat(
frame
for key, frame in orders.items()
if key.keys_by_dimension["region"] == "eu"
)
Every key is a MultiPartitionKey at runtime, so keys_by_dimension works on it. Annotate the keys as dg.MultiPartitionKey, as above, so a type checker accepts that read. dict[str, pl.LazyFrame] also works, but a type checker then needs a cast before each read.
To depend on one dimension only, use a partition mapping. An asset partitioned by day alone can depend on the multi-partitioned asset through dg.MultiToSingleDimensionPartitionMapping(partition_dimension_name="day"). It then receives only that day’s regions: two entries instead of four, still keyed 2026-01-02|eu and 2026-01-02|us.
A partition with no data
Some partitions have no source data and never will. For example, take a monthly-by-distributor grid where one distributor left the marketplace two years ago. Its old partitions have real data and must stay, and its recent ones have no source file.
Such a partition is neither a failure nor an empty table. Return None and the asset skips: the package validates nothing, nothing materializes, and the run succeeds. The partition stays unmaterialized, instead of materializing zero rows or failing.
from pathlib import Path
class SupplierReport(dy.Schema):
supplier_id = dy.String(primary_key=True)
reported_on = dy.Date(nullable=False)
monthly = dg.MultiPartitionsDefinition({
"month": dg.MonthlyPartitionsDefinition(start_date="2026-01-01"),
"distributor": dg.StaticPartitionsDefinition(["acme", "globex"]),
})
def source_path(key: dg.MultiPartitionKey) -> Path:
cell = key.keys_by_dimension
return Path(f"raw/supplier_reports/{cell['distributor']}/{cell['month']}.parquet")
@dd.asset(SupplierReport, partitions_def=monthly)
def supplier_reports(context: dg.AssetExecutionContext) -> pl.DataFrame | None:
path = source_path(context.partition_key)
if not path.exists():
return None
return pl.read_parquet(path)
You write the check for whether the source exists. The decorator never catches FileNotFoundError, or any other exception, to skip for you. A file that is absent on purpose and a misconfigured path raise the same error, so the decorator cannot distinguish them. Return None for the first, and let the error propagate for the second.
A forgotten return also returns None, so it skips the asset. The run succeeds, and the missing materialization is the only sign of the mistake.
Every check still reports on a skipped run, and passes. Dagster requires a result for every check spec, so the package runs the rules over Schema.create_empty(). Each rule then runs over zero rows, so none fails. Dagster records these results with no target_materialization_data, so it never links a passing check on a skipped partition to the last materialization that had data.
Automation
dd.asset passes automation_condition and freshness_policy to dg.asset unchanged, and they apply to the asset only.
You cannot automate on the quarantine. The writer writes the invalid rows inside the asset’s own step. The quarantine has no materialization events. So dg.AutomationCondition.eager() on a downstream asset never requests a run because of the quarantine, however you declare the dependency.
Automate on the check instead. Each failed rule is a failed asset check with its own history. You can also add your own asset check to the asset:
@dg.asset_check(asset=orders, name="triage_needed")
def triage_needed() -> dg.AssetCheckResult: ...
Or read the quarantine on a schedule, through the spec quarantine_spec returns.