Portable data pipelines with Dagster + Ibis: making migrations less painful
One Ibis expression, many engines. A small Dagster repo showing how engine migrations can become a config change instead of a rewrite. There are a lot of companies moving from in-house hosting to cloud providers. And with that comes migrations… For data pipelines, particularly, there’s a question of whether the pipeline should also be migrated to use the data warehouse solution provided by the cloud provider. For instance, since our team is migrating to GCP, we had a discussion on whether we should migrate our pipelines running on PySpark to BigQuery or not. We considered multiple dimensions, like cost, effort, and maintainability. This decision involves rewriting a substantial codebase, like Spark DataFrame calls, There is a lot of discussion in the European Union about moving away from US Companies (see this article from Reuters for instance). There is no guarantee that our migration to GCP will be the last, and with that, no guarantee that we won’t have to rewrite it all again. That forced a question: is there a way to write data transformations so that switching the execution engine is a configuration change, not a rewrite? This post is about a small proof-of-concept repo I built to try to answer this question. While researching about alternatives, we found Ibis, a potential solution to minimize the migration effort. It is a dataframe-style expression API that compiles the same code to DuckDB SQL, Spark SQL, BigQuery SQL, Polars, and many other backends, totaling 20+ backends. Ibis minimizes the migration effort, but it’s just a “translation layer”. We needed structure and orchestration, which can be provided by a single tool, instead of using Airflow with an internal framework for structure. Dagster can solve that with many features provided out-of-the-box. It is an orchestrator built around software-defined assets, giving us the structure dbt users are used to: models, lineage, tests, schedules. The thesis: port your logic to Ibis once, and engine migrations become config diffs. Not “migrate to BigQuery for free this time” (the port to Ibis is real work), but “this is the last engine migration that ever requires a rewrite.” Repo: A deliberately ordinary bronze → silver → gold pipeline: * we’ll come back to that asterisk ( The entire engine/env selection is two environment variables: (The repo wraps these as Same code. Same asset graph. Same checks. Different engine and different storage, selected by deployment config, the documented Dagster pattern ( No This is the Dagster-native answer to dbt’s “where does this model land” or Kedro’s catalog. Assets declare inputs/outputs as So “the same pipeline, but sources are parquet in a lake in prod” is literally just: The thing dbt users ask about first is ( We implemented the usual dbt generic tests And since Dagster is an orchestrator, scheduling is native too: the repo ships a This is the part that sells it. The same This isn’t just a demo trick. Here’s where we keep the demo honest, because the claim above deserves an asterisk: “backend-agnostic” means portable across the operations each backend can express, not that every expression runs everywhere. The repo deliberately includes Ibis’s Polars backend has no window-function translation (Polars natively has Two things make this acceptable, even good: A more subtle gotcha we found the hard way: Before betting a codebase on a portability claim, check Ibis’s per-backend operations support matrix. Fair question. dbt solves a big chunk of this. The honest comparison: If your transforms are pure SQL and your targets are SQL warehouses, dbt is simpler and battle-tested and you should consider it. The Dagster + Ibis combo wins when: your logic outgrows SQL, you want one lineage graph spanning tables and non-table work, or (our case) you’re tired of paying a rewrite tax on every engine migration. Portable data pipelines with Dagster + Ibis: making migrations less painful
flowchart LR
code["same pipeline code
assets + ibis expressions
(written once)"]
code -->|"deployment config
swaps engine + storage"| dep
subgraph dep["pick a backend"]
direction TB
a["duckdb
local csv files"] ~~~ b["polars
same csv files"] ~~~ c["pyspark
parquet lake"] ~~~ d["bigquery, trino, ...
one more config entry"]
end
spark.sql(...) strings, .toPandas() escapes. The demo
flowchart LR
subgraph sources["External sources:
- CSV locally
- Parquet/tables in production"]
events_csv
products_csv
end
subgraph bronze["Bronze:
- Land sources into managed tables"]
raw_events
raw_products
end
subgraph silver["Silver:
- Normalize types
- Trim/lower strings
- Dedupe
- Validate"]
cleaned_events
end
subgraph gold
daily_active_users
category_revenue
latest_event_per_user["latest_event_per_user*"]
end
events_csv --> raw_events
products_csv --> raw_products
raw_events --> cleaned_events
cleaned_events --> daily_active_users
cleaned_events --> category_revenue
raw_products --> category_revenue
cleaned_events --> latest_event_per_user
latest_event_per_user). It’s the most interesting part.DAGSTER_DEPLOYMENT_NAME=local uv run dagster dev # duckdb + local CSVs
DAGSTER_DEPLOYMENT_NAME=polars uv run dagster dev # polars + same CSVs
DAGSTER_DEPLOYMENT_NAME=prod uv run dagster dev # pyspark + parquet "lake" sourcesjust dev <local|polars|prod>; dev prod first seeds a local parquet lake under data/lake/ from the CSVs (just seed-lake), and prod also needs the pyspark extra plus a JDK. Note that ibis currently pins pyspark<4.1, which doesn’t run on Java 25; a future ibis release should allow Spark 4.2, the version that adds Java 25 support.)resources_by_deployment, keyed on DAGSTER_DEPLOYMENT_NAME; Dagster+ sets it automatically, and on a self-hosted OSS deployment it’s just another env var on your code location). The three pieces
1. Transforms are pure Ibis expressions
# transforms.py: no imports from any engine, just ibis
def clean_events(raw_events: ir.Table) -> ir.Table:
normalized = raw_events.mutate(
ts=_.ts.cast("timestamp"),
amount=_.amount.cast("float64"),
event_type=_.event_type.lower().strip(),
)
return (
normalized.mutate(
amount=_.amount.fill_null(0.0), date=_.ts.truncate("D")
)
.filter(_.event_id.notnull(), _.user_id.notnull())
.distinct()
)duckdb., no spark., no pl. anywhere. These functions are pure Table -> Table expressions. They don’t even know which backend they’ll run on until Dagster binds a connection at runtime. 2. An IO manager owns all data I/O
ir.Table; IbisIOManager (a ConfigurableIOManager) handles persistence:load_input → con.table(name) for pipeline tables, or con.read_csv/read_parquet/... for external sources configured per deploymenthandle_output → con.create_table(asset_name, expr, overwrite=True) plus metadata (row count, preview, the compiled SQL)# resources_by_deployment: same logical source, different physical read
"local": IbisIOManager(
sources={"events_csv": {"format": "csv", "path": "data/raw_events.csv"}}, ...),
"prod": IbisIOManager(
sources={"events_csv": {"format": "parquet", "path": "${DATA_LAKE}/landing/events/"}}, ...), 3. Quality gates are portable too
dbt test. Dagster’s @asset_check is the direct analog, and because the checks themselves are Ibis expressions, they’re portable as well:@dg.asset_check(asset=cleaned_events, blocking=True)
def cleaned_events_no_null_user_ids(ibis: IbisResource) -> dg.AssetCheckResult:
n = _violations(ibis, "cleaned_events", _.user_id.isnull())
return dg.AssetCheckResult(passed=n == 0, metadata={"null_user_ids": n})_violations is a small helper: int(con.table(table).filter(predicate).count().execute()) on the run’s backend.)not_null, unique, accepted_values, relationships (an anti-join), plus a couple of custom checks. Marking the not-null check blocking=True reproduces dbt build semantics: if it fails, downstream assets don’t materialize.daily_schedule (0 6 * * *) covering the whole job, stopped by default so it can be toggled on from the UI. The “aha”: one expression, three dialects
daily_active_users(clean_events(raw_events)) expression is compiled offline. No engine, no credentials, no JVM:-- duckdb
... DATE_TRUNC('DAY', "t0"."ts") AS "date" ... TRIM(LOWER("t0"."event_type"), ' ')
-- pyspark
... DATE_TRUNC('day', `t0`.`ts`) AS `date` ... TRIM(' ' FROM LOWER(`t0`.`event_type`))
-- bigquery
... TIMESTAMP_TRUNC(`t0`.`ts`, DAY) AS `date` ... TRIM(LOWER(`t0`.`event_type`), ' ')DATE_TRUNC vs TIMESTAMP_TRUNC, double quotes vs backticks, TRIM(x) vs TRIM(' ' FROM x): all the dialect trivia you never want to hand-maintain, generated from one expression.ibis.<backend>.compile() works without connecting, so a pytest file that compiles every transform against every target dialect is a CI guardrail. If someone adds an operation your production engine can’t express, it fails before deployment, not after. The repo wires this in concretely: a prek pre-push hook runs lint (ruff), type-checking (ty), and that test suite on every push. The honest part: portability has edges
latest_event_per_user, built on a window function:events.mutate(rn=ibis.row_number().over(
ibis.window(group_by="user_id", order_by=_.ts.desc()))).filter(_.rn == 0).over(), but Ibis doesn’t map to it yet). When we materialize the graph on Polars we get:latest_event_per_user FAILED:
ibis.common.exceptions.OperationNotDefinedError:
No translation rule for WindowFunctiongroup_by + join (latest_event_per_user_portable in the repo). Clunkier, but runs everywhere.ibis.get_backend(t).name inside a transform tells you which engine you’re bound to. It’s ugly, but contained to one functionibis.row_number() is zero-based. It compiles to ROW_NUMBER() - 1. rn == 1 silently gives you the second-latest row. Abstraction means portable syntax; you still need to learn the portable semantics. Okay, but why not just dbt?
dbt (+ adapters) Dagster + Ibis Transform language SQL Python expressions → SQL/plans Engine portability per-adapter SQL macros one expression → many dialects Non-SQL engines no yes (Polars, DataFusion…) Python-native logic (ML, APIs, files) bolted on first-class (the same graph) Tests dbt test@asset_checkScheduling needs dbt Cloud / external orchestrator native (schedules, sensors) Incremental models dbt run incrementalpartitions + automation policies Ecosystem maturity bigger younger Caveats worth stating
spark.sql/DataFrame/UDF code must be rewritten as Ibis expressions. The payoff is that it’s potentially the last port.row_number story). Also type inference on file reads (we added portable casts to absorb that). Takeaways
deps as ref(), @asset_check as tests, schedules instead of cron.OperationNotDefinedError at translate time, and the repo shows three ways to handle it.
flowchart LR
code["pipeline code
(write once)"]
code -->|"engine: config"| eng["duckdb | polars | pyspark | ..."]
code -->|"storage: config"| st["csv | parquet | warehouse tables"]