Rohit Swami
India Resume ↗

Writing · Data engineering · 8 min read

Ten times faster, mostly by doing less

I rebuilt a graph engine from R to Python and Polars, and it got ten times faster on 80% less memory. Very little of that came from changing languages.

I rebuilt a graph engine for biological pathway data, moving it from R to Python and Polars. The new one runs ten times faster on 80% less memory, and it scaled to more than a hundred thousand entities across over a thousand pathways.

It would be easy to tell that as a story about R being slow. It isn't one. R is a good language for the statistics it was built for. Code is slow for reasons that mostly have nothing to do with the language, and that would be just as slow in Python written the same way. This post is about three of them, because between them they show up in almost every slow data pipeline I've been handed.

1Growing things one row at a time

The most common pattern in slow data code looks harmless:

edges <- data.frame()
for (p in pathways) {
  edges <- rbind(edges, read_edges(p))   # copies every row so far, every time
}

A data frame can't grow in place. Each rbind allocates a new one big enough for the old rows plus the new ones, copies everything across, and lets the old one go. The first iteration copies one chunk, the second two, the hundredth a hundred. The total work grows with the square of the number of chunks, and so does the churn, because for a moment both copies are alive and the garbage collector has to clean up after every one of them.

chunks

rows copied –necessary –waste –

Fig. 1 Each bar is one iteration, and its height is how many rows that iteration copies, at a thousand rows a chunk. Appending in a loop draws a triangle; its area is the work. Binding once copies every row exactly once, at the end.

The fix is to stop growing things. Collect the pieces, then join them once at the end, which copies every row exactly once:

edges <- do.call(rbind, lapply(pathways, read_edges))   # one bind, at the end
edges = pl.concat([read_edges(p) for p in pathways])     # the same idea in Polars

In R the copying is easy to miss because of copy-on-modify: an object is copied only when something changes it while another name still refers to it. That's what makes R's value semantics cheap most of the time, and it's also what lets a harmless-looking loop go quadratic without anything in the code saying so. Python has its own version of the trap. Concatenating pandas DataFrames inside a loop is the famous one.

2Columns, not rows

The second problem is the shape of the work. It's natural to think about a graph one record at a time: for each node, look up its fields, compute something, append the result. It's also the slowest way to ask a computer to do arithmetic, because a CPU doesn't read memory one value at a time. It reads it in cache lines of 64 bytes, and it is very fast when the next thing it needs sits right next to the last.

Polars keeps data by column, in the Apache Arrow memory format. All the values of one column sit together in a single contiguous buffer. Ask for the sum of a column and the CPU streams through exactly the bytes it needs, several values per instruction. Ask the same question of data laid out record by record, and every cache line it loads is mostly fields it didn't want.

cache lines, by row –by column –useful bytes, by row –

Fig. 2 The same 64 records, eight 8-byte fields each, laid out two ways. The highlighted field is weight. Each scan reads the cache lines its sum needs, and both start at the same moment.

This is also why vectorised R is fast. sum(x) on a numeric vector is a tight loop in C over contiguous memory. The slow version is the one that loops in R itself, a value at a time, building results row by row. Moving to Polars pushed the code into column-at-a-time expressions everywhere, which is most of what "vectorising" means in practice.

3Let the query planner do the boring part

The third problem is order. Eager code runs its steps in the order they're written. Read everything, join it, and only then filter to the rows you need and select the columns the next step uses, and every intermediate result carries data that a later step throws away.

Polars has a lazy API. Instead of running each step immediately it builds a query plan, and optimises the whole plan before it executes anything. Two of those optimisations do most of the work. Predicate pushdown moves filters as close to the data source as it can, so rows that will be dropped are never read or joined. Projection pushdown does the same for columns, so a Parquet scan reads only the columns the query uses. Then the plan runs in parallel across every core.

import polars as pl

counts = (
    pl.scan_parquet("edges.parquet")
      .join(pl.scan_parquet("nodes.parquet"), on="node_id")
      .filter(pl.col("organism") == "hsa")
      .group_by("pathway")
      .agg(pl.len().alias("edges"))
      .collect()
)

Written like that, the filter comes after the join, which is where a person would naturally put it. The optimiser moves it down into the scan of the table it belongs to. .explain() prints the plan Polars intends to run, and reading it is worth the minute whenever a query surprises you.

