Improving MapReduce with HashFold
stevekrenzel.com
stevekrenzel.com
I've got a small prototype that I used to solve a few problems (most notably: http://www.facebook.com/careers/puzzles.php?puzzle_id=8).
Any criticisms are welcome. If you have any recommendations on presenting the concepts clearer, I'm open to those as well as I'm not sure if I did the explanation justice.
Go on with your explanation. I'm not grokking it yet.
It's written in python. You can run it with './peaktraffic input'. It has 5 different use cases for HashFold: a counter, two different filters, an adjacency list builder, and a cluster finder. The README explains the solution, but does so in terms of MapReduce.
This code wasn't written with the intent of public consumption, so take it as you may.
Any questions, let me know.
i.e. Assume we have three nodes and one key, "our_key", and values for "our_key" are distributed across the 3 nodes:
node1 : our_key : [ 1, 2, 3, 4, 5]
node2 : our_key : [ 6, 7, 8, 9, 10]
node3 : out_key : [11, 12, 13, 14, 15]
and our fold function ( which must be associative) is +. First we fold each node:
node1 : fold(+, [ 1, 2, 3, 4, 5]) = our_key : [15]
node2 : fold(+, [ 6, 7, 8, 9, 10]) = our_key : [40]
node3 : fold(+, [11, 12, 13, 14, 15]) = our_key : [65]
When each node is done with its local inputs, the key-value pairs are redistributed so that each key and its values are on a single node. In this case, assume our hash function moved "our_key" to node 2. Node 2 would now have:
node2 : our_key : [15, 40, 65]
And we would fold against that and get:
node2: our_key : 120
This is how we achieve massive scalability. The nodes can do all of their processing in parallel without any intercommunication except for the key-redistribution phase.
Two key things to point out are that HashFold does the folding as the data becomes available (we don't wait for all values to start folding... this reduces memory usage because we don't need to store the list of values, only the result of the fold as it progresses), and the framework described here is sufficient to perform any computation that MapReduce can do which should make the transition easier (as a lot of major data processing tasks are currently framed for MapReduce)
Do you know Haskell a bit? Try to write your examples in real Haskell, a friendly type system helps a lot to clean up your pseudo-code.
Furthermore, I should have explained the types better. As described, the HashFold framework would require the output of the map, the arguments of the fold, and the output of the fold all be the same type. This isn't as limiting as it sounds and if it turns out to be limiting for some problem domains, there are solutions to each of those.
Not so fast. Yes, mapreduce stores all "unapplied" key-value pairs at the reducer. However, HashFold does as well, the big difference being that HashFold will start applying pairs as it sees them. While that's a win on associative functions, it's at best a tie on unassociative functions.
In many cases you can have a significantly lower memory profile. In a worst case it's the same as MapReduce (as you said). I think having the additional flexibility with memory, in addition to the simpler architecture and performance attributes, makes HashFold an attractive alternative.