---------------------------------------------------------------------- This is the API documentation for the dagster_dataframely library. ---------------------------------------------------------------------- ## Core `dd.asset` builds an asset from a schema, and `quarantine_spec` adds that asset's quarantine to the asset graph. asset(schema: type[dataframely.schema.Schema], /, *, quarantine: bool = False, check_granularity: Granularity | None = None, schema_rules: SchemaRules | None = None, max_failure_samples: int | None = None, statistics: bool | None = None, row_sample: int | None = None, name: str | None = None, key_prefix: str | collections.abc.Sequence[str] | None = None, ins: collections.abc.Mapping[str, dagster._core.definitions.assets.job.asset_in.AssetIn] | None = None, deps: collections.abc.Iterable[typing.Union[dagster._core.definitions.asset_key.AssetKey, str, collections.abc.Sequence[str], ForwardRef('AssetSpec'), ForwardRef('AssetsDefinition'), ForwardRef('SourceAsset'), ForwardRef('AssetDep')]] | None = None, metadata: collections.abc.Mapping[str, typing.Any] | None = None, tags: collections.abc.Mapping[str, str] | None = None, description: str | None = None, config_schema: collections.abc.Mapping[str, typing.Any] | None = None, required_resource_keys: collections.abc.Set[str] | None = None, resource_defs: collections.abc.Mapping[str, object] | None = None, hooks: collections.abc.Set[dagster._core.definitions.hook_definition.HookDefinition] | None = None, io_manager_key: str | None = None, partitions_def: Optional[dagster._core.definitions.partitions.definition.partitions_definition.PartitionsDefinition[str]] = None, op_tags: collections.abc.Mapping[str, typing.Any] | None = None, group_name: str | None = None, automation_condition: Union[dagster._core.definitions.declarative_automation.automation_condition.AutomationCondition[dagster._core.definitions.asset_key.AssetKey], dagster._core.definitions.declarative_automation.automation_condition.AutomationCondition[dagster._core.definitions.asset_key.AssetKey | dagster._core.definitions.asset_key.AssetCheckKey], NoneType] = None, freshness_policy: dagster._core.definitions.freshness.FreshnessPolicy | None = None, backfill_policy: dagster._core.definitions.backfill_policy.BackfillPolicy | None = None, retry_policy: dagster._core.definitions.policy.RetryPolicy | None = None, code_version: str | None = None, owners: collections.abc.Sequence[str] | None = None, kinds: collections.abc.Set[str] | None = None, pool: str | None = None) -> collections.abc.Callable[[collections.abc.Callable[..., typing.Union[polars.dataframe.frame.DataFrame, polars.lazyframe.frame.LazyFrame, dagster._core.definitions.result.MaterializeResult[polars.dataframe.frame.DataFrame], dagster._core.definitions.result.MaterializeResult[polars.lazyframe.frame.LazyFrame], NoneType]]], dagster._core.definitions.assets.definition.assets_definition.AssetsDefinition] Turn the decorated function into an asset that validates the frame it returns against `schema`. 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 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 `_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 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 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 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 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 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 The quarantine's asset key has the same prefix. metadata Definition metadata. The schema's `dagster/column_schema` entry replaces a key of the same name. description `None` uses the schema's docstring, or the decorated function's docstring if the schema has none. io_manager_key In a run, `delegating_writer` passes the invalid rows to the same IO manager. partitions_def 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 -------- ```python #| echo: false #| output: false import dataframely as dy import polars as pl import dagster_dataframely as dd ``` At the default `rule` granularity, the asset has one check per rule, plus the column-schema check: ```python 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] ``` The check specs exist before the asset first runs, so the catalog lists every check before its first result. quarantine_spec(schema: type[dataframely.schema.Schema], asset: dagster._core.definitions.assets.definition.assets_definition.AssetsDefinition | dagster._core.definitions.asset_key.AssetKey | str | collections.abc.Sequence[str], *, partitions_def: dagster._core.definitions.partitions.definition.partitions_definition.PartitionsDefinition | None = None) -> dagster._core.definitions.assets.definition.asset_spec.AssetSpec Return an asset spec that adds an asset's quarantine to the asset graph. The asset writes the quarantine whether or not the spec exists. The spec never receives a materialization event, because nothing materializes it. Parameters ---------- asset 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 The quarantine's partitions, when `asset` is a key. Returns ------- A spec keyed `_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 -------- ```python #| echo: false #| output: false import dataframely as dy import polars as pl import dagster_dataframely as dd ``` ```python 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 ``` Pass the spec, not its key, to `dg.Definitions(assets=[...])` with `orders`. A downstream asset can then declare the quarantine as an input. Granularity Type alias. Type aliases are created through the type statement:: type Alias = int In this example, Alias and int will be treated equivalently by static type checkers. At runtime, Alias is an instance of TypeAliasType. The __name__ attribute holds the name of the type alias. The value of the type alias is stored in the __value__ attribute. It is evaluated lazily, so the value is computed only if the attribute is accessed. Type aliases can also be generic:: type ListOrSet[T] = list[T] | set[T] In this case, the type parameters of the alias are stored in the __type_params__ attribute. See PEP 695 for more information. SchemaRules Type alias. Type aliases are created through the type statement:: type Alias = int In this example, Alias and int will be treated equivalently by static type checkers. At runtime, Alias is an instance of TypeAliasType. The __name__ attribute holds the name of the type alias. The value of the type alias is stored in the __value__ attribute. It is evaluated lazily, so the value is computed only if the attribute is accessed. Type aliases can also be generic:: type ListOrSet[T] = list[T] | set[T] In this case, the type parameters of the alias are stored in the __type_params__ attribute. See PEP 695 for more information. ## Hand-wiring `dd.asset` assembles these parts. Use them to build a `@dg.asset` yourself, where `dd.asset` does not fit. schema_metadata(schema: type[dataframely.schema.Schema]) -> dict[str, dagster._core.definitions.metadata.table.TableSchema] Return the definition metadata that fills an asset's Columns tab from the schema. Returns ------- A one-entry mapping from `dagster/column_schema` to `table_schema(schema)`, to pass to `dg.asset(metadata=...)`. 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. table_schema(schema: type[dataframely.schema.Schema]) -> dagster._core.definitions.metadata.table.TableSchema Return the schema as the `dg.TableSchema` that Dagster's Columns tab shows. Each column has its dtype, description, nullability, uniqueness and column constraints. Its tags come from `dy.Column(metadata=...)`, with each value converted to a string. The primary key is a table constraint, so only `unique=True` marks a key column unique. Returns ------- A table schema with the columns in the schema's order. 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. check_specs(schema: type[dataframely.schema.Schema], *, asset: str | dagster._core.definitions.asset_key.AssetKey, check_granularity: Granularity | None = None, schema_rules: SchemaRules | None = None) -> list[dagster._core.definitions.asset_checks.asset_check_spec.AssetCheckSpec] Return the asset check specs for the schema's rules, plus the column-schema check. Parameters ---------- check_granularity How many checks report the schema's rules. 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 Which checks report the schema-level rules at `column` granularity. `None` uses `DAGSTER_DATAFRAMELY_SCHEMA_RULES` if set, else `collapsed`. Returns ------- The column-schema check's spec first, then one spec per rule set. Raises ------ InvalidSettingError A setting's argument or environment variable has a value the setting does not allow. 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. check_results(schema: type[dataframely.schema.Schema], frame: polars.dataframe.frame.DataFrame | polars.lazyframe.frame.LazyFrame, *, asset_key: dagster._core.definitions.asset_key.AssetKey, severity: dagster._core.definitions.asset_checks.asset_check_spec.AssetCheckSeverity, check_granularity: Granularity | None = None, schema_rules: SchemaRules | None = None, max_failure_samples: int | None = None) -> collections.abc.Iterator[dagster._core.definitions.asset_checks.asset_check_result.AssetCheckResult] Yield a result for each check `check_specs` declares, without writing any rows. Parameters ---------- severity The severity of every failing check except the column-schema check, which is always `ERROR`. check_granularity Pass the value `check_specs` received, or the results name checks the asset does not declare. `None` uses `DAGSTER_DATAFRAMELY_CHECK_GRANULARITY` if set, else `rule`. schema_rules Pass the value `check_specs` received, as for `check_granularity`. `None` uses `DAGSTER_DATAFRAMELY_SCHEMA_RULES` if set, else `collapsed`. max_failure_samples A rule's check metadata has at most this many rows that failed it. `None` uses `DAGSTER_DATAFRAMELY_MAX_FAILURE_SAMPLES` if set, else `5`. Yields ------ The column-schema check's result, then one result per rule set, in the order of `check_specs`. Raises ------ InvalidSettingError A setting's argument or environment variable has a value the setting does not allow. 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. ColumnSchemaError The frame's columns or dtypes do not match the schema. It yields the column-schema check's failing result first. check_name(rule_name: str) -> str Return the asset check name for the Dataframely rule named `rule_name`. The quarantine's rule columns and the checks at `rule` granularity have this name. Returns ------- `dy_rule__` followed by `rule_name` with `|` replaced by `__`: `amount|min` becomes `dy_rule__amount__min`. validation_results(schema: type[dataframely.schema.Schema], frame: polars.dataframe.frame.DataFrame | polars.lazyframe.frame.LazyFrame | None, *, valid_key: dagster._core.definitions.asset_key.AssetKey, quarantine_writer: QuarantineWriter | None = None, check_granularity: Granularity | None = None, schema_rules: SchemaRules | None = None, max_failure_samples: int | None = None, statistics: bool | None = None, row_sample: int | None = None) -> collections.abc.Iterator[typing.Union[dagster._core.definitions.result.MaterializeResult[polars.dataframe.frame.DataFrame], dagster._core.definitions.asset_checks.asset_check_result.AssetCheckResult]] Validate a frame and yield the asset's materialization and check results. Parameters ---------- frame The frame to validate, or `None` to skip. valid_key The key the run materializes the valid rows under: the asset's own key, such as `context.asset_key`. A key the asset does not have fails the step with `DagsterInvariantViolationError` on the first yield, and the asset writes nothing. quarantine_writer The writer for the invalid rows, called at most once, with the frame `quarantine_frame` returns. With `None`, any invalid row fails the run with `ValidationAbortError`, and the asset writes nothing. check_granularity Pass the value `check_specs` received, or the results name checks the asset does not declare. `None` uses `DAGSTER_DATAFRAMELY_CHECK_GRANULARITY` if set, else `rule`. schema_rules Pass the value `check_specs` received, as for `check_granularity`. `None` uses `DAGSTER_DATAFRAMELY_SCHEMA_RULES` if set, else `collapsed`. max_failure_samples A rule's check metadata has at most this many rows that failed it. `None` uses `DAGSTER_DATAFRAMELY_MAX_FAILURE_SAMPLES` if set, else `5`. statistics Whether the materialization metadata has statistics of the written rows. `None` uses `DAGSTER_DATAFRAMELY_STATISTICS` if set, else `True`. row_sample The materialization metadata has at most this many valid rows, and this many invalid rows. `None` uses `DAGSTER_DATAFRAMELY_ROW_SAMPLE` if set, else `5`. Yields ------ The asset's `dg.MaterializeResult`, if the run writes its table, then one `dg.AssetCheckResult` per check. On a skip, it yields only the check results, and every check passes. It writes the quarantine instead of yielding it. 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. InvalidSettingError A setting's argument or environment variable has a value the setting does not allow. DagsterInvariantViolationError `frame` is neither a Polars frame nor `None`. ColumnSchemaError The frame's columns or dtypes do not match the schema. It yields the column-schema check's failing result first, and writes nothing. ValidationAbortError Rows failed validation and `quarantine_writer` is `None`, so the package writes nothing. NoValidRowsError Every row failed validation. It writes the quarantine, and not the asset's table. Iterator(*args, **kwargs) QuarantineWriter Type alias. Type aliases are created through the type statement:: type Alias = int In this example, Alias and int will be treated equivalently by static type checkers. At runtime, Alias is an instance of TypeAliasType. The __name__ attribute holds the name of the type alias. The value of the type alias is stored in the __value__ attribute. It is evaluated lazily, so the value is computed only if the attribute is accessed. Type aliases can also be generic:: type ListOrSet[T] = list[T] | set[T] In this case, the type parameters of the alias are stored in the __type_params__ attribute. See PEP 695 for more information. delegating_writer(context: dagster._core.execution.context.asset_execution_context.AssetExecutionContext) -> QuarantineWriter Return a writer that passes the invalid rows to the asset's own IO manager. The manager stores the rows under the asset key `_quarantine`, as it stores any other asset: `PolarsParquetIOManager` writes a parquet file beside the table, and `DuckDBPolarsIOManager` writes a table beside it. The writer discards the metadata the manager emits for this write. Parameters ---------- context Its asset definition must have exactly one asset, as a `@dg.asset` does. Returns ------- A writer that takes the invalid rows and returns the quarantine's asset key as a string. Raises ------ dagster._core.errors.DagsterInvalidPropertyError There is no step, which means a direct invocation. Use `file_writer` there. file_writer(key: dagster._core.definitions.asset_key.AssetKey, quarantine_dir: upath.core.UPath | pathlib.Path | str, partition_key: str | None = None) -> QuarantineWriter Return a writer that writes the invalid rows to a parquet file under `quarantine_dir`. Use it in a direct invocation, where `delegating_writer` raises because there is no step. The file is at `quarantine_path(key, quarantine_dir, partition_key)`. Parameters ---------- key The asset key of the asset whose rows failed validation, not the quarantine's key. Returns ------- A writer that takes the invalid rows and returns the file's path as a string. quarantine_frame(schema: type[dataframely.schema.Schema], failure: dataframely.filter_result.FailureInfo) -> polars.dataframe.frame.DataFrame Return the invalid rows as a writer receives them, with one rule column per rule. Returns ------- The original columns in their own order, then a `String` rule column for every rule, named by `check_name`. A rule column's value is `valid`, `invalid` or `unknown`. 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. quarantine_path(key: dagster._core.definitions.asset_key.AssetKey, quarantine_dir: upath.core.UPath | pathlib.Path | str, partition_key: str | None = None) -> upath.core.UPath Return the path of the parquet file `file_writer` writes an asset's invalid rows to. The layout matches `UPathIOManager`'s, so with `quarantine_dir` set to a manager's `base_dir`, the file is beside the asset's own file. Parameters ---------- key The asset key of the asset whose rows failed validation, not the quarantine's key. quarantine_dir Any path `UPath` accepts, such as a local directory or `s3://bucket/prefix`. Credentials for a remote path come from the environment. partition_key A `dg.MultiPartitionKey` becomes its dimension keys joined by `/`, in dimension-name order. Escaping a `..` and a leading `/` keeps the file under `quarantine_dir`. Returns ------- `//.../_quarantine.parquet`, or `//.../_quarantine/.parquet` for a partition. `quarantine_path` creates no parent directories. validate_quarantine_key(context: dagster._core.execution.context.asset_execution_context.AssetExecutionContext) -> None Raise if another asset already materializes the quarantine's asset key. Call it at the start of a hand-wired asset's function when the asset writes a quarantine, as `dd.asset` does. In a direct invocation it does nothing. Raises ------ QuarantineKeyCollisionError Another asset in the code location materializes `_quarantine`. ## Errors This package raises these errors. Each one subclasses `DagsterDataframelyError`. DagsterDataframelyError Base class for every error this package raises. CollectionNotSupportedError(collection_name: str) -> None `schema=` received a `dy.Collection`. Raised at decoration time. ReservedColumnError(schema_name: str, columns: list[str]) -> None A column name is in the reserved `dy_` namespace. Raised at definition time, because the column would share a name with a check or a rule column this package generates. InvalidColumnNameError(schema_name: str, columns: list[str]) -> None A column name has a character Dagster does not allow in an asset check name. Raised at definition time (ADR-0008). Dagster allows only `A-Za-z0-9_`, and this package builds one check name per rule from the column name. CheckNameCollisionError(schema_name: str, first: str, second: str, name: str) -> None Two rules produce the same asset check name. Raised at definition time, before Dagster's own duplicate-check error, which does not name the rules. InvalidSettingError(setting: str, value: str, allowed: collections.abc.Sequence[str] | str, *, source: str, env_var: str, takes_argument: bool = True) -> None A setting has a value it does not allow. Raised when the setting resolves, so the message names the source of the value. MaterializeResultValueError(asset: str) -> None A returned `dg.MaterializeResult` has no frame in `value`. Raised before the column-schema check, because there is no frame to check. MaterializeResultFieldError(asset: str, field: str) -> None A returned `dg.MaterializeResult` sets `asset_key` or `check_results`. Raised before the column-schema check. The decorator sets both fields itself. ColumnSchemaError(schema_name: str, problems: collections.abc.Sequence[collections.abc.Mapping[str, str]]) -> None A frame's columns or dtypes do not match the schema. This is a bug in the pipeline, not bad data, so the run fails before `Schema.filter` runs, and the asset writes nothing. ValidationAbortError(schema_name: str, invalid_count: int, counts: collections.abc.Mapping[str, int]) -> None Rows failed validation and the asset declares no quarantine, so it writes nothing. Without a quarantine, every row has to be valid. No setting drops the invalid rows and writes the rest: to drop rows, filter them in the asset body. NoValidRowsError(schema_name: str, invalid_count: int, counts: collections.abc.Mapping[str, int], address: str) -> None Every row failed validation, so the run wrote only the quarantine. Nothing writes the asset's table, so an empty table never replaces the last-known-good one. QuarantineKeyCollisionError(asset: str, quarantine: str) -> None Another asset already materializes the quarantine's asset key. Raised before the decorated function runs, on every run of an asset with `quarantine=True` (ADR-0007). Otherwise both assets would write to the same key, and the second write would replace the first. QuarantineDirError(asset: str) -> None Invalid rows need writing, and nothing sets `quarantine_dir`. Raised only when you call an asset with `quarantine=True` directly, because a run always has an IO manager to write the rows. There is no default directory, so nothing writes the rows somewhere nobody chose. ---------------------------------------------------------------------- This is the User Guide documentation for the package. ---------------------------------------------------------------------- ## Guide ### Declaring an asset ```{python} #| label: setup-declaring-an-asset #| code-fold: true #| code-summary: "Setup: the demo's `Orders` schema, which every example on this page uses" #| output: false import dagster as dg import dataframely as dy import polars as pl import dagster_dataframely as dd from dagster_dataframely_demo.schema import Orders ``` The schema is the only required argument, and it is positional. Every other argument is `quarantine`, one of the [Settings](settings.qmd), or a `@dg.asset` parameter under the same name. You write the decorated function as you would for `@dg.asset`: ```{python} #| label: declaring-an-asset-1 @dd.asset(Orders, quarantine=True, group_name="sales") def orders(context: dg.AssetExecutionContext, raw_orders: pl.DataFrame) -> pl.DataFrame: context.log.info("%d rows arrived", raw_orders.height) return raw_orders ``` Dagster passes upstream assets to it by parameter name, as it does for a plain `@dg.asset`. Use `ins=` and `deps=` where a parameter name is not enough. Declare a `context` parameter to use the partition key, the log, resources or configuration. This package does its work at two times: when the module imports, and when the asset runs. ```{mermaid} %%| label: declaring-an-asset-diagram flowchart TB subgraph load["when the module imports"] D["@dd.asset(Orders)"] --> T["Columns tab"] D --> C["a check spec per rule set"] end subgraph run["when the asset runs"] direction TB F["decorated function returns a frame"] --> CS["column-schema check"] CS --> FL["Schema.filter"] FL --> V["valid rows, to the asset's IO manager"] FL --> I["invalid rows, to the quarantine"] FL --> R["a result per check"] end load --> run ``` Two things exist as soon as the module imports, before any run: - **The Columns tab**, built from the schema. The tab shows dtypes, descriptions, nullability, uniqueness, the primary key at table level, and every other column constraint beside its column. - **A check spec per rule set**, so the catalog lists the asset's checks by name before the first run. The asset's description comes from the schema's docstring, as [Naming](naming.qmd) describes. A run checks the column schema of the frame the decorated function returns, then filters the frame with `Schema.filter`. The asset's IO manager writes the valid rows. A writer writes the invalid rows to the quarantine. Each check reports a result. Not every run writes both the valid and the invalid rows: [The failure policy](the-failure-policy.qmd) lists the six outcomes and what each one writes. ![The Checks tab after a run with no failing rows, at the default granularity: one check per rule, each with its own history. The selected check's description is its rendered constraint, and its metadata has the rule's name and expression.](../assets/images/check-list.png) `dd.asset` builds a `@dg.asset` and accepts its parameters under the same names, except the six below. A test fails if `@dg.asset` adds a parameter that `dd.asset` lacks, or removes one that `dd.asset` passes on. `config_schema` is narrower than Dagster's: it takes only a mapping, not the union of six types that `@dg.asset` accepts. A type checker rejects the five legacy forms, but nothing rejects them at run time. | `@dg.asset` parameter | why it is not here | what to write instead | | --- | --- | --- | | `check_specs` | they come from the schema's rules | nothing | | `key` | `key_prefix` and `name` set it | `key_prefix=` plus `name=` | | `output_required` | the outcomes that write nothing must end the step without an output | nothing | | `dagster_type` | declined ([#3](https://github.com/ozanozbeker/dagster-dataframely/issues/3)): it runs before the IO manager and cannot set a severity, so it cannot apply the failure policy | the schema, which this package validates against | | `is_virtual` | a virtual asset has no function to decorate | nothing | | `io_manager_def` | `resource_defs` already sets it | `resource_defs={"io_manager": ...}` | ## The schema is a single `dy.Schema`, never a `dy.Collection` Passing a Collection raises `CollectionNotSupportedError` at decoration time. Declare one asset per member instead, each with the member's own schema. The parts under [Hand-wiring](hand-wiring.qmd) cannot validate a Collection either, because `validation_results` takes a single schema. ## What the decorated function returns It can return five things: ```text pl.DataFrame dg.MaterializeResult[pl.DataFrame] pl.LazyFrame dg.MaterializeResult[pl.LazyFrame] None ``` **The returned value determines what happens, not the annotation.** The same `Schema.filter` call filters a `LazyFrame` and a `DataFrame`. A `dg.MaterializeResult` is unwrapped first. `@dg.asset` enforces its return annotation, because it infers the output's `dagster_type` from it. `dd.asset` does not. The asset always stores a `DataFrame`, because this package collects the valid rows even when the function returns a `LazyFrame`. Annotate the function anyway, so that a type checker enforces the annotation. Parameterize a returned result, because a strict type checker rejects a bare `dg.MaterializeResult` as an implicit `Any`. Return a `dg.MaterializeResult` instead of a bare frame when you want to add metadata, tags or a data version to the materialization. Prefer it to `context.add_asset_metadata`: it is the only one of the two that works when you call the asset directly. [Attaching your own metadata](what-a-run-produces.qmd#attaching-your-own-metadata) describes both. Return `None` to skip, as [A partition with no data](partitioning.qmd#a-partition-with-no-data) describes. To skip, return `None` itself: `dg.MaterializeResult(value=None)` raises `MaterializeResultValueError`. Return anything else and the run fails with `dg.DagsterInvariantViolationError` before the column-schema check. ### The failure policy ```{python} #| label: setup-the-failure-policy #| code-fold: true #| code-summary: "Setup: the demo's `Orders` schema, which every example on this page uses" #| output: false import dagster as dg import dataframely as dy import polars as pl import dagster_dataframely as dd from dagster_dataframely_demo.schema import Orders ``` Whether the asset has `quarantine=True` determines what a run does with invalid rows. There is no lenient mode and no strict flag. | what the function returned | the table | the quarantine | checks | run | | --- | --- | --- | --- | --- | | columns or dtypes that are not the schema's | not written | not written | the column-schema check fails, blocking | fails, `ColumnSchemaError` | | every row valid | written | not written | all pass | succeeds | | some rows failed, no quarantine declared | not written | n/a | fail at `ERROR` | fails, `ValidationAbortError` | | some rows failed, quarantine declared | the valid rows | the invalid rows | fail at `WARN` | succeeds | | every row failed, quarantine declared | not written | every row | fail at `ERROR` | fails, `NoValidRowsError` | | `None`, meaning no source data | skipped | not written | all pass | succeeds | Each row of the table is one of the six outcomes a run can have. The outcome, not the rule, determines the severity of the rules' checks, so they all have the same severity in a run. When the column-schema check fails, its severity is always `ERROR`. ![The Checks tab of an asset with no quarantine declared, after some rows failed. The failing checks have severity `ERROR`, not the `WARN` they have when the asset declares a quarantine. That is the difference between the third and fourth rows of the table.](../assets/images/severity-error.png) ![The run log shows the step failure. The `ValidationAbortError` message lists how many rows failed each rule. The counts add up to more than the number of invalid rows, because one row can fail several rules.](../assets/images/error-validation-abort.png) **Without a quarantine, every row has to be valid.** A run with even one failing row writes nothing, so the last-known-good table is unchanged. No setting discards the invalid rows and writes the rest. To drop invalid rows, drop them yourself in the decorated function: ```{python} #| label: the-failure-policy-1 @dd.asset(Orders) def orders(raw_orders: pl.DataFrame) -> pl.DataFrame: valid, _ = Orders.filter(raw_orders) return valid ``` **With a quarantine, a run where every row fails still fails.** It writes no table rather than an empty one, so an empty table never replaces the last-known-good table. It writes every row to the quarantine before it raises `NoValidRowsError`, so you can read the rows after the run fails. ## The package never casts The column-schema check compares your frame's dtypes with the schema's and fails the run on a mismatch. `Schema.filter` runs with `cast=False`. A dtype mismatch is a bug in the pipeline, and a silent cast can write wrong values. For example, a cast from `Float64` to a declared `Int64` turns `1.9` into `1`, and every rule still passes. The check runs before `Schema.filter`, without executing the plan. It reports every mismatched column at once. **Extra columns and column order are not a mismatch.** The check compares only the columns the schema declares, and `Schema.filter` drops any other column. The table has the schema's columns in the schema's order, whatever order the frame had: ```{python} #| label: the-failure-policy-2 @dd.asset(Orders) def orders(raw_orders: pl.DataFrame) -> pl.DataFrame: # `net` is working state, and `raw_orders` carries columns `Orders` never declared. # Both are dropped. What materializes is the schema's thirteen columns, in its order. return raw_orders.with_columns(net=pl.col("amount") * 0.8) ``` So a decorated function that returns extra columns needs no `select`, and no `Schema.cast` to drop them. **To cast, call `Schema.cast` in the decorated function:** ```{python} #| label: the-failure-policy-3 @dd.asset(Orders) def orders(raw_orders: pl.DataFrame) -> pl.DataFrame: return Orders.cast(raw_orders) ``` This is the intended way to cast, not a workaround, and few assets need it. An asset needs it when its source returns a dtype that differs from the schema's, most often on a read from a warehouse. DuckDB returns `DECIMAL` where the schema declares `Float64`, so an asset that reads from `DuckDBPolarsIOManager` fails the column-schema check on its first run unless it casts. If you cast in every asset, check whether the dtypes differ at all, because dropping extra columns never needs a cast ([#88](https://github.com/ozanozbeker/dagster-dataframely/issues/88)). The package casts only columns it generates. It casts the quarantine's rule columns from `Enum` to `String`, because an `Enum` column makes the Delta writer panic with a Rust `unreachable!()`. ### Where invalid rows go ```{python} #| label: setup-where-invalid-rows-go #| code-fold: true #| code-summary: "Setup: the demo's `Orders` schema, which every example on this page uses" #| output: false 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](https://github.com/ozanozbeker/dagster-dataframely/blob/main/docs/pre-1.0.md#adr-0006-the-assets-own-io-manager-writes-the-quarantine)). The same declaration writes the quarantine to a different place under each IO manager: ```{python} #| label: where-invalid-rows-go-1 @dd.asset(Orders, quarantine=True) def orders(raw_orders: pl.DataFrame) -> pl.DataFrame: return raw_orders ``` | IO manager | where it writes the invalid rows | | --- | --- | | `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](what-a-run-produces.qmd). **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](testing-an-asset.qmd). 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`: ```{python} #| label: where-invalid-rows-go-2 defs = dg.Definitions( assets=[raw_orders, orders, dd.quarantine_spec(Orders, orders)], resources={"io_manager": PolarsParquetIOManager(base_dir="data/warehouse")}, ) ``` ![The lineage view shows the quarantine's node. Its one edge comes from `marketing_orders`, the asset with `quarantine=True`, not from the upstream `raw_marketplace_orders`. The node has a dashed border because nothing materializes it.](../assets/images/quarantine-lineage.png) `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](partitioning.qmd#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. ```bash 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: ```text //.../_quarantine.parquet //.../_quarantine/.parquet ``` Each partition is a separate file in a `_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`. ### What a run produces ```{python} #| label: setup-what-a-run-produces #| code-fold: true #| code-summary: "Setup: the demo's `Orders` schema, which every example on this page uses" #| output: false import dagster as dg import dataframely as dy import polars as pl import dagster_dataframely as dd from dagster_dataframely_demo.schema import Orders ``` You launch a run as you would for any Dagster asset: from the UI, from a schedule or a sensor, or with `dg.materialize([orders], resources={...})`. This package reports a run's results in two places: the table's materialization and the asset checks. ## The table's materialization The row count is under Dagster's own key, `dagster/row_count`. Every other key this package writes is in the `dataframely/` namespace, so these keys sort together, apart from Dagster's keys and your IO manager's. | key | what it holds | | --- | --- | | `dagster/row_count` | how many rows passed validation, not how many the function returned | | `dataframely/valid_sample` | the first few of those rows | | `dataframely/valid_statistics/` | one table per dtype group present | | `dataframely/quarantine_address` | where the writer wrote the invalid rows, on the runs that wrote any | | `dataframely/invalid_count` | how many rows failed at least one rule | | `dataframely/invalid_sample` | the first few of those, rule columns included | | `dataframely/invalid_by_rules` | which sets of rules the rows failed together, largest group first | The last four are absent when no rows failed. `dataframely/invalid_by_rules` counts each invalid row once, under the set of rules it failed. So one bad upstream field that fails three rules appears as one group, not as three unrelated counts. It names rules as the quarantine's rule columns do, so its names match the check names only at `rule` granularity. ![The materialization of a run where 8 rows failed, and the writer wrote them to the quarantine. The third row of `dataframely/invalid_by_rules` is one row that failed three rules, counted once, not three times.](../assets/images/invalid-by-rules.png) Every value that contains rows is a Dagster table value, not markdown, so the UI shows a sortable table. The two counts are integers, and the quarantine address is a string. ![The materialization of a run where no rows failed: the row count, one statistics table per dtype group present, the row sample, and the IO manager's own keys below them.](../assets/images/materialization-metadata.png) ## The statistics tables There is one table for each dtype group present in the written rows. Each table has one row per column, in the frame's column order. A column whose dtype is in no group, such as a `List`, `Struct` or `Array`, appears in no table, because only a count and a null count would apply to it. | group | dtypes | columns | | --- | --- | --- | | `numeric` | `Int*`, `UInt*`, `Float*`, `Decimal` | `count`, `null_count`, `mean`, `std`, `min`, `p50`, `max` | | `temporal` | `Date`, `Datetime`, `Time`, `Duration` | `count`, `null_count`, `min`, `max`, `span` | | `string` | `String`, `Categorical`, `Enum`, `Binary` | `count`, `null_count`, `n_unique`, `min_len`, `max_len`, `n_empty` | | `boolean` | `Boolean` | `count`, `null_count`, `n_true`, `n_false`, `true_rate` | The package computes `mean`, `std`, `p50` and `true_rate`, so it rounds them to four decimal places. `min` and `max` are values from the data, so the table shows them exactly. The one exception is a `Decimal`: the table shows it as a float, so it rounds any value with more digits than a float holds. Lengths are in bytes, the unit of Dataframely's `max_length` on a `String`, so this table matches the column constraints. The table shows a `span` in Polars' own duration form, such as `8d` or `1m 30s`, not ISO-8601. **The string group has no statistic that shows values from the data**, at any setting. A `min` or `max` on an email column would write real addresses into the event log, which a deployment shares, anyone can export, and nothing deletes. Lengths and cardinality detect the same problems. The package computes no statistics for the invalid rows. ## The checks The column-schema check reports first and is blocking. Each rule set then has one check. | check metadata | on which check | | --- | --- | | `dy_rule`, `dy_rule__expr` | a check reporting for one rule | | `dy_failed_count` | a check reporting for one rule, when any row failed | | `dy_failed_sample` | any check, when any row failed | | `dy_rules` | a collapsed check: a row per rule with its failure count and expression | | `dy_schema__errors` | the column-schema check, when it fails: a row per mismatched column | A bound appears in `dy_rule__expr`, not in the check name, so changing a `min` neither renames the check nor starts a new check history. In a collapsed check, each row of `dy_failed_sample` has a `dy_rule` column with the rule that row failed. A collapsed check has no total failure count. Counts are per rule, and one row can fail several rules, so their sum is not a row count. `max_failure_samples` applies per rule, not per check. So in a collapsed check, a rule that a thousand rows failed cannot fill the sample and leave out a rule that one row failed. A run that raises `NoValidRowsError` yields no materialization, so the package copies `dataframely/quarantine_address` onto every check result instead. The Columns tab is not one of these two places. It comes from the asset definition, so the catalog shows it before the first run and still shows it after a failed run. ## Attaching your own metadata Use the context to read the run. Use the return value to write the materialization. **Return a `dg.MaterializeResult` with the frame as its `value`.** `@dg.asset` accepts the same return value. It is the only supported way to set a materialization's tags and data version. It also works under direct invocation: ```{python} #| label: what-a-run-produces-1 @dd.asset(Orders) def orders(raw_orders: pl.DataFrame) -> dg.MaterializeResult[pl.DataFrame]: return dg.MaterializeResult( value=raw_orders, metadata={"source": "stripe", "extracted_at": "2026-08-13"}, data_version=dg.DataVersion("2026-08-13"), tags={"run/flavour": "backfill"}, ) ``` `value` is the frame to validate, and you have to set it. The package copies `metadata`, `data_version` and `tags` onto the table's materialization. The quarantine has no materialization event, so none of these fields apply to it. Setting the result's `asset_key` or `check_results` raises `MaterializeResultFieldError`, which names the field. The decorator sets the asset key from its declaration and the check results from the schema's rules. **This package's own metadata keys take precedence.** If you return `dagster/row_count` or a key in the `dataframely/` namespace, the package's value replaces yours. The package copies every other key you return onto the materialization unchanged. ![An asset whose returned result set `source`, `extract/rows` and `extract/window`, and also set `dagster/row_count` to 999. The materialization shows the three keys unchanged, and `dagster/row_count` shows 12, the number of rows written.](../assets/images/returned-result-metadata.png) ### Using the context `context.add_asset_metadata` is the other way to attach metadata, and older Dagster examples use it. The package does not block it. The decorator builds one asset with one key, so the call works without an `asset_key`: ```{python} #| label: what-a-run-produces-2 @dd.asset(Orders) def orders(context: dg.AssetExecutionContext, raw_orders: pl.DataFrame) -> pl.DataFrame: context.add_asset_metadata({"source": "stripe"}) return raw_orders ``` **It overrides this package's own keys, unlike a returned result.** Dagster merges the context's metadata last, after the package yields its results. So `context.add_asset_metadata({"dagster/row_count": 999})` shows 999 in the catalog for a table with two rows. The package cannot prevent this. **This package does not support `context.set_data_version`.** It is not marked `@public` in Dagster, and it does not work under direct invocation. Return `data_version=` on a `dg.MaterializeResult` instead, which sets the same event tags. Avoid `context.add_output_metadata`. Every asset check is an output, and this package always declares the column-schema check. So an asset built with `dd.asset` always has several outputs, and the call raises: ```text DagsterInvariantViolationError: Attempted to add metadata without providing output_name, but multiple outputs exist. Please provide an output_name to the invocation of `context.add_output_metadata`. ``` Naming the output works: Dagster names the asset's output `result`, whatever you call the asset and under every `key_prefix`. That name is internal to Dagster, and the catalog never shows it. Use `add_asset_metadata` instead. ::: {.callout-note} This package has tested only `add_asset_metadata`, `set_data_version` and `add_output_metadata` on Dagster's context, and it will test no others. It guarantees nothing about the rest of the context, in this or any future Dagster version. If you find another method that works and is worth documenting, [open an issue](https://github.com/ozanozbeker/dagster-dataframely/issues). ::: ### Testing an asset ```{python} #| label: setup-testing-an-asset #| code-fold: true #| code-summary: "Setup: the demo's `Orders` schema, which every example on this page uses" #| output: false 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 marketplace_orders, storefront_orders from pathlib import Path daily = dg.DailyPartitionsDefinition(start_date="2026-01-01") Path("raw/orders").mkdir(parents=True, exist_ok=True) storefront_orders().write_parquet("raw/orders/2026-01-02.parquet") ``` To test an asset, call it. Direct invocation is Dagster's documented way to unit-test an asset. It needs no run, no IO manager and no instance. A call returns the same materializations and check results a run yields, as ordinary Python objects. The validated frame is the materialization's `value`: ```{python} #| label: testing-an-asset-1 import pytest @pytest.fixture def quarantine_dir(tmp_path, monkeypatch): monkeypatch.setenv("DAGSTER_DATAFRAMELY_QUARANTINE_DIR", str(tmp_path)) return tmp_path @dd.asset(Orders, quarantine=True) def orders(raw_orders: pl.DataFrame) -> pl.DataFrame: return raw_orders def test_orders_quarantines_the_bad_lines(quarantine_dir): events = list(orders(dg.build_asset_context(), marketplace_orders())) tables = { event.asset_key: event.value for event in events if isinstance(event, dg.MaterializeResult) } checks = { event.check_name: event.passed for event in events if isinstance(event, dg.AssetCheckResult) } assert tables[dg.AssetKey(["orders"])].height == 12 assert not checks["dy_rule__amount__min"] assert pl.read_parquet(quarantine_dir / "orders_quarantine.parquet").height == 8 ``` ```{python} #| label: testing-an-asset-4 #| echo: false #| output: false # A cell cannot run a pytest test, so the three assertions above would otherwise ship # unchecked on the one page that teaches testing. Call the body directly: only the fixture # carries a decorator, and what it supplies is a directory and an environment variable. import os import tempfile with tempfile.TemporaryDirectory() as tmp: os.environ["DAGSTER_DATAFRAMELY_QUARANTINE_DIR"] = tmp test_orders_quarantines_the_bad_lines(Path(tmp)) ``` **A writer writes the quarantine, and the call does not yield it.** The call yields one materialization, for the table, so the test reads the invalid rows from the parquet file, not from an event. In a run too, the IO manager writes the invalid rows, and Dagster records no materialization for the quarantine. **An asset with `quarantine=True` takes a `context` parameter**, even when your function does not declare one. Pass the context first and the upstream frames after it, as Dagster does. A call has no IO manager, so it writes the invalid rows to the directory in `DAGSTER_DATAFRAMELY_QUARANTINE_DIR` instead. It reads that variable only when some rows are invalid. An asset without `quarantine=True` needs neither the context nor the directory: ```{python} #| label: testing-an-asset-2 @dd.asset(Orders) def plain_orders(raw_orders: pl.DataFrame) -> pl.DataFrame: return raw_orders clean = storefront_orders() events = list(plain_orders(clean)) ``` The call yields a separate `dg.AssetCheckResult` for every check the asset declares, with its asset key set. So a call reports the same check names against the same asset keys as a run. If the decorated function declares a `context` of its own, build one with `dg.build_asset_context()`. A partitioned root asset has no upstream inputs, so its call passes only a context with the partition key. `daily` is the `dg.DailyPartitionsDefinition` from [Partitioning](partitioning.qmd), declared in the setup cell above: ```{python} #| label: testing-an-asset-3 @dd.asset(Orders, partitions_def=daily) def daily_orders(context: dg.AssetExecutionContext) -> pl.DataFrame: return pl.read_parquet(f"raw/orders/{context.partition_key}.parquet") events = list(daily_orders(dg.build_asset_context(partition_key="2026-01-02"))) ``` An asset that aborts raises the error from the call. Test the failure policy with `pytest.raises(dd.errors.ValidationAbortError)`, and a frame whose columns or dtypes differ from the schema with `pytest.raises(dd.errors.ColumnSchemaError)`. **Direct invocation does not support metadata added through the context.** `context.add_asset_metadata` raises when you call the asset directly: ```text AttributeError: 'DirectAssetExecutionContext' object has no attribute '_step_execution_context' ``` A plain `@dg.asset` raises the same error, so the limit is Dagster's, not this package's. Return a `dg.MaterializeResult` instead: a call yields it with its metadata, tags and data version, so a test can read them. ### Partitioning ```{python} #| label: setup-partitioning #| code-fold: true #| code-summary: "Setup: the demo's `Orders` schema, which every example on this page uses" #| output: false 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} #| label: partitioning-1 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} #| label: partitioning-2 @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} #| label: partitioning-3 @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} #| label: partitioning-4 @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} #| label: partitioning-5 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. ![The Partitions tab of an asset partitioned by day and by region, with the partition `2026-08-01|us` selected. Validation ran per partition, so the row count and the statistics on the right are for that partition only.](../assets/images/partitions-grid.png) `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} #| label: partitioning-6 @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} #| label: partitioning-8 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} #| label: partitioning-7 @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` returns. ### `LazyFrame`s ```{python} #| label: setup-lazyframes #| code-fold: true #| code-summary: "Setup: the demo's `Orders` schema, which every example on this page uses" #| output: false import dagster as dg import dataframely as dy import polars as pl import dagster_dataframely as dd from dagster_dataframely_demo.schema import Orders ``` A `pl.LazyFrame` appears in two places: an input your asset reads, and the frame it returns. The IO manager returns a `LazyFrame` or a `DataFrame` for an input, depending on the parameter's annotation. Validation collects a returned `LazyFrame` once, in `Schema.filter`. ```{mermaid} %%| label: lazyframes-diagram flowchart TB subgraph validated["@dd.asset"] direction TB V1["LazyFrame returned"] --> V2["column-schema check"] V2 --> V3["Schema.filter, streaming engine"] V3 --> V4["valid rows and invalid rows, in memory"] V4 --> V5["per-rule checks"] V5 --> V6["DataFrame passed to the IO manager"] end ``` In order: the decorated function returns a `LazyFrame`, the column-schema check runs, `Schema.filter` collects the plan on the streaming engine, the package holds the valid rows and the invalid rows in memory, the checks report, and the IO manager receives a `DataFrame`. ## Reads dispatch on the annotation Annotate an input `pl.LazyFrame` and the IO manager returns a scan that has not run yet. Polars then pushes a `filter` or `select` down into the scan, so it drops those rows and columns while it reads the file: ```{python} #| label: lazyframes-1 @dd.asset(Orders) def recent(orders: pl.LazyFrame) -> pl.LazyFrame: return orders.filter(pl.col("amount") > 100) ``` The IO manager reads the input, not this package, so this depends on the IO manager you bound. The annotation works the same on a plain `@dg.asset`. Annotate `pl.DataFrame` instead and the IO manager reads the whole file. ## Validation collects the plan `Schema.filter` takes a plan and returns two: one for the valid rows and one for the invalid rows. This package collects both in one `collect_all` call on the streaming engine, so Polars reads the source once. **That call collects a returned `LazyFrame` once, and validates it the same way as a returned `DataFrame`.** ```{python} #| label: lazyframes-2 @dd.asset(Orders) def orders(raw_orders: pl.LazyFrame) -> pl.LazyFrame: return raw_orders.filter(pl.col("amount") > 0) ``` So your joins, filters and aggregations run on the streaming engine. Polars falls back to the in-memory engine for any operation the streaming engine does not support, so the engine choice never makes a plan fail. The streaming engine collects only the plan's output into memory, not its intermediate results. Peak memory is that output plus one boolean column per rule, held once while `Schema.filter` separates the valid rows from the invalid rows. This helps a plan with a large intermediate result, such as a join that multiplies rows before a filter removes most of them. The in-memory engine would hold every row that join produced. A returned `DataFrame` goes through the same call after `.lazy()`, which copies no data. The column-schema check runs first and reads only `collect_schema()`, which does not run the plan. So a frame whose columns or dtypes differ from the schema raises `ColumnSchemaError` before the plan runs. The plan runs on the streaming engine, but the write does not stream: the IO manager receives a `DataFrame` in memory. The package counts the valid rows and the invalid rows to determine the outcome, and two outcomes do not write the table. Streaming straight to storage would write the table before the package counted the rows. A plain `@dg.asset` has no validation, checks or statistics, so it can stream end to end without collecting its result into memory. The measurements are in [`docs/pre-1.0.md`](https://github.com/ozanozbeker/dagster-dataframely/blob/main/docs/pre-1.0.md#lazy-validation-and-lazy-storage). Schema on write needs the result in memory, so it suits silver and gold tables better than data at ingestion scale. ## Reference ### Settings ```{python} #| label: setup-settings #| code-fold: true #| code-summary: "Setup: the demo's `Orders` schema, which every example on this page uses" #| output: false import dagster as dg import dataframely as dy import polars as pl import dagster_dataframely as dd from dagster_dataframely_demo.schema import Orders ``` This package resolves each setting from three sources, each overriding the one before: the package default, then an environment variable, then the argument to `dd.asset`. An environment variable sets a default for the code location, and an argument overrides it for one asset. Each variable is `DAGSTER_DATAFRAMELY_` plus the setting's name in upper case. | setting | what it sets | default | | --- | --- | --- | | `check_granularity` | how many checks the package collapses the schema's rules into: `rule`, `column` or `schema` | `rule` | | `schema_rules` | how the package reports the schema-level rules at `column` granularity: `collapsed` or `per_rule` | `collapsed` | | `statistics` | whether each materialization has statistics for the rows it wrote | `true` | | `max_failure_samples` | how many of the rows that failed a rule appear in that rule's check | `5` | | `row_sample` | how many valid rows and how many invalid rows appear in the materialization | `5` | | `quarantine_dir` | the directory the package writes invalid rows to when you call the asset instead of running it | unset: a call with invalid rows raises `QuarantineDirError` | The package resolves the settings once, when you declare the asset, and the check specs and the validation both use those values. `quarantine_dir` has two sources, not three, because `dd.asset` takes no argument for it ([ADR-0006](https://github.com/ozanozbeker/dagster-dataframely/blob/main/docs/pre-1.0.md#adr-0006-the-assets-own-io-manager-writes-the-quarantine)). The package reads its value when the writer writes the rows, not when you declare the asset. An empty value still raises `InvalidSettingError` when you declare the asset, not on the first call with invalid rows. For example, `DAGSTER_DATAFRAMELY_QUARANTINE_DIR=${SCRATCH}` is empty when `SCRATCH` is unset. The package validates every source when it resolves the setting, including the package default and the argument. An invalid value raises `InvalidSettingError`, which names the value and its source. So `statistics="false"` raises instead of turning statistics on, and `max_failure_samples=True` raises instead of meaning one row. **A flag's environment variable accepts `true` or `false`, and nothing else.** Case does not matter, so `TRUE` is the same as `true`. `1`, `yes` and `on` raise `InvalidSettingError`. There is no fourth source and no `set_default_*()` function. ## Changing `check_granularity` starts a new check history `rule` gives each rule its own check and its own history. `column` gives one check per column with rules, `dy_col__`, so a 40-column schema's check list stays readable. A column's check reports for every rule on that column. For a `Struct` column, Dataframely generates one `inner__nullability` rule per field, so a ten-field struct has ten rules and one check. `schema` gives one check, `dy_schema__rules`, for every rule. `schema_rules` sets how the package reports the schema-level rules at `column` granularity, because they belong to no single column. `collapsed` puts them all in one check, `dy_schema__rules`, and `per_rule` gives each one its own check. The setting has no effect at `rule` or `schema` granularity. A schema with no rules has only the column-schema check, at every granularity. That check is always present and always blocking. **Changing `check_granularity` on an asset that has already run starts a new check history.** The asset no longer reports the old checks, so their history ends with the last run before the change. The new checks start with no history. Nothing moves the old history to the new checks, so choose the granularity before you deploy the asset. ## Statistics and both samples are on by default Each materialization has summary statistics for the rows it wrote. [The statistics tables](what-a-run-produces.qmd#the-statistics-tables) lists their columns. ::: {.callout-important} Two of the settings write **real rows of your data to the Dagster event log**, and both are on by default. A deployment shares one event log, and you can export it. This package redacts nothing. If a column contains an email address, a name or an account number, the package writes that value to the log, where it stays. ::: | setting | what it writes | where | | --- | --- | --- | | `max_failure_samples` | up to this many of the rows that failed each rule | that rule's asset check, under `dy_failed_sample` | | `row_sample` | up to this many valid rows and this many invalid rows | the materialization, under `dataframely/valid_sample` and `dataframely/invalid_sample` | A sample is the first rows of the frame, not a random selection. A sample is absent, never empty: with no rows to show, the package writes no metadata key for it. Dataframely's comparable setting defaults to `0`. This package defaults to `5`, because a count shows how many rows failed a rule but not what they contain. Set either setting to `0` to turn it off. Per asset: ```{python} #| label: settings-1 @dd.asset(Orders, max_failure_samples=0, row_sample=0) def orders(raw_orders: pl.DataFrame) -> pl.DataFrame: return raw_orders ``` Or set it for a whole code location, in the deployment's environment: ```bash DAGSTER_DATAFRAMELY_MAX_FAILURE_SAMPLES=0 DAGSTER_DATAFRAMELY_ROW_SAMPLE=0 ``` Turning the samples off does not change `statistics`, which is a separate setting. A sample shows a `Decimal` as a string, with every digit. The statistics tables convert it to a float for display. ### Naming ```{python} #| label: setup-naming #| code-fold: true #| code-summary: "Setup: the demo's `Orders` schema, which every example on this page uses" #| output: false import dagster as dg import dataframely as dy import polars as pl import dagster_dataframely as dd from dagster_dataframely_demo.schema import Orders ``` This page covers four things this package sets: the asset's description, each check's description, the name of the asset's op, and the namespaces of the names and keys it generates. ## The description comes from the schema `dd.asset` takes the asset's description from the first of these you set: the `description=` argument, the schema's docstring, and the decorated function's docstring, which is Dagster's own default. ```{python} #| label: naming-1 class Orders(dy.Schema): """A customer order line: one row per product on one order.""" @dd.asset(Orders) def orders() -> pl.DataFrame: """Joins the two extracts and drops the test accounts.""" ... ``` The catalog shows `A customer order line: one row per product on one order.` as the description. The schema's docstring takes precedence over the function's, because it describes the table, not the code. `dd.asset` reads the schema's own docstring only, and ignores one inherited from a parent class, including `dy.Schema`'s. The spec from `quarantine_spec` has its own description: `Invalid rows from , with one column per rule.` ## Where a check's text comes from The Dagster UI shows a schema's rules in three places, and each has its own fallback order. | place | first choice | then | then | | --- | --- | --- | --- | | Columns tab | the rendered column constraint | the rule name | never the docstring | | check name | `dy_rule__` at `rule` granularity, `dy_col__` or `dy_schema__rules` at the others | never the docstring | | | check description | the rule's docstring | ` ` | the whole name Dataframely reports, column part included | | collapsed check description | each rule's rendered constraint | the rule name | never the docstring | The package renders a rule's value, such as the bound of `min`, in the column constraint and in each check result's `dy_rule__expr` metadata, but never in the check name. So changing a bound does not rename the check or start a new check history. The package renders column constraints as operators: `>= 0.0`, `length <= 64 bytes`, `matches ^[a-z]+$`, `in (a, b)`, `aligned to 1h`. It omits `nullability` and `unique`, because Dagster has separate fields for them on the column. It omits `inf` and `nan` too, because Dataframely adds them to every float column by default. All four still appear in the check name and the check description. A `check=` with a bare lambda has no name, so the package renders it as `custom check` everywhere. Name your checks with a dict, such as `check={"lowercase": ...}`, and the package shows the key instead. ## Dagster names the op after the whole asset key `@dd.asset(Orders, key_prefix="sales", name="orders")` creates an op named `sales__orders`, the same way `@dg.asset` names its op. An op name must be unique within a code location, and an asset name need not be. With the asset name alone, two assets named `orders` under different prefixes would have two ops with the same name. The op name is also the step key and the key under `ops:` in run config, so both use the whole asset key: ```yaml ops: sales__orders: config: threshold: 4 ``` ## The reserved namespaces There are three. **`dy_`** starts every check name, every rule column in the quarantine, and every check metadata key except `dataframely/quarantine_address`. These are the column-schema check `dy_schema__columns`, each rule's check `dy_rule__`, the collapsed checks `dy_col__` and `dy_schema__rules`, and the metadata keys `dy_rule`, `dy_rule__expr`, `dy_rules`, `dy_failed_count`, `dy_failed_sample` and `dy_schema__errors`. A check name becomes an op output name, which Dagster checks against `^[A-Za-z0-9_]+$`. So a check name cannot contain a slash. A rule's check name is `dy_rule__` plus the name Dataframely gives the rule, with `|` replaced by `__`: `amount|min` becomes `dy_rule__amount__min`. `dd.wiring.check_name("amount|min")` returns that check name, so a test does not have to build it by hand. **`dataframely/`** starts every materialization metadata key this package names. A metadata key can contain a slash, so this namespace has the same form as Dagster's own `dagster/`. This package's keys then sort together, apart from Dagster's keys and your IO manager's. `dataframely/quarantine_address` is the one `dataframely/` key that also appears on check results. When every row fails validation, the run raises `NoValidRowsError` and has no materialization. So each check result has the quarantine address under that key. Every public function that takes a schema raises three errors for names this package cannot use ([ADR-0008](https://github.com/ozanozbeker/dagster-dataframely/blob/main/docs/pre-1.0.md#adr-0008-every-public-function-that-takes-a-schema-validates-it)). It raises `ReservedColumnError` for a column name that starts with `dy_`. It raises `InvalidColumnNameError` for a column name with characters Dagster does not allow in a check name, which is anything outside `A-Za-z0-9_`. It raises `CheckNameCollisionError` for two rules that produce the same check name. So `dd.asset` raises them where you declare the asset, and a hand-wired asset raises them in the first of these functions it calls. A column name with such characters almost always comes from `dy.Column(alias=...)`, which Dataframely provides for a name that is not a Python identifier. A `|` is one of them, and Dataframely also uses it to separate a column from its rule name. **`_quarantine`** is the third: the asset key a writer writes the invalid rows under. You also choose asset keys in this key space, so another asset can already have this key. Before the decorated function runs, a run of an asset with `quarantine=True` checks that no other asset in the code location materializes this key. If one does, the run raises `QuarantineKeyCollisionError`. No setting changes any of the three, so they are the same in every project. ### Errors ```{python} #| label: setup-errors #| code-fold: true #| code-summary: "Setup: the demo's `Orders` schema, which every example on this page uses" #| output: false import dagster as dg import dataframely as dy import polars as pl import dagster_dataframely as dd from dagster_dataframely_demo.schema import Orders ``` Every error this package raises is in `dd.errors` and subclasses `dd.errors.DagsterDataframelyError`. Catch one by name, or catch `DagsterDataframelyError` to catch every one. Every message names what failed and how to fix it. | error | raised | what to do | | --- | --- | --- | | `CollectionNotSupportedError` | at decoration | pass a `dy.Schema`; declare one asset per Collection member, each with the member's own schema | | `ReservedColumnError` | at decoration, and by every public function that takes a schema | rename the column that starts with `dy_` | | `InvalidColumnNameError` | the same | rename the column, or change the `alias=` that sets its name, so the name uses only `A-Za-z0-9_`, the characters Dagster allows in a check name | | `CheckNameCollisionError` | the same | rename one of the two rules that produce the same check name | | `InvalidSettingError` | when the package resolves the setting | fix the value at the source the message names | | `MaterializeResultValueError` | before the column-schema check | set `value=` to the frame; or return the frame and call `context.add_asset_metadata`; or write a plain `@dg.asset` and call `dd.wiring.schema_metadata` | | `MaterializeResultFieldError` | the same | remove `asset_key=` or `check_results=`, which the decorator sets itself | | `ColumnSchemaError` | after the column-schema check fails | fix the function that produced the frame, or cast with `Schema.cast` in the asset body | | `ValidationAbortError` | after `Schema.filter`, when rows failed and the asset declares no quarantine | fix the rows upstream, write them to a quarantine with `quarantine=True`, or drop them in the asset body | | `NoValidRowsError` | after `Schema.filter`, when every row failed validation and the asset declares a quarantine | read the invalid rows at the quarantine address in the message | | `QuarantineKeyCollisionError` | before the decorated function runs, on every run of an asset with `quarantine=True` | rename the other asset, or remove `quarantine=True`; if the other asset is your quarantine table, delete it and use `quarantine_spec` | | `QuarantineDirError` | under direct invocation only, when there are invalid rows to write | set `DAGSTER_DATAFRAMELY_QUARANTINE_DIR`, or run the asset | `ColumnSchemaError` names every mismatched column at once, with its expected and actual dtype. The failing check's `dy_schema__errors` table lists the same columns. ![A run log where the `quantity` column is `Int64`. The failing `dy_schema__columns` check has `dy_schema__errors`, which shows the expected and actual dtype. The step failure below it names the same column in the `ColumnSchemaError` message.](../assets/images/error-column-schema.png) `ValidationAbortError` and `NoValidRowsError` both give the failure count per rule. The counts can add up to more than the number of invalid rows, because one row can fail several rules. ## Three ways to get this wrong **Do not declare an asset keyed `_quarantine` if `` has `quarantine=True`.** `orders` with `quarantine=True` writes its invalid rows to the asset key `orders_quarantine`, through the same IO manager that writes `orders`. An asset of your own keyed `orders_quarantine` has the same address. Both assets would write to the same table or file, and the second write would replace the first. So every run of `orders` checks for such an asset before the decorated function runs, and raises: ```text QuarantineKeyCollisionError: 'orders' declares `quarantine=True`, so its invalid rows are written to 'orders_quarantine', which another asset in this code location already materializes. Rename that asset, or remove `quarantine=True` from 'orders'. If that asset is your own quarantine table, delete it and use `quarantine_spec` instead, which adds the quarantine itself to the asset graph. ``` **A `from __future__ import annotations` in your own module breaks an annotated `context` parameter.** Under PEP 563, Dagster receives every annotation as a string. Dagster compares the `context` annotation with the real classes, so it raises this error for both `context: dg.AssetExecutionContext` and `context: AssetExecutionContext`: ```text DagsterInvalidDefinitionError: Cannot annotate `context` parameter with type dg.AssetExecutionContext. `context` must be annotated with AssetExecutionContext, AssetCheckExecutionContext, OpExecutionContext, or left blank. ``` The restriction is Dagster's, not this package's. `@dg.asset` raises the same error for the same annotation. Both accept an unannotated `context`, which is the only option in that message that still works under PEP 563: ```{python} #| label: errors-1 @dd.asset(Orders) def orders(context) -> pl.DataFrame: context.log.info("run %s", context.run_id) return pl.read_parquet("raw/orders.parquet") ``` **A `@dy.rule()` body needs its class parameter.** Without it, Python still defines the schema class and the asset, and raises no error. The run then fails when this package reads the rule's expression for the check metadata: ```text TypeError: Orders.amount_is_positive() takes 0 positional arguments but 1 was given ``` `@dy.rule()` is a classmethod-style decorator, so the body takes `cls`: ```{python} #| label: errors-2 class Orders(dy.Schema): status = dy.Enum(["new", "paid", "shipped", "cancelled"], nullable=False) amount = dy.Decimal(10, 2, nullable=False) @dy.rule() def paid_orders_have_amount(cls) -> pl.Expr: """Require a positive amount on a paid line.""" return (cls.status.col != "paid") | (cls.amount.col > 0) ``` The docstring becomes that check's description in the catalog. ### Hand-wiring ```{python} #| label: setup-hand-wiring #| code-fold: true #| code-summary: "Setup: the demo's `Orders` schema, which every example on this page uses" #| output: false 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 def orders_frame() -> pl.LazyFrame: return storefront_orders().lazy() ``` Use `dd.asset` where you can. One declaration fills the Columns tab, reports every rule as an asset check, filters the rows and writes the invalid rows to the quarantine. `dd.asset` calls functions that `dd.wiring` also exports. Use them where `dd.asset` cannot build the asset: to attach a schema to an asset declared some other way, or to report checks in a way `dd.asset` does not offer. Together, the functions do less than `dd.asset`: you pass them the settings and asset keys that `dd.asset` resolves at definition time. This package will add nothing to `dd.wiring` to make reassembling the decorator easier. To see how `dd.asset` calls the functions, read its source, which is short. In the examples, `orders_frame()` represents whatever produces your frame. ## The parts **The Columns tab.** `schema_metadata(schema)` returns the definition metadata that an asset built with `dd.asset` declares. Pass it as the `metadata` of any `@dg.asset`, and the Columns tab shows the schema: dtypes, descriptions, nullability, uniqueness, the primary key at table level, and every other column constraint next to its column. It returns a mapping with one entry, so you can merge it with your own `metadata`. `dd.asset` merges it last, so its entry takes precedence over any `dagster/column_schema` in yours. `table_schema(schema)` returns the `dg.TableSchema` in that entry. Put it under `dagster/column_schema` in a metadata dict you build yourself. Column tags come from `dy.Column(metadata=...)`, which Dataframely stores but does not use. `table_schema` converts tag values to strings, because `dg.TableColumn.tags` is `Mapping[str, str]` and Dagster rejects any other type at definition time. The Columns tab reads `unique` from each column's own flag, never from `primary_key`. Dataframely checks a primary key with one schema-level rule and does not set the key columns' `unique` flag. For a composite key, marking each column unique would show a constraint that no rule checks. **The checks.** `check_specs(schema, asset=key)` returns a spec for every check the schema defines, including the column-schema check. `check_results(schema, frame, asset_key=key, severity=...)` yields a result for each of those checks and writes nothing. Pass the same `check_granularity` and `schema_rules` to `check_specs` and to the function that yields the results, or to neither. Those settings determine the check names. Dagster matches each result to its spec by name, so different values fail the step. **Validation and the write.** An asset built with `dd.asset` runs `validation_results(schema, frame, valid_key=key, quarantine_writer=...)` after its decorated function. It runs the column-schema check and `Schema.filter`, writes the quarantine, and yields the materialization and every check result. Its return type is `dd.wiring.AssetYield`. The writer writes the quarantine, and `validation_results` yields no result for it. You make one choice: whether to pass `quarantine_writer`. The data determines which of the six outcomes a run has. ```{mermaid} %%| label: hand-wiring-outcomes flowchart TD F["frame passed to validation_results"] --> N{"None?"} N -- yes --> SKIP["nothing written
every check passes
run succeeds"] N -- no --> CS{"column schema matches?"} CS -- no --> DRIFT["nothing written
dy_schema__columns fails, ERROR
rules never evaluated
run fails, ColumnSchemaError"] CS -- yes --> FILTER["Schema.filter"] FILTER --> ANY{"any invalid rows?"} ANY -- no --> CLEAN["table written
every check passes
run succeeds"] ANY -- yes --> Q{"quarantine_writer passed?"} Q -- no --> ABORT["nothing written
failing checks ERROR
run fails, ValidationAbortError"] Q -- yes --> SURV{"any valid rows left?"} SURV -- no --> NONE["quarantine written, table not written
failing checks ERROR
run fails, NoValidRowsError"] SURV -- yes --> PARTIAL["table and quarantine written
failing checks WARN
run succeeds"] ``` [The failure policy](the-failure-policy.qmd) lists the same six outcomes as a table, with what each one writes. **The quarantine.** `validation_results` takes a `QuarantineWriter`: a function that receives the invalid rows, writes them and returns the quarantine address. The address is a string, not a path, because it can name a database table. In a run, `dd.asset` uses `delegating_writer(context)`. It passes the rows to the asset's own IO manager, under the asset key `_quarantine`. That IO manager writes them as it writes any asset with that key. The IO manager receives a copy of the step's output context with only the asset key changed, so a database IO manager still reads its connection settings from it. Nothing records the metadata the IO manager adds during that write, so the package reports the quarantine address itself. `file_writer(key, quarantine_dir, partition_key)` is for direct invocation, which has no step and so no IO manager. It writes only parquet. `validate_quarantine_key(context)` fails the run with `QuarantineKeyCollisionError` when another asset in the code location already materializes `_quarantine`. `dd.wiring` also exports two lower-level functions. `quarantine_frame(schema, failure)` returns the frame every writer receives: the invalid rows, then one `String` rule column per rule, named in the reserved namespace. Use it in an asset that runs `Schema.filter` itself, to write the same frame `dd.asset` would. `quarantine_path(key, quarantine_dir, partition_key)` returns the path `file_writer` writes to, so a test can check the path without writing a file. ## The decorator is a `@dg.asset`, a writer and `validation_results` Pass the schema's metadata and check specs to `@dg.asset`, and call `validation_results` in the body: ```{python} #| label: hand-wiring-1 @dg.asset( metadata=dd.wiring.schema_metadata(Orders), check_specs=dd.wiring.check_specs(Orders, asset="orders"), output_required=False, ) def orders(context: dg.AssetExecutionContext) -> dd.wiring.AssetYield: dd.wiring.validate_quarantine_key(context) yield from dd.wiring.validation_results( Orders, orders_frame(), valid_key=context.asset_key, quarantine_writer=dd.wiring.delegating_writer(context), ) ``` This gives the asset the Columns tab, one check per rule, the row filter and the quarantine. For a single-output asset, `context.asset_key` is the only key you need. When the frame matches the schema, a run of that asset makes these calls: ```{mermaid} %%| label: hand-wiring-sequence sequenceDiagram participant A as orders participant V as validation_results participant W as delegating_writer participant D as Dagster A->>V: Orders, frame, valid_key, quarantine_writer V->>V: column-schema check V->>V: Schema.filter opt some rows failed V->>W: quarantine frame W-->>V: quarantine address end opt the table was written V-->>D: MaterializeResult end V-->>D: an AssetCheckResult per check ``` `orders` passes the frame and the writer to `validation_results`. If rows failed, `validation_results` calls the writer with the quarantine frame. The writer returns the quarantine address. `validation_results` then yields the `MaterializeResult` if it wrote the asset's table, and one result for each check. `orders` passes each one to Dagster with `yield from`. `output_required=False` lets the step end without yielding the materialization, which happens on a column-schema mismatch, on `ValidationAbortError`, on `NoValidRowsError` and on the skip. Without it, every path that does not raise has to yield the output. Without a `quarantine_writer`, any invalid row fails the run with `ValidationAbortError`, as in an asset built with `dd.asset` and `quarantine=False`. `delegating_writer` needs a step, so it raises `dagster._core.errors.DagsterInvalidPropertyError` under direct invocation. There, `dd.asset` uses `file_writer(context.asset_key, quarantine_dir, partition_key)`, and you can do the same. Call `validate_quarantine_key(context)` before the rest of the body, as `dd.asset` does on every run. Under direct invocation it does nothing, because there is no asset graph to check. ## Split the checks off entirely The example above passes its frame to `validation_results`, which writes the table and yields the check results in one step. To separate them, return the frame from an ordinary asset, and put the checks in a `@dg.multi_asset_check` that reads the table back through the IO manager: ```{python} #| label: hand-wiring-2 from collections.abc import Iterator KEY = dg.AssetKey(["orders"]) @dg.asset(metadata=dd.wiring.schema_metadata(Orders)) def orders() -> pl.LazyFrame: return orders_frame() @dg.multi_asset_check(specs=dd.wiring.check_specs(Orders, asset=KEY)) def orders_checks(orders: pl.LazyFrame) -> Iterator[dg.AssetCheckResult]: yield from dd.wiring.check_results( Orders, orders, asset_key=KEY, severity=dg.AssetCheckSeverity.WARN ) ``` ![The asset's Checks tab. The checks ran in their own `@dg.multi_asset_check`, against a table the IO manager had already written. The failing checks have the `WARN` severity passed to `check_results`, and the rows that failed them are in that table, not in a quarantine.](../assets/images/split-checks.png) Use this when the write must not depend on the check results. The IO manager writes whatever the asset returns, lazy or eager, and the checks run afterwards against the written table. Unlike with `dd.asset`, the IO manager writes the table before validation, so it can contain invalid rows. Only a failing check reports them. `check_results` does what `validation_results` does, without writing anything. It yields no materialization, writes no quarantine, and never raises `ValidationAbortError` or `NoValidRowsError`. `dd.asset` sets two things that you pass to `check_results` yourself: - **`severity`.** `validation_results` sets it from whether it wrote the table: `WARN` if it did, `ERROR` if it did not. Here the IO manager always writes the table, but with the invalid rows in it, which is worse than the outcomes `validation_results` reports as `ERROR`. So pass `WARN` to report failures on a table the IO manager wrote anyway, or `ERROR` to give them the severity of a failed run. Every rule's check in the step gets that severity. - **Use the same settings as the specs.** Pass `check_granularity` and `schema_rules` to both calls, or to neither. `check_results` still raises `ColumnSchemaError` when the column schema does not match. The column-schema check fails, the step fails on the error, and `check_results` yields no result for any rule, because it evaluated no rule. ## Without the package at all Every example above uses `dd.wiring`. This is the same asset without this package: the Columns tab, one check per rule, the column-schema check, the row filter and the quarantine, all written by hand. ```{python} #| label: hand-wiring-3 VALID = dg.AssetKey(["orders"]) QUARANTINE = dg.AssetKey(["orders_quarantine"]) COLUMNS = Orders.columns() # Private in Dataframely: nothing public lists a schema's rules before it runs. RULES = list(Orders._validation_rules(with_cast=False)) def check_name(rule: str) -> str: return f"dy_rule__{rule.replace('|', '__')}" COLUMN_SCHEMA = dg.TableSchema( columns=[ dg.TableColumn( name=name, type=str(column.dtype), description=column.description, constraints=dg.TableColumnConstraints( nullable=column.nullable, unique=column.unique ), ) for name, column in COLUMNS.items() ] ) @dg.multi_asset( outs={ "orders": dg.AssetOut( metadata={"dagster/column_schema": COLUMN_SCHEMA}, is_required=False ), "orders_quarantine": dg.AssetOut(is_required=False), }, check_specs=[ dg.AssetCheckSpec(name="dy_schema__columns", asset=VALID, blocking=True), *(dg.AssetCheckSpec(name=check_name(rule), asset=VALID) for rule in RULES), ], ) def orders(): # Eager, because `Schema.filter` returns a `LazyFrame` for a lazy frame and every # count below is a `len()`. The decorator does this collect for you. frame = orders_frame().collect() drift = { name: (column.dtype, frame.schema.get(name)) for name, column in COLUMNS.items() if frame.schema.get(name) != column.dtype } yield dg.AssetCheckResult( check_name="dy_schema__columns", asset_key=VALID, passed=not drift ) if drift: raise ValueError(f"{Orders.__name__} does not match the frame: {drift}") valid, failure = Orders.filter(frame, cast=False) counts = failure.counts() aborting = bool(len(failure)) and not len(valid) if len(valid): yield dg.MaterializeResult( asset_key=VALID, value=valid, metadata={"dagster/row_count": len(valid)} ) if len(failure): rule_columns = {rule: check_name(rule) for rule in RULES} invalid = failure.details().rename(rule_columns) yield dg.MaterializeResult( asset_key=QUARANTINE, value=invalid.with_columns( pl.col(name).cast(pl.String) for name in rule_columns.values() ), metadata={"dagster/row_count": len(failure)}, ) for rule in RULES: yield dg.AssetCheckResult( check_name=check_name(rule), asset_key=VALID, passed=not counts.get(rule), severity=dg.AssetCheckSeverity.ERROR if aborting else dg.AssetCheckSeverity.WARN, ) ``` That runs, and it writes both tables. It differs from `dd.asset` in these ways: - It calls `Orders._validation_rules`, which is private. Dataframely has no public way to list a schema's rules before validation, so every check name and rule column here depends on an API with no deprecation guarantee. This package calls the same private API, and a characterization test checks it. So a change in Dataframely fails one of this package's tests, not your assets. - The Columns tab shows dtypes, descriptions, nullability and uniqueness. The rest of what `Orders` declares is missing: no `>= 0`, no regex, no length bound, no table-level primary key, no column tags. - The checks have no descriptions, so a failing check shows the rule's name but not what the rule checks. - No statistics, no row sample, no failure samples and no `invalid_by_rules` table, so a failing check shows nothing about the rows that failed it. - No `check_granularity`, so a 40-column schema has around 40 checks or more, and you cannot collapse its rules into fewer checks. - You have to collect and filter a `LazyFrame` yourself. - It has three outcomes, not six. A run where every row failed succeeds and skips the asset's table. `dd.asset` fails that run with `NoValidRowsError`, because `quarantine=True` lets a run succeed when some rows fail, not when every row fails. - The quarantine is a second output, so a run that aborts never writes it, which is when you most need its rows. - Nothing checks the names. A column already named `dy_rule__amount__min`, or two rules that produce the same check name, cause a silent collision instead of an error at definition time. - With a `key_prefix`, you build both asset keys yourself, and nothing checks them. ### Or you could just do ```{python} #| label: hand-wiring-4 @dd.asset(Orders, quarantine=True) def orders(raw_orders: pl.DataFrame) -> pl.DataFrame: return raw_orders ```