Montu Mia's System Design
Data Partitioning

The Great Data Carve-Up

Data partitioning and sharding

After setting up the caching layer (Redis) on the backend, BiralTube started running like new again. The ghostly lag vanished, and users were over the moon. A few months went by like this. Demand for cat videos online is sky-high to begin with, and once people saw how fast BiralTube had gotten, the user count just shot up.

But does happiness ever last long for Montu? The more users grew, the more that old ghostly jam started peeking out again. This time, though, the problem was a little different. On his monitoring dashboard, Montu noticed that reading (Read) data was no trouble at all, the caching layer was handling everything. But the moment a user uploaded a new video, or wrote "Meow meow!" in a comment on some cute cat video and hit like (a Write Operation), the server would start to choke.

The reason was simple. With Redis, Montu had made 'reading' data fast, sure, but for 'writing' data, or saving it permanently, there was still just that one poor database! Under the weight of billions of likes, comments, and video metadata, the database's tables had grown so huge that writing even a single new row was setting off a full-on traffic jam.

tensed-montu

How was he going to clear this jam? "Nope, no more pretending I'm the expert, let me just go to Boltu." No sooner thought than done. Montu sprinted off to Boltu's place.

Seeing Montu out of breath, Boltu gave a little smile. "What now? Write operations on BiralTube are jammed up, right? I've been picking up a bit of lag in your system for a few days, that's when I knew the beating that poor database was taking! Come on, pull up a chair. Today I'll teach you how to wheel a database's tables into the operating theater and cut them open."

Carving up the data

Boltu's operating theater

Boltu began. "Listen, with the CDN and caching you used earlier, you brought the read load on the database down to nearly zero. Since your system was read-heavy, everything ran smooth as butter all this time. But the 'write' operations all keep slamming into that one database instance. Your one database has grown to the size of an elephant now, so how's it supposed to move? Poor thing's overloaded."

— "So Boltu, how do I save this elephant?" Despair was written all over Montu's face.

— "Oh, there's a way, of course! If there's a problem in engineering, there's a solution. To fix this, you have to cut your giant database into several smaller pieces and keep them on separate servers. That makes the data easier to manage, and no single server has to carry the whole world's load alone. The process of breaking a database into these small logical pieces is called Data Partitioning. And when those pieces are placed on separate physical servers, each chunk gets the affectionate name Shard."

operation-theatre

Montu's eyes lit up. "Oh, brilliant! Raising a few small horses is way easier than raising one giant elephant. But Boltu, how do I do this splitting? And how does my backend even know which user's data lives on which server (or shard)?"

Boltu gave Montu a pat on the back. "Why so antsy? I'll tell you everything, take it slow. Listen carefully, there are quite a few styles of data-carving too."

On this page