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
from dagster_dataframely_demo._data import storefront_orders
from dagster_polars import PolarsParquetIOManager
@dg.asset
def raw_orders() -> pl.DataFrame:
return storefront_orders()
quarantine=True is the whole setup. There is no directory to configure and no second asset to declare.
The asset’s own IO manager writes the invalid rows, under the asset’s key with _quarantine added to the end. It writes the quarantine of analytics/orders to analytics/orders_quarantine. This works with any IO manager that stores by asset key (ADR-0006).
The same declaration writes the quarantine to a different place under each IO manager:
@dd.asset(Orders, quarantine=True)
def orders(raw_orders: pl.DataFrame) -> pl.DataFrame:
return raw_orders
PolarsParquetIOManager |
orders_quarantine.parquet, in the same directory as orders.parquet |
DuckDBPolarsIOManager |
the table orders_quarantine, in the same schema as orders |
The first is a UPathIOManager and the second is a DbIOManager. Nearly every first-party IO manager subclasses one of the two.
The quarantine has the original columns, then one String rule column per rule, with valid, invalid or unknown in each row. A rule column takes its rule’s name, so it matches the check’s name only at the default rule granularity. At column or schema granularity, use the rule column, not a check, to find the rows that failed a rule. Check samples have at most max_failure_samples rows per rule, so the rule columns are the only complete record of which rules each row failed.
The table’s materialization has the quarantine address, the number of invalid rows, and which sets of rules they failed together. See What a run produces.
A partitioned asset on a database IO manager needs partition_expr. DbIOManager raises an error without it, so you already set it for the asset’s own table. The quarantine write uses the same value.
quarantine=True adds a context parameter to the asset, whether or not the decorated function declares one. The package needs the context to write the quarantine. So direct invocation of the asset needs a dg.build_asset_context() as the first argument. See Testing an asset. An asset without a quarantine keeps the signature you wrote.
The quarantine is not an asset
It is a record of a run. It is not in the asset graph, so no asset can depend on it until you add it with quarantine_spec:
defs = dg.Definitions(
assets=[raw_orders, orders, dd.quarantine_spec(Orders, orders)],
resources={"io_manager": PolarsParquetIOManager(base_dir="data/warehouse")},
)
quarantine_spec returns a dg.AssetSpec keyed orders_quarantine that depends on orders and has its own Columns tab. It has no compute function, because the decorator already writes the rows. The IO manager that wrote the rows also loads them for a downstream asset that names the quarantine as an input.
The quarantine’s Columns tab has no column constraints. It lists the schema’s columns with their dtypes, descriptions and tags, then one String column per rule. The invalid rows can fail any of the schema’s constraints, so a not null on a column full of nulls would be false. The primary key is absent too, because the writer writes rows with a duplicate key to the quarantine.
Pass quarantine_spec the asset definition, and the spec uses its partitions. Pass an asset key instead, and set the partitions with partitions_def=. Passing both an asset definition and partitions_def= raises dg.DagsterInvariantViolationError.
The quarantine never receives a materialization event, even with a spec. If you want to automate on invalid rows, see Automation.
quarantine_dir is for calling, not running
Direct invocation has no step, so there is no IO manager to write the invalid rows. Instead, file_writer writes them under the directory DAGSTER_DATAFRAMELY_QUARANTINE_DIR sets. A run does not use this directory.
DAGSTER_DATAFRAMELY_QUARANTINE_DIR=/tmp/quarantine
The package reads the directory when it writes the invalid rows, not when the module imports. So a monkeypatch.setenv in a test takes effect even though Python imported the module that defines the asset earlier. A call where every row is valid never reads it.
If the variable is unset, a call with invalid rows raises QuarantineDirError. There is no default directory.
There is no per-asset override.
Under that directory, the file path follows UPathIOManager’s layout, built from the asset key and the partition key:
<quarantine_dir>/<key part>/.../<name>_quarantine.parquet
<quarantine_dir>/<key part>/.../<name>_quarantine/<partition>.parquet
Each partition is a separate file in a <name>_quarantine/ directory, so a backfill of one partition rewrites one file. The file is always parquet, so it keeps every dtype the schema declares.
The package formats a multi-partition key as its dimension keys joined by /, in dimension-name order. It escapes a partition key containing .. or a leading /, instead of rejecting it, so no partition can write a file outside quarantine_dir.