Hand-wiring

Setup: the demo’s Orders schema, which every example on this page uses
import dagster as dg
import dataframely as dy
import polars as pl

import dagster_dataframely as dd
from dagster_dataframely_demo.schema import Orders
from dagster_dataframely_demo._data import storefront_orders


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.

flowchart TD
    F["frame passed to validation_results"] --> N{"None?"}
    N -- yes --> SKIP["nothing written<br/>every check passes<br/>run succeeds"]
    N -- no --> CS{"column schema matches?"}
    CS -- no --> DRIFT["nothing written<br/>dy_schema__columns fails, ERROR<br/>rules never evaluated<br/>run fails, ColumnSchemaError"]
    CS -- yes --> FILTER["Schema.filter"]
    FILTER --> ANY{"any invalid rows?"}
    ANY -- no --> CLEAN["table written<br/>every check passes<br/>run succeeds"]
    ANY -- yes --> Q{"quarantine_writer passed?"}
    Q -- no --> ABORT["nothing written<br/>failing checks ERROR<br/>run fails, ValidationAbortError"]
    Q -- yes --> SURV{"any valid rows left?"}
    SURV -- no --> NONE["quarantine written, table not written<br/>failing checks ERROR<br/>run fails, NoValidRowsError"]
    SURV -- yes --> PARTIAL["table and quarantine written<br/>failing checks WARN<br/>run succeeds"]

The failure policy 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 <name>_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 <name>_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:

@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:

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:

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.

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.

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

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