Skip to content

How Does Facebook's Distributed Like Counter Work?

Naseebullah AhmadiSenior Software Engineer, London

A celebrity posts and a million people tap the heart in ten minutes. One row with a number in it can't take that. The fix splits the problem in two: an exact record of who liked what, and a count that is allowed to run a few seconds behind, spread across slots, batched through a log, and cached for everyone who only wants to read it.

12 min read
#engineering
In one line

A like is two different facts. "Did this person like this post?" has to be exact, and it's a row in a table with a unique constraint. "How many people liked it?" only has to be close, so it's allowed to be slow and a little wrong. Spread its writes across slots or batch them through a log, cache it for readers, and recount it from the exact table now and then to fix any drift.

A celebrity posts a photo. Within ten minutes a million people have tapped the heart, and every one of them sees a number under it that seems to keep up. On average that's about 1,700 likes a second on one post, with the first minute far higher.

Facebook has published parts of how its storage works (the TAO paper describes likes as edges in a graph, with a count kept next to them), but not the whole path a like takes. So this is how you'd build a counter that behaves like theirs, and why each piece is there.

Start with the version that looks fine

A posts table with a like_count column, and an endpoint that bumps it:

@itsnas TypeScript
// POST /posts/:id/like
async function like(req: Request) {
  await db`
    UPDATE posts SET like_count = like_count + 1
    WHERE id = ${req.params.id}
  `
  return new Response(null, { status: 204 })
}
main
Nas (@itsnas)
One row, one number, one lock

This has two problems, and only one of them is about scale.

The first is correctness. Tap the heart twice, or let the client retry a request that timed out, and the count goes up twice. The table knows how many likes happened, not who did them, so it can't tell a second like from a repeated one.

The second is the hot row. In #postgres (or any database with row locks), an UPDATE holds the row's lock until its transaction commits. A second UPDATE on the same row waits for it. So every like on the celebrity's post queues behind the one before it, and the row's throughput is capped by how fast one commit can finish: somewhere from a few hundred to a couple of thousand a second, depending on disk and replication. That's below the 1,700 a second average, and far below the first-minute peak. Meanwhile every other post on the site is fine, because the load isn't spread out. It all lands on one key.

A like is two facts

Pull those two problems apart and they want different storage.

"Did Sam like this post?" has to be exact. Sam sees a filled heart or an empty one, and it has to match what Sam did. That's a row per like, with a unique constraint on the pair:

@itsnas TypeScript
async function like(req: Request) {
  const [inserted] = await db`
    INSERT INTO likes (post_id, user_id)
    VALUES (${req.params.id}, ${req.user.id})
    ON CONFLICT (post_id, user_id) DO NOTHING
    RETURNING post_id
  `
  if (inserted)
    await likeEvents.publish({ postId: req.params.id, delta: 1 })
  return new Response(null, { status: 204 })
}
main
Nas (@itsnas)
The exact fact: one row per person per post

A double tap or a retry now hits the constraint and does nothing, which fixes the first problem for good. It also helps with the second: a million likes are a million different rows, and inserting different rows doesn't make anyone wait on a lock the way updating one row does. Unlike is the mirror image: delete the row, and publish -1 only if a row was actually deleted.

"How many people liked it?" is a different kind of question. Nobody can tell 1,204,311 from 1,204,388, and the page shows "1.2M" anyway. That number is allowed to be a few seconds behind and briefly off by a little. Having that slack is what lets the counter scale, and the event published above is how the count hears about each like without the like waiting on it.

Spread the writes out: sharded counters

The count still has one hot key. The first way to cool it is to split it. Instead of one row per post, keep a fixed number of slots, add each like to a random one, and sum the slots to read the total:

  1. Index 0: 150612
  2. Index 1: 150490
  3. Index 2: 150533 (+1)
  4. Index 3: 150577
  5. Index 4: 150401
  6. Index 5: 150622 (+1)
  7. Index 6: 150548
  8. Index 7: 150528
One post's count split over 8 slots. Each like lands on a random slot, so no single row takes every write. The total is the sum: 1,204,311.
@itsnas TypeScript
const SLOTS = 8
 
async function addLikes(postId: string, delta: number) {
  const slot = Math.floor(Math.random() * SLOTS)
  await db`
    INSERT INTO like_counts (post_id, slot, likes)
    VALUES (${postId}, ${slot}, ${delta})
    ON CONFLICT (post_id, slot)
    DO UPDATE SET likes = like_counts.likes + ${delta}
  `
}
 
async function likeCount(postId: string) {
  const [{ total }] = await db`
    SELECT COALESCE(SUM(likes), 0) AS total
    FROM like_counts WHERE post_id = ${postId}
  `
  return total
}
main
Nas (@itsnas)
Writes pick a slot; reads add them up

