Map Reduce: A simple introduction (2010)
ksat.me
ksat.me
In your log files you have the IP which you can use to geolocate those users.
These counts can be caluclated in parallel.
We do this by first "mapping" each line in a log file extracting out the IP. The output of map will be a geolocated hash by region.
Reduce's goal is to calculate an answer based on these keys.
If we think about the output of map being the geo location, we can then collect a count by geo located hash of where our users are coming from.
Map: Map the raw input to some key value space
Reduce: Reduce the output of the map to some aggregate result, think of this like the sql aggregation functions
Hope that helps!
With the books spread over a distributed, in-memory data grid of 20 or so nodes, the system sent a command object to each node to do the calculations in parallel in the same process context as the position data. The data was partitioned so that all positions for each book were in the same node. When the calculations were complete, the reduce process consolidated the results.
VaR requires storing interim results to roll the calculation up the hierarchy of books, but that was straightforward with this approach.
Any time you can parallelize the calculations, map-reduce is worth considering.
Incidentally, the core idea (like all good software ideas) was present in Lisp decades before Google popularized it. The words map and reduce are even in the language.
Your map function scans the rows and outputs the kv pair (zipcode+","+income, csv line). All the csv lines in a group go to the same instance of the reduce function, where you can run any code you like (compute averages, do deep learning, etc). The output is the results of what you want to compute for each group.
This is a pretty simple example, but does demonstrate where the power of mr comes from -- arbitrary functions in the map and reduce functions that are allowed restricted one-way communication from the mapper to the reducer. It also should help you understand the glib "you can implement SQL on mapreduce" comment below, which is what Apache Hive does.
In this analogy, the pregnant woman is the node and the baby is the processed data (of course).
MapReduce is all about dividing tasks as small as they can be and then executing those tasks in parallel in several nodes. Most examples are of counting words because it's very straight forward to explain that you can give one page of a book to N people, have them count the words, and afterwards merge sort the sums into M counters and just sum it again.
https://github.com/snowplow/snowplow
https://github.com/PredictionIO/PredictionIO
We (Snowplow) use MapReduce primarily to:1. Scale our event enrichment process horizontally - raw events come in, we validate them, enrich them (IP -> geo etc), store them. With MapReduce, we just throw more boxes at the enrichment process for larger users (we enrich 200m events in ~90 mins on 6 x c3.2xlarges, spot cost of $0.58)
2. Do easy recomputations across user's full history of raw events - e.g. we add a new enrichment or a user's business logic changes, we can rerun over their full history going back to 2012
Hope this helps!
I think what makes this example so confusing (as well as the other common hadoop example where you do a partial count in the mapper "hello":3, "goodbye":2) is that the keys mean nothing to the final result (unless you really wanted to know the frequency of each word) - they are only used to shard the work.
Sorry I don't have a real-world example (crunching log files as you mentioned is a common one, but it's really just the same thing: counting frequencies)
The shop is now more popular, and a bit bigger too. The volume increases to say 2000 transactions a day, and you say, well, let's use mysql to store this. You are still comfortably generating the report at the end of the year using simple sql queries.
Now suddenly, they decide go really big, to expand and open more stores across the city or state, say about 500 stores. They also expand the items in the stores from just few hundreds to few thousands. They project the transactions across all stores to be in the range of 1,000,000, on average containing about 10 items in each. They also want more reports on which products are doing good/bad, what are the buying habits, what's the trend over Thanksgiving, do year-to-year comparisons on various matrix. And the mysql solution no longer works - storing 10,000,000 rows on daily basis and running those sql queries is turning out to be practically impossible. In a year, you now have 3,650,000,000 rows that needs to be joined with 100,000 items. Yes, you can add more space and more resources to the machine, but running all those queries are now taking hours or days instead of seconds.
This is the point where Hadoop/Map-Reduce comes to rescue. You now have a cluster of, say, 5 machines, each having, say, 64G or RAM and 1 PB of storage, but still costing just around 25000$.
Since you are not familiar with Java or Map - Reduce, and/or don't have time to learn it, you decide to use Hive - an important tool of the Hadoop ecosystem among others - that still let's you access and process the data in the familiar sql query way - but generating map-reduce jobs on your behalf on the nodes. They split the processing across different nodes, bring required outputs together, may run through other map-reduce jobs if required, and ultimately, give you the results that you can use produce those reports - in a reasonable time, of course.
This is very simplified but real-life business use case.
The first part covers MapReduce, the rest you can skip.
The idea is that it's the simple "hello world" of teaching mapreduce.
Likewise for teaching new programming language syntax, the idea of literally displaying the phrase "hello world" is not useful for explaining more real world business uses. Since everybody presumably already knows what "hello world" means, they can ignore that string and instead, pay attention to the surrounding syntax (printf, WriteLine, println, puts, echo, etc) of whatever new programming language they're trying to learn.
Since counting words is very easy to do without mapreduce (using dictionaries or associative arrays) and it also doesn't require any particular business domain knowledge, people can ignore it and just concentrate on the structure of setting up mapreduce.
In that context, the "uselessness" of counting word frequencies makes it easier to isolate the learning of mapreduce.
glennengstrand.info/analytics/oss
Here you will see a blog, paper, and github repo of some map reduce jobs that take openly available San Francisco crime data and load it into an OLAP cube.
While I can agree that pointing out the spelling and/or grammar mistakes could be constructive, calling the article "un-shareable" because of them seems terribly over-the-top.
I can appreciate numerous spelling/grammar errors making you analyze an article with a bit more scrutiny. However, I can't think of a single example of an article with those kind of errors where I couldn't figure out from the content whether I thought the author really knew what they were talking about.
At any rate, I think there is a big difference between an article needing a little bit more scrutiny, and the article being "unsharable".
Just to be clear, I don't want to give the impression that I don't think spelling/grammar matter (even for blog posts). I just think it is easy to get so pedantic about it that you place far too much weight on them.
I have my own set of pet peeves, and am probably guilty of allowing the violation of one of them to taint my view of an article too quickly and too often.
And of course, there are always those articles that are so bad grammatically that it looks like a first-grader wrote them (but I don't think that's the kind of article we were discussing).
I try not to judge an article/blog post on the quality of the grammar when the ideas they are trying to convey are solid and useful.
There are some exceptionally fantastic minds out there that are able to succinctly explain, sometimes difficult, subjects.
I appreciate them for at least trying, and if they explained it in a language that is not their first language, all the more respect to them.
Which is a stream based processing system.
I have seen one paper on a streaming MapReduce solution though.
Have you seen any references to it in the wild other than the Google Research paper?
They do definitely seem to have switched from MapReduce though at least - http://www.theregister.co.uk/2010/09/09/google_caffeine_expl...
You can find it at http://www.reddit.com/r/explainlikeimfive
http://www.michaelnielsen.org/ddi/how-the-bitcoin-protocol-a...
He designs his own cryptocoin to explain the choices made by bitcoin and how it works.
Interesting article, thanks.