dplyneage answers a simple question about your data pipelines:
where did each column come from? You pipe a dplyr/dbplyr query
(or pass raw SQL) into extract_lineage(), and render the
answer as an interactive diagram with lineage_flow().
This vignette starts with the smallest possible example and works up to the cases where lineage gets hard: joins with ambiguous columns, CTEs, and columns computed from several sources at once.
Installation
Install the released version from CRAN:
install.packages("dplyneage")Or the development version from GitHub:
pak::pak("tgerke/dplyneage")dbplyr, dtplyr, and arrow pipelines are analyzed entirely in R, so
for those there is nothing else to install. Raw SQL strings and duckplyr
frames (see Other lazy backends) are
analyzed by sqlglot, dplyneage’s one Python dependency: install the
reticulate package to enable that engine, and sqlglot itself is
provisioned automatically the first time it’s needed. See
vignette("python-integration") if you manage your own
Python environment.
Your first lineage diagram
First, create a small in-memory DuckDB database to work with:
library(dplyneage)
library(dplyr)
library(duckdb)
con <- dbConnect(duckdb(), ":memory:")
customers <- tibble(
id = 1:5,
name = c("Alice", "Bob", "Charlie", "Diana", "Eve"),
email = paste0(tolower(name), "@example.com")
)
orders <- tibble(
order_id = 1:10,
customer_id = rep(1:5, each = 2),
amount = c(100, 150, 200, 75, 300, 125, 180, 90, 250, 160)
)
copy_to(con, customers, "customers", overwrite = TRUE)
copy_to(con, orders, "orders", overwrite = TRUE)The simplest lineage there is: two columns selected from one table.
tbl(con, "customers") |>
select(id, name) |>
extract_lineage() |>
lineage_flow(height = "300px")Each output column connects back to the source column it came from.
Try dragging the tables around, zooming with the mouse wheel, and
hovering a column to highlight its connections. Clicking a column
isolates its trace cone, everything upstream and downstream of it, and
Escape (or a background click) releases it. The legend in the corner
names the colors; lineage_flow() also offers
theme = "dark", a minimap, and a PNG download
button in the zoom controls.
A realistic pipeline
Lineage becomes useful once transformations pile up. Here is a join followed by an aggregation:
tbl(con, "customers") |>
left_join(tbl(con, "orders"), by = c("id" = "customer_id")) |>
group_by(id, name) |>
summarise(total_spent = sum(amount, na.rm = TRUE), .groups = "drop") |>
extract_lineage() |>
lineage_flow(height = "400px")Notice that total_spent traces back to
orders.amount, not to customers, even though
amount appears unqualified in the generated SQL. When you
pass a dbplyr table, dplyneage doesn’t parse SQL at all: it walks the
pipeline’s own query tree, which records exactly which table each column
came from, so attribution is always right.
If a pipeline embeds raw SQL with dbplyr::sql(), the
query tree can’t see inside that string, so
extract_lineage() hands the whole query to sqlglot instead
(with a message). You can also force a specific engine with
engine = "r" or engine = "sqlglot"; see
?extract_lineage.
Computed columns
Most columns in a pipeline are computed rather than selected, and
mutate() is where that happens. Lineage treats a computed
column like any other: its edges start at every column the expression
reads. This pipeline builds a few of them the usual ways: a plain
expression, a case_when() written against a column defined
two lines earlier, an across() that names its outputs, and
a per-customer share computed with .by.
computed <- tbl(con, "orders") |>
mutate(
amount_eur = amount * 0.92,
tier = case_when(
amount_eur >= 180 ~ "large",
amount_eur >= 90 ~ "medium",
TRUE ~ "small"
),
across(c(amount, amount_eur), round, .names = "{.col}_rounded")
) |>
mutate(
customer_share = amount / sum(amount, na.rm = TRUE),
.by = customer_id
) |>
extract_lineage()
lineage_flow(computed, height = "400px")Every new column traces to orders.amount, including
tier, which was written against amount_eur
rather than amount: a reference to a column defined earlier
in the same mutate() resolves through it to the source. The
two across() outputs each trace to their own input.
customer_share also draws an edge from
customer_id, because on a SQL backend the .by
column becomes the PARTITION BY of a window function and so
sits inside the expression that computes the share; the
sum() in that expression makes the edge an aggregation,
which the diagram animates.
lineage_edges() gives the same lineage as a table, one
row per edge, with each computed column’s classification and defining
expression:
lineage_edges(computed)
#> source_table source_column target_table target_column transformation
#> 1 orders order_id output order_id identity
#> 2 orders customer_id output customer_id identity
#> 3 orders amount output amount identity
#> 4 orders amount output amount_eur transformation
#> 5 orders amount output tier transformation
#> 6 orders amount output amount_rounded transformation
#> 7 orders amount output amount_eur_rounded transformation
#> 8 orders amount output customer_share aggregation
#> 9 orders customer_id output customer_share aggregation
#> expression
#> 1 order_id
#> 2 customer_id
#> 3 amount
#> 4 amount * 0.92
#> 5 case_when(amount_eur >= 180 ~ "large", amount_eur >= 90 ~ "medium", TRUE ~ "small")
#> 6 round(amount)
#> 7 round(amount_eur)
#> 8 amount/sum(amount, na.rm = TRUE)
#> 9 amount/sum(amount, na.rm = TRUE)The rest of the family follows the same rule. Overwriting a column
(mutate(amount = amount * 0.92)) keeps an edge from the
source column, now classified as a transformation instead of an
identity. transmute() and
mutate(.keep = "none") drop the passthrough columns and
their identity edges with them, and mutate(x = NULL) drops
one. rename() and relocate() change nothing
about provenance, so their columns keep identity edges to the original
names. All of this holds on tbl_lazy() frames, dtplyr, and
arrow as well as on a database, with one backend difference: data.table
and Acero have no window clause, so on dtplyr and arrow a
.by column stays an indirect source (see Other lazy backends).
Where lineage gets hard
These are the cases that break naive lineage tools. dplyneage handles them because both engines resolve the full query structure rather than pattern-matching column names: dbplyr pipelines through their query tree, raw SQL through sqlglot’s lineage engine. We’ll use raw SQL here to keep the examples compact.
Tracing through CTEs
Columns are traced through intermediate CTEs back to the
base tables. recent is transparent, and amount
correctly attributes to orders:
extract_lineage("
WITH recent AS (
SELECT customer_id, amount FROM orders WHERE order_date > '2024-01-01'
)
SELECT customer_id, SUM(amount) AS total FROM recent GROUP BY customer_id
") |>
lineage_flow(height = "300px")Columns with multiple sources
A computed column can come from several tables at once.
COALESCE over a full join gets an edge from both
sources:
extract_lineage("
SELECT COALESCE(u.email, a.email) AS email
FROM users u FULL JOIN archive a ON u.id = a.id
") |>
lineage_flow(height = "300px")The same applies to arithmetic across tables
(o.amount * r.rate), CASE expressions, and
both branches of a UNION. dbplyr pipelines get the same
treatment: a full_join() key column traces to both sides,
and a union_all() column to every branch.
Expanding SELECT *
With a schema available, SELECT * expands to real
columns. dbplyr input never needs a schema (the pipeline itself knows
its columns), but for raw SQL, pass one yourself:
extract_lineage(
"SELECT * FROM customers",
schema = list(customers = c("id", "name", "email"))
) |>
lineage_flow(height = "300px")Raw SQL and schemas
As the examples above show, extract_lineage() accepts a
SQL string directly, useful for auditing queries you didn’t write in R.
Two things to know:
-
Qualified columns (
o.amount) always resolve correctly, schema or not. - Unqualified columns need a schema to be attributed with certainty. Pass a named list mapping each table to its columns:
extract_lineage(
"SELECT c.name, order_date
FROM customers c JOIN orders o ON c.id = o.customer_id",
schema = list(
customers = c("id", "name", "email"),
orders = c("order_id", "customer_id", "order_date", "amount")
)
) |>
lineage_flow(height = "300px")Without a schema, SELECT * cannot be expanded and
produces a warning rather than a silently empty diagram.
Queries in other SQL dialects work by setting
dialect:
extract_lineage(query, dialect = "postgres")
extract_lineage(query, dialect = "snowflake")Local data frames
extract_lineage() reads lineage from a lazy
query, the tree dbplyr builds up before anything touches the database. A
pipeline on a plain tibble has no such tree: dplyr executes each verb
immediately, so by the time you could ask about lineage, only the result
is left.
The lightest fix is dbplyr::tbl_lazy(), which wraps the
frame in a lazy table backed by no database at all. Such a pipeline
can’t be collected, but lineage never runs the query, so a diagram
doesn’t need it to be:
sales <- dbplyr::tbl_lazy(
data.frame(
customer_id = c(1, 1, 2),
amount = c(100, 250, 40)
),
name = "sales"
)
sales |>
group_by(customer_id) |>
summarise(total = sum(amount, na.rm = TRUE)) |>
extract_lineage() |>
lineage_flow(height = "300px")When the pipeline should stay runnable, use a real (if throwaway)
database instead. dbplyr::memdb_frame() builds the data in
in-memory SQLite (install the RSQLite package once), and
copy_to(dbplyr::memdb(), df, name = "df") does the same
copy for a frame you already have. Lineage depends only on the
pipeline’s structure, never on the data, so for large frames copying a
slice is enough:
copy_to(dbplyr::memdb(), head(df), name = "df") yields the
same diagram as copying every row.
The duckdb connection from earlier in this vignette works just as
well (copy_to(con, df)); tbl_lazy() is simply
the fastest route when no connection exists and none is needed.
Other lazy backends
dbplyr is not the only backend that builds a query tree before
executing. dtplyr records
steps toward data.table code, duckplyr holds a duckdb
relation behind an ordinary-looking tibble, and arrow queries compile to
Acero plans. extract_lineage() reads all of them, and the
pipe looks the same:
dtplyr::lazy_dt(
data.frame(
customer_id = c(1, 1, 2),
amount = c(100, 250, 40)
),
name = "sales"
) |>
group_by(customer_id) |>
summarise(total = sum(amount, na.rm = TRUE)) |>
extract_lineage() |>
lineage_flow(height = "300px")dtplyr and arrow pipelines are walked natively in R, like dbplyr’s.
duckplyr is the exception: its relation is a duckdb object with no
R-readable structure, so the engine renders it to duckdb SQL and
analyzes that with sqlglot, which means duckplyr lineage needs
reticulate where the others need nothing. The metadata records which
engine ran and a matching dialect label ("data.table" for
dtplyr, "arrow" for arrow, "duckdb" for
duckplyr).
Three backend habits affect the diagram you get:
- Name your frames where the backend lets you.
lazy_dt(df, name = "sales")putssaleson the source node; without it dtplyr auto-names frames_DT1,_DT2, and so on. In-memory duckplyr frames and arrow tables have no name slot at all, so their nodes come out asdf/df_1andarrow_table/arrow_table_2. File-backed sources (duckplyr::read_csv_duckdb(),arrow::open_dataset()) use the file path, which reads much better. - data.table and Acero have no
OVERclause, so a windowed column’s grouping keys stay indirect on dtplyr and arrow (dashedgroup_byedges underinclude_indirect = TRUE) instead of becoming direct sources the way rendered SQL makes them. - duckplyr hands anything it cannot translate back to eager dplyr
without a message:
group_by()followed bymutate(),rowwise(), and functions with no duckdb translation (case_when(),paste0(),between(), andrename_with()are common ones). The result comes back as an ordinary tibble, or as a duckplyr frame with no relation left behind it, andextract_lineage()says so either way.Sys.setenv(DUCKPLYR_FALLBACK_INFO = TRUE)makes duckplyr name the step that fell back; extract lineage from the step before it, or rewrite it (a.bymutate whose aggregates passna.rm = TRUEstays lazy, for instance).
Column labels travel on every backend: label attributes
survive lazy_dt() and arrow_table(), and arrow
schemas contribute column types to the hover cards for free.
Column labels
Column names say what a thing is called; labels say what it means. R
data often has them already: haven and labelled store a
label attribute on every column imported from SAS, SPSS, or
Stata, and database tables keep the same idea as column comments.
extract_lineage() collects both (label
attributes from local frames, comments from a live duckdb or postgres
connection), and a labels argument covers everything else,
winning over the automatic sources:
sales_df <- data.frame(
customer_id = c(1, 1, 2),
amount = c(100, 250, 40)
)
attr(sales_df$customer_id, "label") <- "Customer surrogate key"
attr(sales_df$amount, "label") <- "Order amount in USD"
lineage <- dbplyr::tbl_lazy(sales_df, name = "sales") |>
group_by(customer_id) |>
summarise(total = sum(amount, na.rm = TRUE)) |>
extract_lineage(labels = list(output = c(total = "Total spent per customer")))
doc <- jsonlite::fromJSON(lineage_json(lineage), simplifyVector = FALSE)
doc$nodes[[1]]$labels # sales: both labels, captured from the attributes
#> $customer_id
#> [1] "Customer surrogate key"
#>
#> $amount
#> [1] "Order amount in USD"
doc$nodes[[2]]$labels # output: total hand-labelled, customer_id propagated
#> $total
#> [1] "Total spent per customer"
#>
#> $customer_id
#> [1] "Customer surrogate key"That receipt shows the whole pattern. customer_id passes
through the group_by() unchanged, so its label travels to
the output node on its own: labels (and types, when captured) propagate
along identity edges, the way dbt Catalog carries descriptions through
passthrough columns. total is a computed column no
propagation can describe, so it gets a labels entry keyed
by "output", the synthetic name of a single query’s result;
in a multi-model pipeline, the model’s name plays that role. Had two
sources fed one column conflicting labels, it would have stayed bare:
missing beats wrong.
In the diagram, those labels become hover cards. Try it on this
rendering: customer_id shows its label on both nodes, and
total shows the hand-supplied one.
lineage_flow(lineage, height = "300px")lineage_json() records the same labels as per-node
labels and types maps, and each OpenLineage
schema-facet field gains a description; the OpenLineage
article shows the rendered event.
Building diagrams by hand
For documentation or design sketches, skip extraction entirely and build the diagram yourself:
nodes <- list(
create_table_node(
table_name = "customers",
columns = c("id", "name", "email"),
x = 0, y = 100,
table_type = "source"
),
create_table_node(
table_name = "customer_summary",
columns = c("customer_id", "full_name", "contact"),
x = 500, y = 100,
table_type = "target"
)
)
edges <- list(
create_column_edge("customers", "id", "customer_summary", "customer_id"),
create_column_edge("customers", "name", "customer_summary", "full_name"),
create_column_edge("customers", "email", "customer_summary", "contact")
)
lineage_flow(nodes, edges, height = "300px")Nodes come in three types: "source" (blue),
"transform" (orange), and "target" (green).
Edges accept a label (e.g. "SUM()") and
animated = TRUE for emphasis.
lineage_example() renders a complete hand-built diagram you
can use as a template.
Exporting lineage
Diagrams answer questions interactively; sometimes you need the same lineage as plain data. Two exporters cover the common cases.
lineage_json() serializes the nodes, edges, and metadata
to a small, stable JSON document, stamped with a
format_version so consumers can rely on its shape (see
?lineage_json for the full schema). Because the output is
deterministic, you can commit it alongside your pipeline code and let CI
diff it: if a refactor silently changes where a column comes from, the
diff shows it before it ships; the lineage
checks in CI article turns that into a one-call gate with
lineage_check(). It is also the natural handoff format for
data catalogs or anything scriptable with jq.
lineage <- tbl(con, "customers") |>
left_join(tbl(con, "orders"), by = c("id" = "customer_id")) |>
group_by(id, name) |>
summarise(total_spent = sum(amount, na.rm = TRUE), .groups = "drop") |>
extract_lineage()
lineage_json(lineage)
#> {
#> "format_version": 1,
#> "metadata": {
#> "dialect": "duckdb",
#> "engine": "r",
#> "models": {
#> "output": {
#> "sql": "SELECT id, \"name\", SUM(amount) AS total_spent\nFROM (\n SELECT customers.*, order_id, amount\n FROM customers\n LEFT JOIN orders\n ON (customers.id = orders.customer_id)\n) AS q01\nGROUP BY id, \"name\"",
#> "engine": "r",
#> "dialect": "duckdb",
#> "namespace": "duckdb"
#> }
#> },
#> "node_count": 3,
#> "edge_count": 3
#> },
#> "nodes": [
#> {
#> "id": "customers",
#> "type": "source",
#> "columns": ["id", "name"]
#> },
#> {
#> "id": "orders",
#> "type": "source",
#> "columns": ["amount"]
#> },
#> {
#> "id": "output",
#> "type": "target",
#> "columns": ["id", "name", "total_spent"]
#> }
#> ],
#> "edges": [
#> {
#> "source": "customers",
#> "source_column": "id",
#> "target": "output",
#> "target_column": "id",
#> "transformation": "identity",
#> "expression": "id"
#> },
#> {
#> "source": "customers",
#> "source_column": "name",
#> "target": "output",
#> "target_column": "name",
#> "transformation": "identity",
#> "expression": "name"
#> },
#> {
#> "source": "orders",
#> "source_column": "amount",
#> "target": "output",
#> "target_column": "total_spent",
#> "transformation": "aggregation",
#> "expression": "sum(amount, na.rm = TRUE)"
#> }
#> ]
#> }lineage_graphml() writes GraphML, the XML format that
graph tools speak: igraph, Gephi, and yEd all open it directly. Every
column becomes its own node, which is what makes real graph queries
possible. The classic one is impact analysis (“if
orders.amount changes, which outputs are affected?”) or its
reverse, tracing an output back to every source column that feeds
it:
path <- tempfile(fileext = ".graphml")
lineage_graphml(lineage, path)
g <- igraph::read_graph(path, format = "graphml")
# Everything upstream of total_spent
igraph::subcomponent(g, "output.total_spent", mode = "in")
#> + 2/6 vertices, named, from 55ab4d9:
#> [1] output.total_spent orders.amount
# Everything downstream of orders.amount
igraph::subcomponent(g, "orders.amount", mode = "out")
#> + 2/6 vertices, named, from 55ab4d9:
#> [1] orders.amount output.total_spentBoth functions return the serialized string when called without
path, so they compose in pipes and tests.
You don’t need a graph library for the everyday questions, though.
lineage_upstream() and lineage_downstream()
answer them straight from the lineage object (pass
"orders.amount" for one column, or "orders" to
trace every column of a table at once), and
lineage_unused() lists the columns with no path to any
target, which is how dead weight in a multi-model pipeline shows up.
Next steps
- The lineage
checks in CI article sets up
lineage_check()as a GitHub Actions gate that fails pull requests on breaking provenance changes -
vignette("python-integration")explains how the Python dependency is managed, and how to use your own environment - The function reference documents every argument
- Found a query that traces incorrectly? Please open an issue with the SQL