Each slot is still a hot row, but it only gets its share of the load. How many slots you need falls out of two numbers: the peak rate and what one row can take.

N=⌈peak writes per secondwrites one row can take per second⌉N = \left\lceil \frac{\text{peak writes per second}}{\text{writes one row can take per second}} \right\rceil

(1)  Slots needed for a hot key.

A peak of 10,000 likes a second against rows that each manage 1,000 needs 10 slots. The cost moves to reads: a total is now a sum over NN rows instead of one. That's why you don't give every post 64 slots just in case. Most posts get a handful of likes a day and one slot is plenty; only posts that start getting hot need more. Past a few hundred writes a second per key, reach for the next idea instead.

Spread the writes over time: batching

Slots spread the writes across rows. Batching spreads them across time, and it goes much further. A thousand likes in the same second don't need a thousand updates. They need one +1000.

That's what the likeEvents stream is for. Put the events on a log like #kafka, partitioned by post id, so every event for one post goes to the same consumer. That consumer adds them up in memory and writes the total once a second:

@itsnas TypeScript
const pending = new Map<string, number>()
 
function onLikeEvent({ postId, delta }: LikeEvent) {
  pending.set(postId, (pending.get(postId) ?? 0) + delta)
}
 
async function flush() {
  const batch = [...pending]
  pending.clear()
 
  for (const [postId, delta] of batch)
    await db`
      UPDATE post_counts SET likes = likes + ${delta}
      WHERE post_id = ${postId}
    `
  await consumer.commitOffsets()
}
 
setInterval(flush, 1000)
main
Nas (@itsnas)
A thousand +1s become one UPDATE

Now the hottest post on the site costs one UPDATE a second, whether it's getting ten likes or ten thousand. And because one consumer owns each post, no two flushes ever race on that row (the only other writer is the repair job at the end of this post), so the slots from the last section usually aren't needed at all.

The price is freshness: the stored count trails the real one by up to the flush interval, plus however far the consumer is behind. For a like count, a second or two is invisible.

  1. Liker to Like API: like
  2. Like API to likes (exact): INSERT, unique
  3. Like API to Event log: +1 event
  4. Event log to Aggregator: by post id
  5. Aggregator to post_counts: +N per second
  6. post_counts to Count cache: refresh
  7. Viewers to Count cache: read count
The like is stored exactly and at once. The count hears about it through the log and catches up a second later.

Reading the count

Writes were the hard part, but reads are the bigger number: a post with a million likes has been viewed many millions of times, and every view shows the count. None of those views should reach the database.

Keep the count in a cache in front of post_counts (Facebook's is built on #memcached, described in Scaling Memcache at Facebook), and let the aggregator refresh it after each flush. A viewer's read is one cache hit. The number can lag the database by a moment, which is the same slack the count already has.

The one person who notices lag is the one who just liked the post. They tap the heart and expect to see 1,204,312, not 1,204,311 for another second. That's handled on their screen, not in the counter: the client fills the heart and adds one to the number it's showing as soon as the tap happens, and the heart's state comes from the likes table (which is exact) the next time the page loads. Everyone else is reading a number that's a second old, and has no way to tell.

Fixing the drift

Every shortcut above can leave the count slightly wrong: a publish lost between the insert and the log, a batch counted twice after a crash. None of that touches the likes table, which is still exact. So the count can always be rebuilt from it:

@itsnas TypeScript
async function reconcile(postId: string) {
  await db`
    UPDATE post_counts
    SET likes = (SELECT COUNT(*) FROM likes WHERE post_id = ${postId})
    WHERE post_id = ${postId}
  `
}
main
Nas (@itsnas)
The exact table corrects the fast one

Run it on a schedule for posts that got likes recently, not for every post on the site. Counting a post's likes is an index scan over its rows, which is cheap for most posts and a few seconds for the largest, so it's a background job, never part of a request. It's also slightly racy: likes that arrive while the scan runs, or an aggregator flush that lands in the middle of it, can be counted twice or missed. Running it again later settles that, the same way it settles everything else.

Putting it together

PieceAnswersExact?Why it's there
likes tableDid this person like it?Yes, unique constraintStops double likes; the source of truth
Sharded counterHow many, with no logYes, but read costs NNSpreads one hot row over NN rows
Log + aggregatorHow many, at any rateSeconds late, can driftTurns thousands of +1s into one write
Count cacheWhat every viewer seesA moment behindKeeps views off the database
Optimistic UIWhat the liker seesFor them, yesHides the lag from the one who'd notice
ReconciliationWhat the count should beYes, eventuallyCorrects whatever drift the rest caused

The heart you tap is exact from the moment you tap it. The number next to it is an estimate that's always catching up, and the design is built so that nobody reading it can tell the difference.