Rohit Swami
India Resume ↗

Writing · Data engineering · 6 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.

4Where 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. Arrow's compact buffers help on top, as do categorical types for the strings that repeat endlessly in biological data: gene and pathway identifiers, organism codes, interaction types.

5Profile 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.

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.