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:

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:

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:
- Index 0: 150612
- Index 1: 150490
- Index 2: 150533 (+1)
- Index 3: 150577
- Index 4: 150401
- Index 5: 150622 (+1)
- Index 6: 150548
- Index 7: 150528

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.
(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 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:

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.
- Liker to Like API: like
- Like API to likes (exact): INSERT, unique
- Like API to Event log: +1 event
- Event log to Aggregator: by post id
- Aggregator to post_counts: +N per second
- post_counts to Count cache: refresh
- Viewers to Count cache: read count
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:

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
| Piece | Answers | Exact? | Why it's there |
|---|---|---|---|
likes table | Did this person like it? | Yes, unique constraint | Stops double likes; the source of truth |
| Sharded counter | How many, with no log | Yes, but read costs | Spreads one hot row over rows |
| Log + aggregator | How many, at any rate | Seconds late, can drift | Turns thousands of +1s into one write |
| Count cache | What every viewer sees | A moment behind | Keeps views off the database |
| Optimistic UI | What the liker sees | For them, yes | Hides the lag from the one who'd notice |
| Reconciliation | What the count should be | Yes, eventually | Corrects 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.