Skip to contents

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") puts sales on 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 as df/df_1 and arrow_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 OVER clause, so a windowed column’s grouping keys stay indirect on dtplyr and arrow (dashed group_by edges under include_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 by mutate(), rowwise(), and functions with no duckdb translation (case_when(), paste0(), between(), and rename_with() are common ones). The result comes back as an ordinary tibble, or as a duckplyr frame with no relation left behind it, and extract_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 .by mutate whose aggregates pass na.rm = TRUE stays 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_spent

Both 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