asset()

Turn the decorated function into an asset that validates the frame it returns against schema.

Usage

Source

asset(
    schema,
    /,
    *,
    quarantine=False,
    check_granularity=None,
    schema_rules=None,
    max_failure_samples=None,
    statistics=None,
    row_sample=None,
    name=None,
    key_prefix=None,
    ins=None,
    deps=None,
    metadata=None,
    tags=None,
    description=None,
    config_schema=None,
    required_resource_keys=None,
    resource_defs=None,
    hooks=None,
    io_manager_key=None,
    partitions_def=None,
    op_tags=None,
    group_name=None,
    automation_condition=None,
    freshness_policy=None,
    backfill_policy=None,
    retry_policy=None,
    code_version=None,
    owners=None,
    kinds=None,
    pool=None
)

The decorated function returns a pl.DataFrame or pl.LazyFrame, a dg.MaterializeResult whose value is one of those, or None to skip.

Every keyword parameter not listed below is @dg.asset’s, passed to it unchanged.

Parameters

quarantine: bool = False

Whether a run writes the valid rows when some rows fail. With True, the run writes the invalid rows to the quarantine, under the asset key <name>_quarantine. A run fails with QuarantineKeyCollisionError if another asset already materializes that key. With False, any invalid row fails the run, and the asset writes nothing. True also adds a context parameter to the asset, whether or not the decorated function declares one.

check_granularity: Granularity | None = None

How many asset checks report the schema’s rules: one per rule at rule, one per column with rules at column, and one for the schema at schema. Changing it on an asset that has already run starts a new check history. None uses DAGSTER_DATAFRAMELY_CHECK_GRANULARITY if set, else rule.

schema_rules: SchemaRules | None = None

Which checks report the schema-level rules at column granularity: collapsed into one check, dy_schema__rules, or per_rule, one check each. None uses DAGSTER_DATAFRAMELY_SCHEMA_RULES if set, else collapsed.

max_failure_samples: int | None = None

A rule’s check metadata has at most this many rows that failed it, under dy_failed_sample. A run writes the rows to the Dagster event log unredacted. dy.Config.set_max_failure_examples does not change this number. None uses DAGSTER_DATAFRAMELY_MAX_FAILURE_SAMPLES if set, else 5.

statistics: bool | None = None

Whether the materialization metadata has statistics of the written rows, one table per dtype group. None uses DAGSTER_DATAFRAMELY_STATISTICS if set, else True.

row_sample: int | None = None

The materialization metadata has at most this many valid rows, and this many invalid rows. A run writes the rows to the Dagster event log unredacted. None uses DAGSTER_DATAFRAMELY_ROW_SAMPLE if set, else 5.

key_prefix: str | Sequence[str] | None = None

The quarantine’s asset key has the same prefix.

metadata: Mapping[str, Any] | None = None

Definition metadata. The schema’s dagster/column_schema entry replaces a key of the same name.

description: str | None = None

None uses the schema’s docstring, or the decorated function’s docstring if the schema has none.

io_manager_key: str | None = None

In a run, delegating_writer passes the invalid rows to the same IO manager.

partitions_def: dg.PartitionsDefinition[str] | None = None
The writer writes the quarantine under the same partition key.

Returns

A decorator that returns a dg.AssetsDefinition with the schema’s check specs and a Columns tab filled from the schema.

Raises

CollectionNotSupportedError

schema is a dy.Collection.

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.

InvalidSettingError
A setting’s argument or environment variable has a value the setting does not allow.

Examples

At the default rule granularity, the asset has one check per rule, plus the column-schema check:

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")


[spec.name for spec in orders.check_specs]
['dy_schema__columns',
 'dy_rule__primary_key',
 'dy_rule__order_id__nullability',
 'dy_rule__amount__nullability',
 'dy_rule__amount__min',
 'dy_rule__amount__inf',
 'dy_rule__amount__nan']

The check specs exist before the asset first runs, so the catalog lists every check before its first result.