Sunday, February 12, 2017

Effect of Different MapReduce Configurations on Identifying Trigrams From Text - Guest Post by Aviral Kumar Gupta

Trigrams - what are they?
An n-gram is typically a contiguous sequence of words extracted from a given text or speech. In the fields of text mining and analytics, n-grams are used to develop analytical models that can be used to carry out a variety of natural language processing tasks like spelling correction, search and text summarization.

Trigrams are a special case of n-grams where n=3, which means we need to pick three consecutive words from a sentence to form a trigram. Suppose the text you are given is:

The quick brown fox jumps over the lazy dog.

Some of the trigrams from this sentence are (The quick brown), (quick brown fox), (the lazy dog), etc. The sequence of trigrams is considered only from left to right and not the other way round. Each sentence has N-2 number of trigrams where N is the number of words in the sentence.

How do we extract the trigrams?
In our class project, we were using MapReduce programming to extract trigrams from the given text on a Hadoop Distributed File System (HDFS). MapReduce works on (key, value) pairs. Trigrams act as the ‘key’ and their count in the data set is the ‘value’ we are estimating. The mapper breaks down the text strings into trigrams and assigns a count of ‘1’ to each trigram. The intermediate keys are subsequently passed to the reducer and it shuffles and sorts them to count the number of occurrences of each trigram in the entire data set.

Here is the ‘Mapper code’ to extract the trigrams:

Given below is the ‘Reducer code’ for counting the number of occurrences for each trigram:


Now that we have extracted and counted the trigrams, we can experiment with the effect of different mapper and reducer configurations and study its effect on execution time.

Aims
In this example, we are interested in analyzing the execution time of the MapReduce program with different number of mappers and reducers involved while running the tasks on the University at Buffalo’s high-performance computing cluster – Lake Effect.
Lake Effect is a subscription-based Infrastructure as a Service (IAAS) cloud provided by the Center for Computational Research (CCR) for UB’s research purposes. Researchers at UB can access this cloud and make use of it for high performance computing requirements like Big Data analytics.

Hadoop was pre-configured on the cluster - in particular, there was one master and two data nodes. Jobs could be submitted remotely to the JobTracker. There were 120 students who were running experiments on this cluster during the course of the semester. This allowed testing the impact of other jobs on the task at hand and experimenting with scalability of the system, albeit on a small scale.

Data set

