Skip to main content

dagster-dataframely

dataframely declares what a polars frame has to look like. Dagster has first-class surfaces for exactly that: the Columns tab and asset checks. This package connects them with one declaration.

import dagster as dg
import dataframely as dy
import polars as pl

import dagster_dataframely as dd


class Orders(dy.Schema):
    order_id = dy.String(primary_key=True)
    amount = dy.Float64(nullable=False, min=0.0)


@dd.dataframely_asset(schema=Orders)
def orders(raw_orders: pl.DataFrame) -> pl.DataFrame:
    return raw_orders.select("order_id", "amount")

That is the whole integration.

Before the asset has ever run, Orders fills the catalog's Columns tab: dtypes, descriptions, nullability, uniqueness, the primary key stated once at table level, and one pill per remaining constraint. Every run reports one asset check per dataframely rule, each with its own pass/fail history, behind a blocking gate check that compares the frame's columns and dtypes against the schema before a single row is filtered. Add quarantine=dg.AssetOut() and the rows the schema rejects land in a sibling asset instead of failing the run, as long as something survives.

The transform keeps plain polars annotations. Upstream assets bind as ordinary parameters, the return may be a DataFrame or a LazyFrame, and there is no context parameter to write: the wrapper reaches the context itself.

Installation

uv add dagster-dataframely

Requires Python 3.12+.

[!NOTE] Pre-1.0. The public surface is pinned by a test rather than held by convention, so it will not move quietly. It can still move: a 0.x minor release is where a breaking change lands. Pin dagster-dataframely>=0.1,<0.2 if that matters to you.

dagster, dataframely and polars are the dependencies. The IO managers import pydantic and universal-pathlib directly, so both are declared too; both already arrive with dagster, so nothing new lands in your environment. Writing to s3://, gs:// or az:// needs that scheme's fsspec filesystem installed alongside: s3fs, gcsfs or adlfs.

The failure policy is the asset's declared shape

There is no lenient mode and no strict flag, deliberately. Declaring a quarantine is the consent to partial data, so what a rejected row costs is visible in the definition and cannot disagree with what the asset declares.

what the transform returned good table quarantine checks run
columns or dtypes that are not the schema's not written not written the gate fails, blocking fails, SchemaGateError
every row valid written skipped, not written empty all pass green
some rows rejected, no quarantine declared not written n/a fail at ERROR fails, ValidationAbortError
some rows rejected, quarantine declared the survivors the rejected rows fail at WARN green
every row rejected, quarantine declared skipped every row fail at ERROR fails, NothingSurvivedError

Without a quarantine, every row has to be good. A run that rejects even one row fails and writes nothing, so the last-known-good table stays in place. Landing the survivors and dropping the rest is the failure this package exists to make visible, so it is not reachable by configuration. To drop rows anyway, drop them in your own asset body, where the drop is a line you wrote:

@dd.dataframely_asset(schema=Orders)
def orders(raw_orders: pl.DataFrame) -> pl.DataFrame:
    good, _ = Orders.filter(raw_orders)
    return good

The quarantine table carries the rejected rows with the original columns, then one String column per rule reading valid, invalid or unknown, named exactly as that rule's asset check. Its materialization also carries a cooccurrence table, so one broken upstream field tripping three rules reads as one row rather than three unrelated counts. It inherits the good asset's key prefix, group and IO manager, and its own dg.AssetOut can override any of them, which is how rejected rows reach a separate storage and ownership domain. Three settings raise QuarantineSettingError instead of being honoured: automation_condition, freshness_policy and code_version cannot differ between two outs of one step, so they belong on the decorator, where they cover both.

The package never casts

The gate compares the frame's dtypes against the schema's and aborts on a mismatch; the filter runs with cast=False. A Duration('ns') arriving where the schema declares Duration('us') is a pipeline defect, and silently widening it is how a thousandfold error reaches a table nobody re-reads.

If you do want conformance, write the cast yourself, in your own asset body, as a line you can see:

@dd.dataframely_asset(schema=Orders)
def orders(raw_orders: pl.DataFrame) -> pl.DataFrame:
    return Orders.cast(raw_orders)

The one cast the package makes is on columns it generated itself: the quarantine's outcome columns go from Enum to String, because a raw Enum panics the Delta writer.

Storage

defs = dg.Definitions(
    assets=[orders],
    resources={
        "io_manager": dd.DataframelyParquetIOManager(
            base_dir="s3://my-bucket/warehouse"
        )
    },
)

DataframelyParquetIOManager and DataframelyCSVIOManager write under a universal-pathlib base_dir, so a local directory and s3://, gs:// or az:// are the same call on credentials from the ambient environment. A dtype the format cannot hold raises UnwritableDtypeError before the write rather than from inside polars.

These are the supported path. They record path, bytes_written and dagster/storage_kind on each materialization and nothing else. No column schema, in particular: the asset definition owns what the data is, and leaving the materialization bucket empty is what keeps the Columns tab showing the schema as declared.

