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.
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.
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
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.
Would love an expert on the subject to chime in.