LazyFrames

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

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.

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:

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

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

Schema on write needs the result in memory, so it suits silver and gold tables better than data at ingestion scale.