What degrades on someone else's manager is one surface. A stock polars IO manager writes its own dagster/column_schema onto the materialization, and with both metadata buckets populated Dagster merges them in the catalog Columns tab: column names come out lowercased and the constraint pills disappear, because the materialization's constraint-free schema becomes the base. Dtypes, descriptions and column tags survive, and the Lineage Metadata accordion is definition-only and never merged, so full fidelity is still one click away. This is documented, not designed around. Parity is a preference, never a constraint.

CSV without the losses

A CSV cell holds text, so five dtypes have nowhere to land. Each is encoded on the way out and decoded on the way back by its declared inverse, so the frame read back compares equal to the frame written.

dtype what the cell holds how it reads back
Duration the integer tick count, in the column's own time unit the ticks cast into the declared Duration
List JSON, wrapped in a one-field object named after the column the wrapper decoded and its one field taken
Array the same wrapper decoded as a List, then cast to the declared Array
Struct a JSON object decoded into the declared Struct
Binary base64 decoded back to Binary

A run log line names the encoded columns on both paths, because that is the only surface that can carry it.

Two cases are refused rather than encoded: Binary and Duration inside a List, Array or Struct. Polars cannot write the first to JSON and cannot read the second back, so encoding either would land a file that no longer round-trips. Object is refused by both managers.

The read needs the schema. The decode reads each column's declared dtype off the carrier the decorator puts in the asset's definition metadata, which dd.schema_metadata also builds for a plain @dg.asset. With no schema in reach the read is an ordinary inferred CSV read, and an encoded column arrives as text. The schema never comes from a sidecar file and never from the data, so it costs no round trip and cannot drift. The carrier holds the live class rather than a copy of it, and a live object does not cross a process boundary, so a manager that only reaches a deserialized definition falls back to the inferred read too.

Prefer parquet unless something downstream needs text. Parquet is self-describing, keeps every dtype natively, and needs no schema to read.

Validation materializes

Schema.filter collects, so a validated frame is a frame in memory. A LazyFrame return is accepted and collected, and the IO managers collect before the write with a warning in the run log. Sinking lazily is not on the supported path; it is tracked in issue #27.

The habitat is post-landing transformation: bronze to silver to gold, where the data is already on your side and the question is whether it is fit to publish. Ingestion-scale and larger-than-memory work belongs to other tools.

Settings

Every knob resolves through three tiers, each overriding the one before: the package default, then an environment variable, then the argument on the asset. A platform engineer sets a house style once for a whole code location, and an asset overrides it where that style is wrong. Each variable is DAGSTER_DATAFRAMELY_ plus the setting's name, upper-cased.

setting what it decides default
check_granularity how far the schema's rules collapse into checks: rule, column or schema rule
multi_column_rules where the rules no single column owns land at column granularity: schema or per_rule schema
statistics whether each materialization carries a profile of what it wrote true
max_failure_samples how many of the rows a rule rejected reach that rule's check 5
row_sample how many of the good table's rows reach its materialization 5

The chain validates on resolve, at every tier including the package's own, so a typo raises InvalidSettingError naming the value and the tier that supplied it rather than quietly becoming something else three modules later.

Changing check_granularity orphans check history

rule gives every rule its own check and its own timeline. column gives one check per rule-bearing column, dy_col__<column>, which is what makes a 40-column schema's check list readable. schema gives a single dy_schema__rules for all of them.

Changing it on an asset that has already run orphans that asset's check history. The old check names stop being reported and their timelines end where the change landed, while the new ones start empty. Nothing migrates them, so choose it before the asset ships rather than after.

Statistics and both samples are on by default

Each materialization carries a skimr-style profile of what it wrote: one table per dtype family present, under stats/numeric, stats/temporal, stats/string and stats/boolean, on the quarantine as well as the good table.

[!IMPORTANT] Two of the settings write real rows of your data into the Dagster event log, and both ship on. The event log is shared across a deployment, it is exportable, and nothing here is redacted. If a column holds an email address, a name or an account number, that value lands in the log and stays there.

setting what it writes where
max_failure_samples up to this many of the rows each rule rejected that rule's asset check, under dy_failed_sample
row_sample up to this many of the rows the good table holds its materialization, under sample

The quarantine carries no row sample, deliberately. Its rows already reach the event log once, through the check that rejected each of them and with the rule attached, so a second copy would carry less and cost the same.

dataframely's own comparable setting defaults to 0, so this package is deliberately the more generous of the two. The reason is that a red check raises exactly one question the counts cannot answer: not that amount|min rejected 43 rows, but what three of those rows held. Paying for that in the event log should be a decision, which is what this section is for.

Setting either to 0 turns it off entirely, and the metadata key is then absent rather than empty. Per asset:

@dd.dataframely_asset(schema=Orders, max_failure_samples=0, row_sample=0)
def orders(raw_orders: pl.DataFrame) -> pl.DataFrame:
    return raw_orders.select("order_id", "amount")

Or once for a whole code location, in the deployment's environment:

