An Unorthodox Approach to Database Design : The Coming of the Shard
highscalability.com
highscalability.com
In Erlang you also get the ability to have the same database on multiple computers or fragment the database into different servers that hold only a certain piece.
It certainly was a trick to unlearn normalization techniques I had been using for years, but once you get he hang of it the data really becomes or efficient.
Effectively, sharding takes advantage of the locality of information. As soon as you need information outside of that locale, you need to make a more expensive request.
Big sites that are primarily read-hungry often don't need to shard at all, instead relying upon replication. Look at Craigslist - they have a single master database, about 20 slaves, and an archive cluster for postings that nobody looks at anymore. Or PlentyOfFish - 1 machine for billions of pageviews. Social sites with lots of user-created content need to shard much earlier - LiveJournal has a fraction of the users as Craigslist, but they had to start sharding around 2002/2003.
There're locality concerns as well, but they tend to be secondary to the write issue. Disks are big, and with good caching you can often avoid hitting the disk at all.
Is that really true? I find it hard to imagine...but something deep down inside of me says "Yes it's possible". Oh my...