While buying that much ram sound costly, that's only looking at capacity. Assuming typical 1u servers and common pricing at the moment, you're paying roughly $1k in capital for each 10GB of ram capacity. However, each of these servers gives you a couple hundred thousand random reads per second. To duplicate that with spinning hard drives would take several hundred spindles at least. SSD's are better, but still, it ends up being a lot of devices to duplicate that io capacity.
This is why virtually everyone at stupendous scale (google, facebook, etc) ends up with very ram centric architectures.
There's a simple way to decide how much of what storage you need [2]. Look at the distribution of access times. With current technology, in very rough terms, if an item is accessed more than once a day, it's more cost effective to store it on SSD. If it's accessed more than once in an hour, it's more cost effective to store it in ram.
[1]: http://perspectives.mvdirona.com/2010/07/01/Velocity2010.asp...
[2]: http://www.cs.cmu.edu/~damon2007/pdf/graefe07fiveminrule.pdf
So far I have never heard of any one running a commercial RDBMS reiterate the MySQL-mantra that you need the entire DB in RAM and I find it a very puzzling attitude to efficient database usage.
Most of the time, less than 10% of your DB represents the active working-set of your data and a good database-system should be able to analyse what is being used, gather statistics about the data, use those statistics to intelligently optimize your queries and minimize the need for disk IO.
For any decent RDBMS, having to match total memory with the database size, is wasting money on RAM for very little extra gain in performance and this doesn't really make much business sense.
A concrete example would be a server I manage. It has 16GBs of RAM and handles around 5000 GBs worth of DBs. Due to a good RDBMS and intelligent caching and use of statistics, that server has a cache-hit ratio of 98%. That means that only 2% of queries entering the system results in actual disk IO.
In order to get the last 2% of queries to get into RAM, I would have to increase the system memory by a factor of 30. At that point just getting another server and setting up replication would be much cheaper and also allow a theoretical doubling of troughput.
Really. Having the DB all in RAM is not something which makes marginal sense at all when you compare it to the costs it involves. No amount of internet-argument can convince me that this goal is anything besides a pure waste of money.
"It has 16GBs of RAM and handles around 5000 GBs worth of DBs. Due to a good RDBMS and intelligent caching and use of statistics, that server has a cache-hit ratio of 98%."
This may be your query distribution but it plainly is not everyone's.
"No amount of internet-argument can convince me that this goal is anything besides a pure waste of money."
If Jim Gray's published work doesn't convince you there's not really much left to discuss.
This is actually one of the big long term challenges we're going to have to deal with @ foursquare. Right now we calculate whether you should be awarded a badge when you check in by examining your entire checkin history (which means it needs to be in ram so we can load it fast). While this works now, as we continue to grow it will become more and more of a problem so we'll have to switch to another method of calculating how badges are awarded. Several different options here, each with pluses and minuses.
E.g. for "you've seen foo 4 times today" you'd keep a list of the foos visited. When a new foo comes in you first remove the foos older than 1 day from the list, and insert the new foo. If the list now contains 4 foos you award the badge.
You can use a bitvector to record the badges that have any active state at all, so for badges that have no state associated with them yet (e.g. no foos seen yet) you only pay 1 bit. Or if very few badges have active state on average you could use a list of badges that have state instead of a bitvector, so that you only pay for badges that actually have active state. So total storage is 1 bit per badge or less + a couple of bytes per badge with active state?
Also I recommend giving @naveen a separate shard. The guy just keeps roaming around in the city. Thank you.
This approach also increases the complexity of adding new badges, which is undesirable for product and business reasons.
It could certainly help in some case though, and it's something we're considering.
To the extent that you do this across the board, you'll have yet another tool to defend against being over capacity.
I suppose if your CPU load is not too high you could even run background processes to monitor "near-badge" status for all users with high enough frequency. This you could really get wild with: people with most number of recent check-ins should get scanned first and of course this list should be updated as a separate entity so it stays small and always in memory.
In this way, the amount of current day's events are pretty small and you don't need to keep all historic events in memory.
EDIT: I thought about this a bit more, and the checkins are probably time-sensitive, so another 32-bit timestamp (Foursquare Epoch) would need to be stored for each checkin.
For example, starting at ~9pm the data about check-ins into library may be safely unloaded to disk (the probability of needing this data is low and reading from disk for whose rare cases would do just fine) and be replaced in memory with data about check-ins to clubs, etc...
http://en.wikipedia.org/wiki/Memoization ; see the list of implementations at the bottom - there are ones for Lisp, Python, Perl, Java, etc.
If your disk is spinning at 6000 rpm, a given sector spins past 100 times per second, which makes lookup time anywhere from 0 to 1/100th of a second, or 1/200th of a second on average. Disks only let you look for a limited number of things at once, so you can do less than 1000 disk seeks per second. Period.
If data is living on disk and you have high query volume, it is really, really easy to blow past that limit. The solution is to shard data on multiple machines in RAM. This gives a fixed cost per unit of RAM. As long as you don't exceed available RAM, it works well. Luckily it isn't hard to monitor available RAM and respond in advance. (They didn't do that in this case.)
If you don't care about latency, or have a lower query volume, then you can live with data on disk, and frequently accessed data in RAM. A few well-designed caching layers in front can give you even more headroom. This is much cheaper. But has more complicated failure modes, and they can be tricky to monitor properly.
If you can fit your db into RAM that's the ideal, but of course you also have to have a reliable, persistent backup. The best current compromise is probably PCI based SSD drives.
John Ousterhout (of Tcl fame) and his group at Stanford are working on a project called RAMCloud that's exploring the feasibility of storing data in RAM at all times, using disk only as backup. See his projects page http://www.stanford.edu/~ouster/cgi-bin/projects.php for more info. IMO, given the continuous improvements in RAM prices and network latencies, this type of setup for permanent storage will be the norm in data centers in a few years.
These new projects really seem a lot like re-hashing problems that have already been solved over and over again - people just keep forgetting to do proper systems engineering in the first place. There is no magic bullet.
http://www.infoq.com/presentations/Scale-at-Facebook
In short they use sharded mysql instances (sans joins) as a key-value store with memcache on top of that.
I don't understand why they are running the whole nosql db on those monster machines - it just defeats the purpose. And sharding architecture that they are mentioning is quite error prone. I would go with read/write seperate channels for things of their nature - it appears they don't have that either.
Monitoring your working set is tough.
Not we're not talking necessarily about a ramdisk here or specialized software (though we could be) - we're talking about proper system design with enough ram and tuning to get the response times you need.