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.xminor release is where a breaking change lands. Pindagster-dataframely>=0.1,<0.2if 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
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)
| File | Size | Uploaded | |
|---|---|---|---|
| dagster_dataframely-0.1.0.tar.gz | 58.6 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| 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 logRelease 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