DAGSTER_DATAFRAMELY_MAX_FAILURE_SAMPLES=0
DAGSTER_DATAFRAMELY_ROW_SAMPLE=0

Turning the samples off leaves statistics on. The string family deliberately carries no value-bearing statistic at any setting, only lengths and cardinality: consenting to summary statistics is not consenting to raw values.

The kit

The decorator is one arrangement of parts the package also exports on their own: check_specs, schema_metadata, table_schema, quarantine_table_schema, quarantine_frame, process and check_name. Reach for them when the decorator's shape is not the shape you need: a schema attached to an asset you did not declare, or an out arrangement the door does not offer. Then the @dg.multi_asset is yours to wire, out of the same parts.

@dg.multi_asset(
    outs={
        "orders": dg.AssetOut(metadata=dd.schema_metadata(Orders), is_required=False)
    },
    check_specs=dd.check_specs(Orders, asset="orders"),
)
def orders():
    yield from dd.process(
        Orders,
        transform(),
        context=dg.AssetExecutionContext.get(),
        good_out="orders",
    )

is_required=False matters: the gate and both abort paths end the step without yielding an out.

This is not a route to dy.Collection support. Hand-wiring a Collection means reimplementing the hardest part of the package rather than assembling it, because process is single-schema by signature. Declare one asset per member instead, each with the member's own schema. Passing a Collection to the decorator raises CollectionNotSupportedError at decoration time.

Two user-side traps

A from __future__ import annotations in your own module breaks a context parameter. Under PEP 563 every annotation reaches Dagster as a string, and its check on the context parameter compares against the real classes, so it rejects context: dg.AssetExecutionContext and context: AssetExecutionContext alike:

DagsterInvalidDefinitionError: Cannot annotate `context` parameter with type dg.AssetExecutionContext.
`context` must be annotated with AssetExecutionContext, AssetCheckExecutionContext, OpExecutionContext, or left blank.

Only the last option in that message survives PEP 563: leave the parameter blank, or take no context at all and reach it with dg.AssetExecutionContext.get(). Decorator call sites avoid this structurally, because the wrapper fetches the context and your transform never takes one. Kit call sites are yours to write, so this one is yours to hit.

A @dy.rule() body needs its class parameter. Written without one, the class still builds and the asset still defines; the run then fails when this package reads the rule's expression for the check metadata:

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:

class Orders(dy.Schema):
    status = dy.String(nullable=False)
    amount = dy.Float64(nullable=False)

    @dy.rule()
    def paid_orders_have_amount(cls) -> pl.Expr:
        """Paid orders must carry a positive amount."""
        return (cls.status.col != "paid") | (cls.amount.col > 0)

The docstring is not decoration: it becomes that check's description in the catalog.

The reserved namespace

Every check name sits under dy_, and so does every quarantine outcome column and every key in check metadata: the gate check dy_schema__dtypes, the rule checks dy_rule__<rule>, the collapsed checks dy_col__<column> and dy_schema__rules. The materialization keys are deliberately outside it, because sample, cooccurrence and stats/* are for a data consumer rather than for this package's bookkeeping. A schema with a column of its own inside the namespace raises ReservedColumnError at definition time, and two rules that rewrite to one check name raise CheckNameCollisionError. The prefix is hardcoded rather than configurable: its whole value is being the same string in every project.

License

Apache-2.0

Metadata

Release files for dagster-dataframely 0.1.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for dagster-dataframely 0.1.0
File Size Uploaded
dagster_dataframely-0.1.0.tar.gz 58.6 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for dagster-dataframely 0.1.0
File Interpreter ABI Platform
dagster_dataframely-0.1.0-py3-none-any.whl Python 3 none any Details

Total release size: 122.3 kB

Release files / dagster_dataframely-0.1.0.tar.gz

Download URL dagster_dataframely-0.1.0.tar.gz
Size 58.6 kB
Tags Source
SHA-256 checksum
How to use checksums
97ac29f8ce4c683b66dff092d4308f61562fae75b228e153e8882ba50fbb06e8
BLAKE2b-256 checksum
How to use checksums
45329378d7b54c63839beeb85d1074699fc86c8545e3e8d60193278975a0e9ae
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Aug 11, 2026.

Transparency log

Release files / dagster_dataframely-0.1.0-py3-none-any.whl

Download URL dagster_dataframely-0.1.0-py3-none-any.whl
Size 63.6 kB
Tags Python 3
SHA-256 checksum
How to use checksums
042c030606edca1fe987d1c253e12d7ea768188e2540b4a79b4380029f801088
BLAKE2b-256 checksum
How to use checksums
87b687aec9ebceeb47e49ae5d0e07975ad556bd02eebb0a0826b63291c3422b7
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Aug 11, 2026.

Transparency log

Release history Release notifications | RSS feed

0.9.0

2 release files

0.8.0

2 release files

0.7.0

2 release files

0.6.0

2 release files

0.5.0

2 release files

0.4.0

2 release files

0.3.0

2 release files

0.2.0

2 release files

This release

0.1.0 This release

2 release files

0.0.1

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page