System Design103 min total · 14 parts
System Design Fundamentals for Interviews: Scalability, Trade-offs, and the Framework Interviewers Actually Grade
Part 6 of 14 · ~5 min
Database Scaling
A single database, same as a single app server, eventually runs out of headroom — but unlike an app server, you can't simply clone it and split traffic between identical copies, because a copy that's drifted out of sync with the original has stopped being a copy of anything real. Scaling a database is really a series of decisions about which correctness guarantees you're willing to loosen in exchange for more room.
Scaling Reads with a Leader and Its Followers
Leader-follower replication hands one database exclusive control over every write — that's the leader — while one or more followers continuously copy whatever changed on the leader and take on read traffic of their own:
Customer writes ----> Leader --(streams every change downstream)--> Follower A
| -> Follower B
| -> Follower C
Customer reads ----> whichever Follower is handy (or the Leader itself)
This lines up almost perfectly with how Fanline's own traffic splits: browsing shows, checking a seating chart, searching by city all dwarf the actual moment someone finishes a purchase — often by fifty or a hundred requests to one. Read replicas let that browsing volume scale out horizontally, adding followers as the audience grows, while write throughput stays capped by whatever the single leader can carry — which is fine, because the leader was never where the pressure was; browsing was. Should the leader go down, a follower usually steps up to take its place, but that promotion isn't instant, and whatever writes hadn't finished copying over at the moment of failure can end up gone for good, depending on exactly how replication was set up beforehand.
The trade baked into this whole design: replication is rarely instant, so a follower can hand back stale data for a short window right after a write lands on the leader. A customer who just added a seat, hits refresh, and gets routed to a follower still catching up might briefly see their own cart as empty — a quieter cousin of the cross-server cart bug from two chapters back, caused by the same root issue: two places claiming to hold the truth that haven't finished agreeing with each other.
Sharding and Partitioning
Replicas scale browsing, but every one of them still carries a full copy of everything, and writes are still bottlenecked on one leader. Sharding goes after the deeper limit by cutting the data itself into pieces spread across separate database instances, each owning only a slice of the rows — which grows both storage and write capacity together, since different shards accept writes for their own slice with zero coordination between them.
| Approach | How rows get split | Strength | Weakness |
|---|---|---|---|
| Range-based | Rows land on a shard based on where their key falls in a range — shows in Jan–Apr on shard A, May–Aug on shard B | Range scans ("everything on this weekend") stay cheap, hitting one or a few contiguous shards | Traffic that isn't evenly spread across those ranges can turn one shard into a hot spot |
| Hash-based | The shard key gets hashed, and the hash decides the shard | Spreads load evenly no matter what pattern exists in the raw keys | A range scan now has to touch every shard, since neighboring keys land all over the place |
Once Fanline expands past a handful of local venues into full regional coverage, it shards by venue region using a hash, deliberately — "what's playing near me this weekend" is asked constantly, while "list everything nationwide in date order" almost never is, so hash-based sharding trades away exactly the query Fanline barely uses for the one it depends on every single page load.
Then Marlow Vance's reunion tour gets announced across eleven cities in one morning, and a single shard — the one holding the two largest metro venues on the routing — absorbs roughly ten times the write volume every other shard sees that week, purely because that's where the demand happened to land. That's a hot shard: the whole scheme was designed around an assumption — that demand spreads roughly evenly once you hash by region — and one extraordinary week simply broke it. Fixing it isn't about adding more shards and hoping the imbalance smooths itself out — it won't, because the imbalance was never about shard count. It's about rethinking which key owns the traffic, which here means carving those two metro venues onto a shard of their own.
The Operational Complexity Sharding Introduces
Sharding tends to be Fanline's last resort, not its first move, because the real costs never fit neatly into a short summary:
- Cross-shard queries that used to be one query against one database now mean hitting several shards and stitching the results together in application code — once a customer's orders span venues sitting on different shards, assembling their full purchase history is a job the database can no longer do for you on its own.
- Cross-shard transactions give up the atomicity a single database hands you automatically — moving a held seat between two shards' inventory during a venue-merge migration needs the same multi-step coordination (two-phase commit, or a saga of reversible steps) that any transaction spanning separate services needs.
- Resharding, once the original scheme stops matching reality — exactly what happened after the Marlow Vance announcement — is a migration nobody enjoys running: rows have to move between shards while the system stays up and serving traffic, ideally borrowing the same consistent-hashing trick from the load-balancing chapter so only the truly necessary slice of data actually has to move.
- Hot shards, as above, show up whenever real traffic doesn't match the pattern the sharding scheme assumed, and the fix is almost always rethinking the shard key itself, not padding out the shard count.
Common mistake: reaching for sharding the moment "the database is struggling" comes up, when the real problem is read traffic that replicas and caching would have absorbed for a fraction of the operational overhead. Fanline stuck with vertical scaling and replicas until write volume and total data size genuinely couldn't fit on one leader anymore — not the moment growth started feeling uncomfortable.