Building a Search Index at Kosmix: How We Discovered the Need for MapReduce
(This post is a reflection on our work at Kosmix from 2004 onwards)
In the previous post we talked about how having the ability to have the data in a distributed file system such as Kosmos File System helped us separate where the data resided from where the computation on the data happened. The distributed file system gave us replication and thus fault tolerance from machine failures but allowing the data to be accessed from any other machine then gave us a lot of flexibility — if the node containing the data was overloaded we could move computation on that data to a different node and load balance effectively. However, whenever possible you do not want to move large amounts of data across machines.
In order to build our search index from our web page repository we essentially had to perform a huge external sort on the web page repository — it is easy to see that the web page repository essentially is tuples of the form (webpage, term) with all the terms of a webpage occurring at one node. The inverted index or the search index needs to invert this and create tuples of the form (term, webpage) with all webpages for a term occurring at one node (simplifying the exposition here, in reality we also store what is called the positions list — where in the document the term occurs). Expressed this way, it is easy to see the indexing can be viewed as a sorting problem or a binning problem. We wrote custom code to perform external sorting or external hashing and had to write special purpose code to break computation into the right sized chunks and potentially redo chunks if there are failures. Also we needed to load balance and take care of stragglers. As part of the indexing we also had to compute things like IDF (inverse document frequency) which counts how many documents a term occurs in — this is an aggregate function on the bin corresponding to the term.
The Mapreduce paradigm is a generalization of the above notion and essentially takes a large input and redistributes (maps) it based on some set of keys and then computes (reduces) some function on all tuples belonging to a single key. Notice that this is identical to what a database system would do internally to sort a table (to perform a sort merge join or to respect order by clause in SQL) or to compute a group-by clause in SQL. It was a convenient abstraction for very large amounts of data that was sitting on external storage on which we needed to compute some aggregate function and/or have the result sorted. It is a very low level abstraction however and requires very careful programming.