Our data set comprises of a sample of 50 different articles from the Chronicling America website (http://chroniclingamerica.loc.gov/). The data sets were extracted using the Web Scraper tool that merged the 50 extracted articles in to a single file. In order to ensure that the reducers are optimally utilized, we replicated 20 of these files using the same tool so that the frequency of trigrams increases. We merged both the files (50 articles and 20 articles) and created a single file containing data of the 70 articles.


In order to run our code using multiple mappers and multiple reducers, we need to configure these changes in our ‘main method’ of the code.


Below is the ‘Main method’ of the program:


The important thing to notice about the code is that it offers the flexibility to set the number of reducers and mappers on which the tasks would be executed on the CCR cluster. The user can assign these inputs as command line arguments at the time of command execution. For e.g:

hadoop jar Trigram.jar trigram \input \output 2 3

2’ in red is the number of reducers and ‘3’ in green is the number of mappers.
 The number of mappers are spawn based on the number of input splits in the data set. Therefore, here the number ‘3’ is actually specifying the number of input splits that the code should create in the input file, in order to create as many mappers. The default splits occur according to the block size, which is 128 MB for CCR, but as our input data set is not so big, we have parametrized the number of input splits so that the user can run the code with multiple mappers.

Empirical Results
We have carried out an experiment whereby the same input file is executed with multiple combinations of reducers and mappers.We calculated the execution time of each run by adding the ‘Total time spent by all maps’, ‘Total time spent by all reducers’ and the ‘CPU time’ (refer Table 1).
Table 1: Total and Average execution times for all the mapper-reducer combinations
Below is the graph where we have mapped the mapper-reducer pair vs the execution time taken by each run with that configuration.





Discussions
Here are our observations:
a.     Best execution time was observed when only 1 mapper and 1 reducer were used, and worst case was when 4 mappers and 4 reducers were used.
b.     For 1 mapper, as the number of reducers are increasing the total time of execution is also increasing.
c.      Similarly, as the number of mappers increase, we can see a steady trend in increase of execution time.
d.     The slope and height of the graph for reducers is almost the same across all combinations with mappers.

From the above observations, we can conclude that the execution time of a program increases with an increase in the number of mappers or reducers, when the data set is not very large.

As our data set was very small, adding more than one reducer only results in increasing the overhead for the program. By increasing the number of reducers, we are adding an extra time cost of instantiating the reduce task as well as the network transfer and parsing time which would be spent in transferring data to different reducers, while there is not enough data for each reducer to work upon. In addition, each reducer creates an output file, so, if multiple reducers are being used, for a small data set, then multiple output files would be created. This would again add to the overhead, as extra I/O operations would be required to create these multiple reducer output files.

In another scenario, if the input data set to be processed were very large, then running the program with a single reducer would have affected the performance adversely. The entire load would have gone to a single reducer. In such cases, using multiple reducer to divide the tasks would have been beneficial.

Hence, the number of reducers used for a map-reduce code is a significant factor. Having too many or too few reducers would hamper the productivity of the program.
While using the CCR cluster for executing our jobs, we can expect long waiting time at some points in time, when the cluster would be busy executing other tasks and our jobs would be in the queue. Although the actual execution time of the programs would be in the order of seconds, but the elapsed time for executing the jobs might be much higher, almost an hour in some cases. In case the cluster fails to respond due to some configuration challenges or power failure, the queued jobs could delay with infinite waiting time.

About the author: Aviral Kumar Gupta is a graduate student in the MIS Department, University at Buffalo, NY.

Friday, February 3, 2017

The metaphysical world of Anne Dillard.

An excerpt: 
Last year I saw three migrating Canada geese flying low over the frozen duck pond where I stood. I heard a heart-stopping blast of speed before I saw them there; I felt the flayed air slap at my face. They thundered across the pond, and back, and back again. I swear I have never seen such speed, such single-mindedness, such flailing of wings. They froze the duck pond as they flew; they rang the air; they disappeared. I think of this now, and my brain vibrates to the blurred bastinado of feathered bone. "Our God shall come" it says in a psalm for Advent, "and shall not keep silence; there shall go before him a consuming fire, and a mighty tempest shall be stirred up round about him." It is the shock I remember. Not only does something come if you wait, but it pours over you like a waterfall, like a tidal wave. You wait in all naturalness without expectation or hope, emptied, translucent, and that which comes rocks and topples you; it will shear, loose, launch, winnow and grind. 

I have gutted on richness and welcome hyssop. This distant silver November sky, these sere branches of trees, shed and bearing their pure and secret colors - this is the real world, not the world gilded and pearled. [...] I am buoyed by a calm and effortless longing, an angled pitch of the will, like the set of the wings of the monarch which climbed a hill by falling still.  
                                          - Anne Dillard, The Pilgrim at Tinker Creek.



Monday, October 31, 2016

Finding bigrams using Map Reduce

Natural language processing and computational linguistics applications often use an n-gram for analyzing textual data. An n-gram is a contiguous sequence of n items from a given text. If n=2,  is called a bigram.
Suppose the text you have is

It is raining outside.

The bigrams from this text are (It is), (is raining), (raining outside). Usually the sequence is considered from left to right and not in the backward direction. So (is,It) would not be regarded as a valid bigram. 

Suppose now, you had a lot of text from which bigrams have to be extracted. Why should this happen? There could be several reasons - you want to understand the semantics of text or conversation; you are interested in studying the co-occurrence statistics of words (a.k.a which words tend to occur together frequently); you want to use bigrams as a primitive for more involved natural language processing tasks and perhaps there are other related problems. 


Now if the text is large (such as a whole book or newspaper articles for a whole month/year) the simple task of extracting bigrams tends to become compute intensive. 



A Distributed File System (DFS) is typically used to store the large volume of data. Hadoop DFS is open source and has a relatively easy learning curve. Hence it is a popular choice for storage. Manipulating this data is primarily done with MapReduce programming. It has been used efficiently at Google, Inc., Yahoo!, Facebook, and many other companies that deal with large data volumes. 

So how to extract our bigrams using MapReduce?

Here is a possibility - the Map task will take in (key, value) pairs. In our example, the Map task could take in and output a series of intermediate keys of the form (bigrams, count). The Reduce task then adds up all the counts associated with a certain bigram. 


Straight forward enough?


Here is the Mapper code to do the job. 

So what does the code do? It simply tokenizes the string and puts together the previous token along with the current token as long as the previous one was not null. It outputs (bigram, count).

Now the Reduce code can work as follows. 
It keeps a count of the number of bigrams seen so far.

Putting it all together - we have the main method as follows.


The main method is slightly more involved. First, we input three arguments - the input path where the text will be stored; the output path where the generated bigrams will be stored (these paths will be on the HDFS) and the last is the number of reducers to use. Lines 13-15 simply state that you are running the MyBigramCount program. 

Job (Line 13) refers to the MapReduce job configuration. It can be usedto present the MapReduce job to the Yarn execution framework which comes built in Apache Hadoop 2.0 and later versions.
FileInputFormat indicates where the input files are available for the Map task and FileOutputFormat gives the location of the output on Lines 19-20. Job is used to set the Mapper, Reducer and Combiner (is used) implementations (Lines 22-24). Additionally, it can be used to specify which Comparator to be used, whether files should be put in the DistributedCache, whether intermediate and/or job outputs are to be compressed (and how), whether job tasks can be executed in a maximum number of attempts per task and other details(Lines 22-28).


To compile the program, execute
hadoop com.sun.tools.javac.Main (path-to the location of MyBigramCount.java)
Assuming there are no compile errors, create a jar file that contains all the class files for MyBigramCount using the command below:
jar cf MyBigramCount.jar MyBigramCount*.class
Execute using the command
hadoop jar MyBigramCount.jar MyBigramCount Input-Dir Output-Dir

Saturday, October 22, 2016

A little indulgence on a Saturday afternoon.

Had been craving for crab curry for quite a while. Experimented with Alaskan king crabs and a Goanese recipe that is well reviewed by crab lovers. Loved the taste. Come home if you want to try some. Warning: It is disappearing fast!

Crab curry

Friday, October 21, 2016

Big Data and the Mortuary

Sounds gross, right?

But here's the story. I get on the "big data" bandwagon.
I develop these nice scalable algorithms for learning from the big data. I test them on data that does not fit in the memory of my personal laptop (otherwise the reviewers for my journal papers do not agree that I am doing GOOD research). I use a cluster for all of this work - which I have spent hours configuring and playing around with.

But then I decide to relocate. With all my tangible possessions. How do I move this large cluster?

Needless to say, I am emotionally attached to it, by now. It was not easy getting it up and running in the first place. Then I did all the dirty data cleaning work on it. That was another 100+ hours of my existence. And now I am facing this scenario -- I do not have the physical space to store the machines; the administrator who spent many sleepless nights with me has decided to take a break.

The research organizations and foundations that helped set up the cluster have not given thought to the problem. They are only concerned with getting useful things out of it still. Once you build an empire, you usually do not want to relocate.

The environmentalists roll their eyes at me. Large scale computing infrastructure in the electronic waste recycling? You serious? Well go find a place to host them.

Yes, true.

Meanwhile, does any one care about what would happen to these large scale repos once no one rides the bandwagon anymore?

I feel like the Lorax in Dr. Seuss's book(http://www.seussville.com/books/book_detail.php?isbn=9780385372022).

The Onceler and the Truffula trees had their way. The Onceler is unrepentant. He  defiantly tells the Lorax that he will keep on "biggering" his business, but at that moment one of his machines fells the very last of the Truffula trees. Without raw materials, the factory shuts down. The Lorax says nothing. Just one sad backward glance and disappears behind the smoggy clouds. Where he last stood is a small monument engraved with a single word: "UNLESS". The Onceler ponders the message for years, in solitude. 

Friday, July 22, 2016

Stir Fried Wild Figs (Dumur) - A Delicacy of the East.

The Indian fig tree, Gular (Ficus racemosa), is unusual in that its figs grow on or close to the tree trunk [Wikipedia]. The fruit is almost never sold commercially in large quantities because of the threat of fig wasps. This however, is quite easily available in local vegetable markets in West Bengal (particularly in Kolkata) and goes by the name -- 'Dumur'.



 The easiest recipe one can try with it is 'Stir Fried Dumur' - Cut the fruit in half or quarters(you may have a white secretion) and pressure cook (two whistles should do the trick) with turmeric and salt.

In a wok, heat oil. Allow 1/2 tea spoon of whole mustard seeds to splutter and add in chopped onions, ginger garlic paste, curry leaves. When onions brown and oil separates, add in the dumur, cumin powder, chili powder and salt to taste. Stir fir for a few minutes and remove from the flame. Serve as a side dish with dal and rice.

A more popular variant is to do the above but add in diced potatoes. This helps increase the quantity, if you want to prepare and keep for a few servings!