Turn the decorated function into an asset that validates the frame it returns against schema.
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.