quarantine_spec()

Return an asset spec that adds an asset’s quarantine to the asset graph.

Usage

Source

quarantine_spec(
    schema,
    asset,
    *,
    partitions_def=None,
)

The asset writes the quarantine whether or not the spec exists. The spec never receives a materialization event, because nothing materializes it.

Parameters

asset: dg.AssetsDefinition | dg.AssetKey | str | Sequence[str]

The asset the invalid rows came from, as its definition or its asset key. Pass a key only for an asset you cannot import. A definition also sets the spec’s partitions.

partitions_def: dg.PartitionsDefinition | None = None
The quarantine’s partitions, when asset is a key.

Returns

A spec keyed <name>_quarantine that depends on asset, to pass to dg.Definitions(assets=[...]).

Raises

ReservedColumnError

A column name is in the reserved dy_ namespace.

InvalidColumnNameError

A column name has a character Dagster does not allow in an asset check name.

CheckNameCollisionError

Two rules produce the same asset check name.

dg.DagsterInvariantViolationError
asset is a definition and the call also passes partitions_def.

Examples

class Orders(dy.Schema):
    order_id = dy.String(primary_key=True)
    amount = dy.Float64(nullable=False, min=0.0)


@dd.asset(Orders, quarantine=True)
def orders(raw_orders: pl.DataFrame) -> pl.DataFrame:
    return raw_orders.select("order_id", "amount")


dd.quarantine_spec(Orders, orders).key
AssetKey(['orders_quarantine'])

Pass the spec, not its key, to dg.Definitions(assets=[...]) with orders. A downstream asset can then declare the quarantine as an input.