# Partitioning


Setup: the demo's `Orders` schema, which every example on this page uses

``` python
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:


``` python
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:


``` python
@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:


``` python
@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:**


``` python
@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:


``` python
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.


<figure class="figure">
<p><img src="../assets/images/partitions-grid.png" class="img-fluid figure-img" /></p>
<figcaption>The Partitions tab of an asset partitioned by day and by region, with the partition <code>2026-08-01|us</code> selected. Validation ran per partition, so the row count and the statistics on the right are for that partition only.</figcaption>
</figure>


`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:

``` text
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:


``` python
@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.


``` python
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:


``` python
@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](../reference/quarantine_spec.md#dagster_dataframely.quarantine_spec) returns.
