DEV Community

Cover image for How LinkedIn Rebuilt Its Feed: FollowFeed Explained
Vinayak Gaur
Vinayak Gaur

Posted on

How LinkedIn Rebuilt Its Feed: FollowFeed Explained

Every time you open LinkedIn, your feed is built from scratch. For 400 million members, in about a tenth of a second.

For years that feed ran on a search engine built for something else, until the cost stopped making sense. A team spent twenty months replacing it. The new system, FollowFeed, is five times faster, holds twenty times more data, and runs on half the servers. The trick behind all three numbers is that your feed is never stored at all. The animated version:

A feed assembled per request

A feed is a question

You follow people and companies. Each one produces content: shares, likes, comments. Call each of those streams a timeline: a list of what that entity did, newest first.

Your feed isn't a thing stored somewhere. It's an answer to a question: out of every timeline I follow, what are the best few dozen items right now? LinkedIn's index holds hundreds of millions of timelines. Answering that question, fast, is the whole problem.

The system they outgrew

The old feed ran on Sensei, a generic search platform built on Lucene. Lucene's core trick is the inverted index: a map from every term to the documents containing it. It's fast if you keep it in memory, and Sensei did. That's the catch.

More content meant more memory, which meant more machines, and membership doubled in three years, from 200 million to 400 million. Sensei was also three projects stitched together (Bobo, Norbert and Zoie), and once LinkedIn built a new search engine, Galene, Sensei stopped being anyone's priority. The maintenance bill stayed.

Push or pull

When someone you follow posts, do you push it or pull it?

  • Push (fan-out-on-write): copy the post into the pre-built feed of every follower, immediately. Reads are trivially cheap, because your feed is already sitting there, finished.
  • Pull (fan-out-on-read): store the post once, in its author's timeline, and assemble the feed when someone asks.

Push stores 62 times more data

Push is fine until you count the copies: LinkedIn measured 62 times more data than the pull design. And a pre-built feed is a frozen ranking. Change the relevance model and you have to re-rank everyone's stored feed.

So FollowFeed pulls. Feeds are built at query time, every time.

Built at query time

Take the code to the data

Pulling has a cost: the work now happens while the member is waiting. Personalizing a feed isn't one lookup. The model wants thousands of candidate records, and the naive version ships all of them across the network to a scoring service. The network is the slow part.

So FollowFeed inverts it. Scoring runs on the index nodes, the machines already holding the data. Each node filters and scores its own records, and returns only the handful that survived.

Move the computation to the data

Timelines on disk

If you're not holding everything in memory, you need storage that's fast on disk. FollowFeed uses RocksDB, an embeddable key-value store built on LevelDB and tuned for SSDs. LinkedIn runs Java and RocksDB is C++, so the team wrote the bindings and contributed them back to open source.

A timeline isn't stored as one key per post (that would be one disk lookup per post). It's a linked list of blobs. Each blob holds a chunk of serialized records plus metadata: the timestamp range it covers and how many records are inside. Records use Avro, so each node keeps a schema registry inside RocksDB.

Seven hundred and twenty

Hundreds of millions of timelines don't fit on one machine. The naive choice is one partition per node, but then adding a node means splitting partitions and moving the pieces.

FollowFeed over-partitions instead: the data is cut into 720 partitions, permanently. The number is chosen for its divisors. 720 divides evenly by 40, 45, 60, 80 and 360, so the cluster can be any of those sizes and every node still gets whole partitions. You move partitions; you never split one. Rebalancing becomes bookkeeping instead of surgery.

720 partitions, fixed for good

Two caches

There's a cache in front of RocksDB, and its shape is the interesting part. Some timeline keys map to long lists of records, others to nothing at all. Size a cache by entry count, and a few huge entries blow the memory budget.

So FollowFeed runs two: a fat cache for keys with records, measured in total records (its weight), and a skinny cache for keys with no records, capped by key count. Both are Guava caches, split into sub-caches so threads don't queue on one lock.

One request, end to end

A request arrives at a broker with the entities you follow, your filter criteria and how many records to return. The broker finds the right index nodes through D2, LinkedIn's name service and load balancer (backed by ZooKeeper), and fans the query out in parallel. Each node pulls candidates from cache or RocksDB, filters them with a small SQL-like grammar (privacy, geography, content type), scores and ranks what's left, and returns its best records. The broker merges, drops duplicates, applies diversity filtering, and returns the top few.

Scoring is where the latency budget goes. The team started with a general-purpose scoring library that interprets the model at runtime, then generated Java code for the model instead: 50 microseconds per record at p99, fifteen times faster than where they started.

Keeping it fed, everywhere

Everything that happens on LinkedIn lands on a Kafka stream, but that stream is partitioned by Kafka's rules, not FollowFeed's 720. A service in the middle, the Partitioner, republishes the firehose onto topics matching the index nodes' partition ranges. Adding a replica means starting a node and pointing it at the same topics. Every timeline is replicated to every datacenter, and each datacenter's serving cluster is sized to the read traffic it actually takes.

Fast at the tail

Broker calls are asynchronous, so one slow node can't hold the others hostage, and if a node is slow the broker fires the same request at up to three replicas and takes the first answer. Inside each node: read-copy-update, fine-grained locking, and tuned garbage collection, caches and networking.

The results

What it bought them

  • Mobile feed queries at p99: about 140 ms, five times faster than Sensei
  • Page load at p90: 150 ms better
  • 20 times more data in the index, on half the hardware

Takeaways

  • A feed can be a query, not a table. Pull avoids 62× duplication and frozen rankings.
  • When the network is the bottleneck, move the computation to the data.
  • Over-partition with a number that has many divisors, so rebalancing never splits a partition.
  • Size caches by what actually costs memory, not by entry count.

Watch the full breakdown: https://youtu.be/WlEdcWpXn7M

Source: FollowFeed: LinkedIn's Feed Made Faster and Smarter, LinkedIn Engineering.

Practice designing systems like this yourself at systemdesignlab.in.

Top comments (1)

Collapse
 
devantibot profile image
DEV ANTIBOT •

You need to verify your account .
Link is in the profile.