Designing a Distributed Counter: Why Global Counts Don’t Scale josedacruz, September 10, 2026September 10, 2026 TL;DR: A single database row that stores a count (likes, views, votes) works fine until enough people hit it at once — then it becomes a bottleneck that slows down or breaks your whole app. The fix isn’t a bigger database, it’s splitting the count into several pieces that get combined later. Here’s how that works and where it can bite you. The problem Imagine you’re building the backend for a short-video app. Every video has a “likes” count shown right under the play button. The obvious way to build this is one column in your database: a row per video, with a likes integer that goes up by one every time someone taps the heart icon. This works great for months. Then one of your creators posts a video that gets picked up by a celebrity account and starts going viral. Suddenly thousands of people are tapping “like” on that exact video, every second, at the same time. Each like triggers something like this in your database: UPDATE videos SET likes = likes + 1 WHERE id = 42; That single line looks harmless. But your database can only apply one update to that exact row at a time — it has to read the current value, add one, and write it back, without two updates stepping on each other. When ten thousand of these arrive in the same second, they all queue up waiting for the same row. Your API starts timing out, not because your servers are overloaded, but because they’re all stuck waiting in line for one tiny piece of data. This is called a “hot row” or “hot key” problem: one specific piece of data becomes so popular that it turns into a traffic jam, even though the rest of your system has plenty of capacity to spare. Why it happens The root cause is that a counter is, by definition, a single shared number. Every single person who likes the video needs to agree on what the current count is before adding to it. That agreement has to happen somewhere, and if that “somewhere” is one row in one database table, then every update has to pass through that one row, one at a time. It’s a bit like a popular checkout lane at a grocery store. Doesn’t matter how many cashiers (servers) you hire elsewhere in the store — if there’s only one register that can ring up a particular item, everyone buying that item stands in the same line. All traffic hitting one row creates a bottleneck; splitting the count across shards spreads the load. Most teams don’t notice this in testing because test traffic is never concentrated enough to matter. It only shows up in production, under real, bursty load — which is exactly when you can least afford an outage. The better approach The standard fix is called a sharded counter (sometimes “distributed counter”). Instead of one row holding the total, you keep several rows — say, 20 — each holding a partial count for the same video. When someone likes the video, your app picks one shard at random (or based on something like the user’s ID) and increments just that one. UPDATE video_like_shards SET count = count + 1 WHERE video_id = 42 AND shard_id = 7; Because there are 20 shards instead of 1, you’ve turned one traffic jam into 20 smaller lines, each handling roughly a twentieth of the load. To show the total likes count on the video, you add up all 20 shards: SELECT SUM(count) FROM video_like_shards WHERE video_id = 42; That sum is a bit more expensive to compute than reading one row, so in practice teams don’t run it on every single page view. Instead, they either cache the total and refresh it every few seconds, or run a background job that keeps a separate “display count” column updated periodically. The person looking at the video doesn’t actually need the count to be accurate to the exact second — “1.2M likes” updating every 10 seconds looks identical to a human as updating instantly. Writes go to a random shard; a background job periodically rolls the shards up into a cached, display-ready total. This pattern shows up everywhere once you know to look for it. Big tech counters for likes, views, and follower counts almost never read live from a single source of truth — they’re built from shards or logs that get summed up asynchronously. The number you see on screen is usually a few seconds (or longer) behind reality, and that’s a deliberate trade-off, not a bug. If your traffic is spiky rather than steady — normal most of the time, but occasionally slammed by a viral moment — you can even make the shard count dynamic: start with 1 or 2 shards for a normal video, and automatically split into more shards if the write rate for that specific video crosses a threshold. This way you’re not paying the overhead of 20 shards for videos nobody watches. More shards raise the ceiling on sustainable write throughput — but with diminishing returns, so there’s rarely a reason to over-shard. Pitfalls to avoid Sharded counters solve the hot-row problem, but they introduce new decisions you have to get right: Don’t shard everything by default. Most counters in your app — like a user’s total post count — will never see enough concurrent writes to need this. Sharding adds real complexity (more rows, aggregation jobs, eventual consistency). Only reach for it once you’ve actually measured contention on a specific counter, not as a precaution everywhere. Decide how “live” the number needs to be. A like count can lag by seconds. A count of items left in stock during a flash sale usually can’t — showing “in stock” when there are actually zero left leads to orders you can’t fulfill. Match your aggregation frequency (or lack of caching at all) to how much staleness is acceptable for that specific number. Watch for double-shard problems during scale-up. If you dynamically add shards for a suddenly-popular item, make sure old and new shard counts both get included in the total — it’s easy to accidentally “lose” counts that were written to a shard your aggregation job doesn’t know about yet. Remember shards don’t fix “how many unique people liked this.” A sharded integer counter is great for “how many,” but if you also need to know exactly who liked something (to show “liked by Alice and 41 others,” or to let someone un-like), you need a separate table of individual like records anyway. The counter and the detailed record are usually two different pieces of data solving two different problems. The bigger lesson generalizes past counters: any time you see “everyone updates the same one thing,” ask what happens when a thousand people try to do it in the same second. Usually the fix isn’t a faster single thing — it’s splitting the thing into pieces that can be updated independently, and accepting that the combined view will be slightly, harmlessly behind reality. Related architecture case studydatabase designdistributed systemsscalabilitysystem design