# Hand-wiring


Setup: the demo's `Orders` schema, which every example on this page uses

``` python
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](../reference/wiring.table_schema.md#dagster_dataframely.wiring.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](../reference/wiring.check_specs.md#dagster_dataframely.wiring.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](../reference/wiring.validation_results.md#dagster_dataframely.wiring.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.


*\[Rich HTML output -- view on the documentation site\]*


[The failure policy](the-failure-policy.md) lists the same six outcomes as a table, with what each one writes.

**The quarantine.** [validation_results](../reference/wiring.validation_results.md#dagster_dataframely.wiring.validation_results) takes a [QuarantineWriter](../reference/wiring.QuarantineWriter.md#dagster_dataframely.wiring.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](../reference/errors.QuarantineKeyCollisionError.md#dagster_dataframely.errors.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](../reference/wiring.file_writer.md#dagster_dataframely.wiring.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](../reference/wiring.validation_results.md#dagster_dataframely.wiring.validation_results)

Pass the schema's metadata and check specs to `@dg.asset`, and call [validation_results](../reference/wiring.validation_results.md#dagster_dataframely.wiring.validation_results) in the body:


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


*\[Rich HTML output -- view on the documentation site\]*


`orders` passes the frame and the writer to [validation_results](../reference/wiring.validation_results.md#dagster_dataframely.wiring.validation_results). If rows failed, [validation_results](../reference/wiring.validation_results.md#dagster_dataframely.wiring.validation_results) calls the writer with the quarantine frame. The writer returns the quarantine address. [validation_results](../reference/wiring.validation_results.md#dagster_dataframely.wiring.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](../reference/errors.ValidationAbortError.md#dagster_dataframely.errors.ValidationAbortError), on [NoValidRowsError](../reference/errors.NoValidRowsError.md#dagster_dataframely.errors.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](../reference/errors.ValidationAbortError.md#dagster_dataframely.errors.ValidationAbortError), as in an asset built with `dd.asset` and `quarantine=False`.

[delegating_writer](../reference/wiring.delegating_writer.md#dagster_dataframely.wiring.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](../reference/wiring.validation_results.md#dagster_dataframely.wiring.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
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
    )
```


<figure class="figure">
<p><img src="../assets/images/split-checks.png" class="img-fluid figure-img" /></p>
<figcaption>The asset's Checks tab. The checks ran in their own <code>@dg.multi_asset_check</code>, against a table the IO manager had already written. The failing checks have the <code>WARN</code> severity passed to [check_results](../reference/wiring.check_results.md#dagster_dataframely.wiring.check_results), and the rows that failed them are in that table, not in a quarantine.</figcaption>
</figure>


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](../reference/wiring.check_results.md#dagster_dataframely.wiring.check_results) does what [validation_results](../reference/wiring.validation_results.md#dagster_dataframely.wiring.validation_results) does, without writing anything. It yields no materialization, writes no quarantine, and never raises [ValidationAbortError](../reference/errors.ValidationAbortError.md#dagster_dataframely.errors.ValidationAbortError) or [NoValidRowsError](../reference/errors.NoValidRowsError.md#dagster_dataframely.errors.NoValidRowsError).

`dd.asset` sets two things that you pass to [check_results](../reference/wiring.check_results.md#dagster_dataframely.wiring.check_results) yourself:

- **`severity`.** [validation_results](../reference/wiring.validation_results.md#dagster_dataframely.wiring.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](../reference/wiring.validation_results.md#dagster_dataframely.wiring.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](../reference/wiring.check_results.md#dagster_dataframely.wiring.check_results) still raises [ColumnSchemaError](../reference/errors.ColumnSchemaError.md#dagster_dataframely.errors.ColumnSchemaError) when the column schema does not match. The column-schema check fails, the step fails on the error, and [check_results](../reference/wiring.check_results.md#dagster_dataframely.wiring.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
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](../reference/errors.NoValidRowsError.md#dagster_dataframely.errors.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
@dd.asset(Orders, quarantine=True)
def orders(raw_orders: pl.DataFrame) -> pl.DataFrame:
    return raw_orders
```
