How Hunch Built a Data-Crunching Monster
readwriteweb.com
readwriteweb.com
It's remarkable how quickly distributed computing has come to the forefront of scaling. 3 years ago, before hadoop and horizontal scaling became huge buzz-ideas, this (buying comically large boxes and heavily optimizing your math libs) was the standard practice.
That being said, there is a reason why people have slowly but surely been moving towards distributed systems: If Hunch needs boxes that big to deal with 500k uniques/month, they are going to need to run some seriously insane hardware when they actually grow a large user base.
The only way I've seen it done in Hadoop is to store a big file listing every edge as a pair and do N passes if you want to do a calculation with a distance of N. Not the most elegant approach but with enough horsepower it'll work.
Would love an expert on the subject to chime in.
http://www.royans.net/arch/pregel-googles-other-data-process...
They do try to group related vertices as much as possible, but if edges span computers, that's okay too. I assume that it works best when the interconnect between computers is low-latency, though.
Unless they've thought of some brilliantly simple method that I haven't, which is certainly possible.
So, yeah, they're going to have a lot of cross-node messaging.
If it's general and not secret sauce, mind posting a brief description? Appreciated
For better or worse, we don't use any graph-specific heuristics beyond general assumptions of the graph's structure (which may or may not be correct, but that's another story...). We're dealing with a massive sparse graph whose vertex set, but not edge set, fits in memory (RAM). Embarrassingly, the edges are stored in a database in MySQL; we might be running the world's largest graph on MySQL, but I think we get better performance from MySQL than we could from any other products. Needless to say, we're not using anything relational. Like a lot of graphs (the Web, social networks, etc.), the number of edges is several orders of magnitude larger than the size of the vertex set, so we plan accordingly.
Perhaps my use of the term "partition" above is incorrect; these aren't perfect partitions (they don't exist, except in theory), but are what's known as "quasi-cliques" (the field of graph theory is rife with jargon).
We're happy with what we've come up with; contrary to most, the limits we face are mostly due to limited storage capacity rather than limited computation time or network capacity, which is essential for handling growth.
Sorry if all of this has been too general; there's not too much I can say specifically without revealing our "secret sauce". :D