rows into the join –cells read –

Fig. 3 The query above, as written and as Polars plans it. Line width is rows flowing between steps. The table sizes are invented for the figure; the shape of the saving isn't.

4Joins are where the time hides

A graph engine joins constantly: edges to the nodes at either end, nodes to the pathways they belong to, interactions to their annotations. How a join is done matters far more than which language does it. The naive way, for every row on one side scan every row on the other, does n × m comparisons, and with a hundred thousand rows on each side that's ten billion of them.

A hash join does n + m. It builds a hash table on the smaller side, keyed on the join column, then streams the larger side through it, and each row finds its matches with a single lookup instead of a scan. Polars does this by default, in parallel across cores, but it's only as good as what reaches it. Pushing filters below the join, as the planner did above, means the hash table is built on the rows that survive rather than all of them.

work so far –at 100k × 100k rows –

Fig. 4 Sixteen edges joined to twelve nodes. The nested loop checks every pair; the hash join drops each node into a bucket by the hash of its id, then sends each edge straight to its bucket. The readout scales both to a hundred thousand rows a side.

5Where the memory goes

Memory follows the same logic as time. Growing a table in a loop means two copies alive at once, over and over. Joining before filtering means building an intermediate table of rows nobody keeps. Reading whole files means holding columns no step uses. Remove those, and most of the peak goes with them.

R adds a trap of its own: copy-on-modify. Change one value in a data frame that something else still refers to, including the copy a function received as its argument, and R copies the whole column, or the whole frame, before writing. It's what makes R code safe to reason about, and it's also why memory can double without a single line that says "copy". Arrow's compact columnar buffers help on top, and so does storing repeated strings once. Biological data repeats strings endlessly: gene and pathway identifiers, organism codes, interaction types. A categorical column keeps each distinct string a single time and stores a small integer code per row instead.

distinct values –plain strings –categorical –

Fig. 5 A model, not a measurement: a million rows of 12-character identifiers, with 16 bytes of per-string overhead for plain strings and a 4-byte code per row for the categorical. The fewer distinct values a column has, the more a categorical saves; with every value distinct, it costs a little more, because each row still carries its code. Hover or tap the chart to read any point.

6A graph is two tables

It helps that a graph fits a dataframe engine better than it looks. A graph is a table of nodes and a table of edges, and most of what a pathway engine computes is a join or a group-by on them. A node's degree is a group-by on the edge table. Its neighbours are a join. Two hops out is the same join twice. Written that way, every step is a column operation the engine can plan, parallelise and keep columnar, instead of a Python loop over objects:

degree = edges.group_by("source").len()                       # how connected each node is

two_hop = (
    edges.join(edges, left_on="target", right_on="source", suffix="_2")
         .select("source", pl.col("target_2").alias("two_hops_away"))
         .unique()
)

7Profile first

None of this is clever, and that's the point. A rewrite like this gets ten times faster because it does less: fewer copies, fewer bytes moved, fewer rows carried through steps that are about to drop them. Changing languages was the occasion for removing that work, not the reason things got fast.

The advice I'd give anyone handed a slow pipeline is the advice I keep relearning myself: profile it, then fix what the profile shows. It's almost never what you'd have guessed, and it is very often something growing in a loop. In R, profvis shows where the time goes line by line, and tracemem tells you when an object has been copied. In Python, py-spy samples a running process without touching its code. And Polars will show you its plan before it runs anything, which is the quickest way to see whether a filter reached the scan:

print(query.explain())          # the optimised plan, without executing it
# FILTER and PROJECT should appear inside the scans, not above the join

The other half of a rewrite is proving it's right, because a pipeline that's ten times faster and occasionally different is worse than the slow one. The only check I trust is running both, side by side, on the same real inputs, and diffing what comes out: the same rows, in any order, with numbers equal to a sensible tolerance. It's slow and boring, and it's the reason a faster engine can replace the old one rather than sit next to it.

I'm Rohit Swami. I build the unglamorous machinery real products run on: data pipelines, real-time services, open-source tools, and products of my own. More about me, or write to me.

The figures on this page are simulations written for it. They run in your browser, and the numbers in them are illustrative unless the text says otherwise.