Showing posts with label hadoop. Show all posts
Showing posts with label hadoop. Show all posts

Monday, December 17, 2012

Hadoop on Windows Azure


Taken from: http://msdn.microsoft.com/en-us/magazine/jj190805.aspx

Hadoop on Windows Azure

Lynn Langit

There’s been a lot of buzz about Hadoop lately, and interest in using it to process extremely large data sets seems to grow day by day. With that in mind, I’m going to show you how to set up a Hadoop cluster on Windows Azure. This article assumes basic familiarity with Hadoop technologies. If you’re new to Hadoop, see “What Is Hadoop?” As of this writing, Hadoop on Windows Azure is in private beta. To get an invitation, visit hadooponazure.com. The beta is compatible with Apache Hadoop (snapshot 0.20.203+).

What Is Hadoop?

Hadoop is an open source library designed to batch-process massive data sets in parallel. It’s based on the Hadoop distributed file system (HDFS), and consists of utilities and libraries for working with data stored in clusters. These batch processes run using a number of different technologies, such as Map/Reduce jobs, and may be written in Java or other higher-level languages, such as Pig. There are also languages that can be used to query data stored in a Hadoop cluster. The most common language for query is HQL via Hive.  For more information, visit hadoop.apache.org.

Setting Up Your Cluster

Once you’re invited to participate in the beta, you can set up your Hadoop cluster. Go to hadooponazure.com and log in with your authorized Windows Live ID. Next, fill out the dialog boxes on the Web portal using the following values:
  1. Cluster (DNS) name: Enter name in the form “<your unique string>.cloudapp.net”.
  2. Cluster size: Choose the number of nodes, from 4 to 32, and their associated storage allocations, from 2TB to 16TB per cluster.
  3. Administrator username and password: Enter a username and password; password complexity restrictions are listed on the page. Once this is set you can connect via remote desktop or via Excel.
  4. Configuration information for a SQL Azure instance: This is an option for storing the Hive Metastore. If it’s selected, you’ll need to supply the URL to your SQL Azure server instance, as well as the name of the target database and login credentials. The login you specify must have the following permissions on the target database: ddl_ddladmin, ddl_datawriter, ddl_datareader.
After you’ve filled in this information, click Request cluster. You’ll see a series of status updates in the Web portal as your cluster (called Isotope in the beta) is being allocated, created and started. For each cluster you allocate, you’ll see many worker nodes and one head node, which is also known as the NameNode.
After some period of time (five to 30 minutes in my experience), the portal will update to show that your cluster is allocated and ready for use. You can then simply explore the Metro-style interface (by clicking the large buttons) to see what types of data processing and management tasks you can perform (see Figure 1). In addition to using the Web portal to interact with your cluster, you might want to open the available ports (closed by default) for FTP or ODBC Server access. I’ll discuss some alternative methods of connecting in a bit.
The Windows Azure Hadoop Portal
Figure 1 The Windows Azure Hadoop Portal
In the Your Cluster section of the portal, you can perform basic administrative tasks such as configuring access to your cluster, importing data and managing the cluster via the interactive console. The interactive console supports JavaScript or Hive. As Figure 1 shows, you can also access the Your Tasks section. Here you can run a MapReduce job (via a .jar file) and see the status of any MapReduce jobs that are running as well as those that have recently completed.
The portal buttons display information about three recently completed MapReduce jobs: C# Streaming Example, Word Count Example and 10GB Terasort Example. Each button shows the status of both the Map and the Reduce portion of each job. There are several other options for viewing the status of running (or completed) MapReduce jobs directly from the portal and via other means of connecting to your cluster, such as Remote Desktop Protocol (RDP).

Connecting to Your Data

You can make data available to your Hadoop on Windows Azure cluster in a number of ways, including directly uploading to your cluster and accessing data stored in other locations.
Although FTP allows uploading theoretically any size data files, best practice is to upload files that are in a lower gigabyte-size range. If you want to run batch jobs on data stored outside of Hadoop, you’ll need to perform a couple of configuration steps first. To set up outside connections, click on the Manage Cluster button on the main portal and then configure the storage locations you wish to use, such as a Windows Azure Blob storage location, a Windows Azure Data Market query result or an Amazon Web Services (AWS) S3 storage location:
  1. To configure a connection to an AWS S3 bucket, enter your security keys (both public and private) so you can access data stored on S3 in your Hadoop cluster.
  2. To work with data from the Windows Azure Data Market, fill in values for username (WLID), passkey (for the data source you wish to query and import), source (extract) query and (destination) Hive table name. Be sure to remove the parameter for default query limit (100 rows) from the query generated by the tools in the Data Market before you enter the query into the text box on your cluster.
  3. To access data stored in Windows Azure Blob storage, you’ll need to enter the storage account name (URL) to the Blob storage locations and your passkey (private key) value.

Running a MapReduce Job

After setting up and verifying your Hadoop cluster, and making your data available, you’ll probably want to start crunching this data by running one or more MapReduce jobs. The question is, how best to start? If you’re new to Hadoop, there are some samples you can run to get a feel of what’s possible. You can view and run any of these by clicking the Samples button on the Web portal.
If you’re experienced with Hadoop techniques and want to run your own MapReduce job, there are several methods. The method you select will depend on your familiarity with the Hadoop tools (such as the Hadoop command prompt) and your preferred language. You can use Java, Pig, JavaScript or C# to write an executable MapReduce job for Hadoop on Windows Azure.
I’ll use the Word Count sample to demonstrate how to run a MapReduce job from the portal using a .jar file. As you might expect, this job counts words for some input—in this example a large text file (the contents of an entire published book)—and outputs the result. Click Samples, then WordCount to open the job configuration page on the portal, as shown in Figure 2.
Setting Up the WordCount Sample
Figure 2 Setting Up the WordCount Sample
You’ll see two configurable parameters for this job, one for the function (word count) and the other for the source data (the text file). The source data (Parameter 1) includes not only the name of the input file, but also the path to its location. This path to the source data file can be text, or it can be “local,” which means that the file is stored on this Hadoop Windows Azure cluster. Alternatively, the source data can be retrieved from AWS S3 (via the S3n:// or the S3:// protocol), from the Windows Azure Blob storage (via the ASV:// protocol) or from the Windows Azure Data Market (by first importing the desired data using a query), or be retrieved directly from the HDFS store. After you enter the path to a remote location, you can click on the verification icon (a triangle) and you should get an OK message if you can connect using the string provided.
After you configure the parameters, click Execute job. You’ll find a number of ways to monitor both job status as the job is executing and job results after the job completes. For example, on the main portal page, the Your Tasks section displays a button with the status of the most recent jobs during execution and after completion. A new button is added for each job, showing the job name, percentage complete for both the Map and the Reduce portions during execution, and the status (OK, failed and so forth) after job completion.
The Job History page, which you can get to from the Manage Your Account section of the main page, provides more detail about the job, including the text (script) used to run the job and the status, with date and time information. You can click the link for each job to get even more information about job execution.
If you decide to run a sample, be sure to read the detailed instructions for that particular sample. Some samples can be run from the Web portal (Your Tasks | Create Job); others require an RDP connection to your cluster.

Using JavaScript to Run Jobs

Click on the Interactive Console button to open the JavaScript console. Here you can run MapReduce jobs by executing .jar files (Java) by running a Pig command from the prompt, or by writing and executing MapReduce jobs directly in JavaScript.
You can also directly upload source data from the js> prompt using the fs.put command. This command opens a dialog box where you can choose a file to upload to your cluster. IIS limits the size of the file you can upload via the JavaScript console to 4GB.
You can also use source data from other remote stores (such as Windows Azure Blobs) or from other cloud vendors. To work with source data from AWS S3, you use a request in the format s3n://<bucket name>/<folder name>.
Using the JavaScript console, you can verify connectivity to your AWS S3 bucket by using the #ls command with the bucket address, like so:
js> #ls s3n://HadoopAzureTest/Books
Found 2 items
-rwxrwxrwx   1          0 2012-03-30 00:20 /Books
-rwxrwxrwx   1    1395667 2012-03-30 00:22 /Books/davinci.txt

When you do, you should get a list of the contents (both folders and files) of your bucket as in this example.
If you’d like to review the contents of the source file before you run your job, you can do so from the console with the #cat command:
js> #cat s3n://HadoopAzureTest/Books/davinci.txt

After you verify that you can connect to your source data, you’ll want to run your MapReduce job. The following is the JavaScript syntax for the Word Count sample MapReduce job (using a .jar file):
  1. var map = function (key, value, context) {
  2.   var words = value.split(/[^a-zA-Z]/);
  3.   for (var i = 0; i < words.length; i++) {
  4.     if (words[i] !== "") {
  5.       context.write(words[i].toLowerCase(), 1);
  6.     }
  7.   }
  8. };
  9. var reduce = function (key, values, context) {
  10.   var sum = 0;
  11.   while (values.hasNext()) {
  12.     sum += parseInt(values.next());
  13.   }
  14.   context.write(key, sum);
  15. };
In the map portion, the script splits the source text into individual words; in the reduce portion, identical words are grouped and then counted. Finally, an output (summary) file with the top words by count (and the count of those words) is produced. To run this WordCount job directly from the interactive JavaScript console, start with the pig keyword to indicate you want to run a Pig job. Next, call the from method, which is where you pass in the location of the source data. In this case, I’ll perform the operation on data stored remotely—in AWS S3.
Now you call the mapReduce method on the Pig job, passing in the name of the file with the JavaScript code for this job, includ­ing the required parameters. The parameters for this job are the method of breaking the text—on each word—and the value and datatype of the reduce aggregation. In this case, the latter is a count (sum) of data type long.
You then specify the output order using the orderBy method, again passing in the parameters; here the count of each group of words will be output in descending order. In the next step, the take method specifies how many aggregated values should be returned—in this case the 10 most commonly occurring words. Finally, you call the to method, passing in the name of the output file you want to generate. Here’s the complete syntax to run this job:
pig.from("s3n://HadoopAzureTest/Books").mapReduce("WordCount.js", "word, count:long").orderBy("count DESC").take(10).to("DaVinciTop10Words.txt")

As the job is running, you’ll see status updates in the browser—the percent complete of first the map and then the reduce job. You can also click a link to open another browser window, where you’ll see more detailed logging about the job progress. Within a couple of minutes, you should see a message indicating the job completed successfully. To further validate the job output, you can then run a series of commands in the JavaScript console.
The first command, fs.read, displays the output file, showing the top 10 words and the total count of each in descending order. The next command, parse, shows the same information and will populate the data variable with the list. The last command, graph.bar, displays a bar graph of the results. Here’s what these commands look like:
js> file = fs.read("DaVinciTop10Words.txt")
js> data = parse(file.data, "word, count:long")
js> graph.bar(data)

An interesting aspect of using JavaScript to execute MapReduce jobs is the terseness of the JavaScript code in comparison to the Java. The MapReduce WordCount sample Java job contains more than 50 lines of code, but the JavaScript sample contains only 10 lines. The functionality of both jobs is similar.

Using C# with Hadoop Streaming

Another way you can run MapReduce jobs in Hadoop on Windows Azure is via C# Streaming. You’ll find an example showing how to do this on the portal. As with the previous example, to try out this sample, you need to upload the needed files (davinci.txt, cat.exe and wc.exe) to a storage location such as HDFS, ASV or S3. You also need to get the IP address of your Hadoop HEADNODE. To get the value using the interactive console, run this command:
js>#cat apps/dist/conf/core-site.xml

Fill in the values on the job runner page; your final command will look something like this:
Hadoop jar hadoop-examples-0.20.203.1-SNAPSHOT.jar
-files "hdfs:///example/apps/wc.exe,hdfs:///example/apps/cat.exe"
-input "/example/data/davinci.txt"
-output "/example/data/StreamingOutput/wc.txt"
-mapper "cat.exe"
-reducer "wc.exe"

In the sample, the mapper and the reducer are executable files that read the input from stdin, line by line, and emit the output to stdout. These files produce a Map/Reduce job, which is submitted to the cluster for execution. Both the mapper file, cat.exe, and the reducer file, wc.exe, are shown in Figure 3.
The Mapper and Reducer Files
Figure 3 The Mapper and Reducer Files
Here’s how the job works. First the mapper file launches as a process on mapper task initialization. If there are multiple mappers, each will launch as a separate process on initialization. In this case, there’s only a single mapper file—cat.exe. On exe­cution, the mapper task converts the input into lines and feeds those lines to the stdin portion of the MapReduce job. Next, the mapper gathers the line outputs from stdout and converts each line into a key/value pair. The default behavior (which can be changed) is that the key is created from the line prefix up to the first tab character and the value is created from the remainder of the line. If there’s no tab in the line, the whole line becomes the key and the value will be null.
After the mapper tasks are complete, each reducer file launches as a separate process on reducer task initialization. On execution, the reducer converts key/value input pairs into lines and feeds those lines to the stdin process. Next, the reducer collects the line-­oriented outputs from the stdout process and converts each line into a key/value pair, which is collected as the output of the reducer.

Using HiveQL to Query a Hive Table

Using the interactive Web console, you can execute a Hive query against Hive tables you’ve defined in your Hadoop cluster. To learn more about Hive, see hive.apache.org.
To use Hive, you first create (and load) a Hive table. Using our WordCount MapReduce sample output file (DavinciTop10­Words.txt), you can execute the following command to create and then verify your new Hive table:
hive> LOAD DATA INPATH
 'hdfs://lynnlangit.cloudapp.net:9000/user/lynnlangit/DaVinciTop10Words.txt'
 OVERWRITE INTO TABLE wordcounttable;
 hive> show tables;
 hive> describe wordcounttable:
 hive> select * from wordcounttable;

Hive syntax is similar to SQL syntax, and HiveQL provides similar query functionality. Keep in mind that all data is case-sensitive by default in Hadoop.

Other Ways to Connect to Your Cluster

Using RDP In addition to working with your cluster via the Web portal, you can also establish a remote desktop connection to the cluster’s NameNode server. To connect via RDP, you click the Remote Desktop button on the portal, then click on the downloaded RDP connection file and, when prompted, enter your administrator username and password. If prompted, open firewall ports on your client machine. After the connection is established, you can work directly with your cluster’s NameNode using the Windows Explorer shell or other tools that are included with the Hadoop installation, much as you would with the default Hadoop experience.
My NameNode server is running Windows Server 2008 R2 Enterprise SP1 on a server with two processors and 14GB of RAM, with Apache Hadoop release 0.20.203.1 snapshot installed. Note that the cluster resources consist of the name node and the associated worker nodes, so the total number of processors for my sample cluster is eight.
The installation includes standard Hadoop management tools, such as the Hadoop Command Shell or command-line interface (CLI), the Hadoop MapReduce job tracker (found at http://[namenode]:50030) and the Hadoop NameNode HDFS (found at http://[namenode]:50070). Using the Hadoop Command Shell you can run MapReduce jobs or other administrative tasks (such as managing your DFS cluster state) via your RDP session.
At this time, you can connect via RDP using only a Windows client machine. Currently, the RDP connection uses a cookie to enable port forwarding. The Remote Desktop Connection for Mac client doesn’t have the ability to use that cookie, so it can’t connect to the virtual machine.
Using the Sqoop Connector Microsoft shipped several connectors for Hadoop to SQL Server in late 2011 (for SQL Server 2008 R2 or later or for SQL Server Parallel Data Warehouse). The Sqoop-based SQL Server connector is designed to let you import or export data between Hadoop on Linux and SQL Server. You can download the connector from bit.ly/JgFmm3. This connector requires that the JDBC driver for SQL Server be installed on the same node as Sqoop. Download the driver at bit.ly/LAIU4F.
You’ll find an example showing how to use Sqoop to import or export data between SQL Azure and HDFS in the samples section of the portal.
Using FTP To use FTP, you first have to open a port, which you can do by clicking the Configure Ports button on the portal and then dragging the slider to open the default port for FTPS (port 2226). To communicate with the FTP server, you’ll need an MD5 hash of the password for your account. Connect via RDP, open the users.conf file, copy the MD5 hash of the password for the account that will be used to transfer files over FTPS and then use this value to connect. Note that the MD5 hash of the password uses a self-signed certificate on the Hadoop server that might not be fully trusted.
You can also open a port for ODBC connections (such as to Excel) in this section of the portal. The default port number for the ODBC Server connections is 10000. For more complex port configurations, though, use an RDP connection to your cluster.
Using the ODBC Driver for Hadoop (to Connect to Excel and PowerPivot) You can download an ODBC driver for Hadoop from the portal Downloads page. This driver, which includes an add-in for Excel, can connect from Hadoop to either Excel or PowerPivot. Figure 4 shows the Hive Pane button that’s added to Excel after you install the add-in. The button exposes a Hive Query pane where you can establish a connection to either a locally hosted Hadoop server or a remote instance. After doing so, you can write and execute a Hive query (via HiveQL) against that cluster and then work with the results that are returned to Excel.

Figure 4 The Hive Query Pane in Excel
You can also connect to Hadoop data using PowerPivot for Excel. To connect to PowerPivot from Hadoop, first create an OLE DB for ODBC connection using the Hive provider. On the Hive Query pane, next connect to the Hadoop cluster using the connection you configured previously, then select the Hive tables (or write a HiveQL query) and return the selected data to PowerPivot.
Be sure to download the correct version of the ODBC driver for your machine hardware and Excel. The driver is available in both 32-bit and 64-bit editions.

Easy and Flexible—but with Some Unknowns

The Hadoop on Windows Azure beta shows several interesting strengths, including:
  • Setup is easy using the intuitive Metro-style Web portal.
  • You get flexible language choices for running MapReduce jobs and data queries. You can run MapReduce jobs using Java, C#, Pig or JavaScript, and queries can be executed using Hive (HiveQL).
  • You can use your existing skills if you’re familiar with Hadoop technologies. This implementation is compliant with Apache Hadoop snapshot 0.203+.
  • There are a variety of connectivity options, including an ODBC driver (SQL Server/Excel), RDP and other clients, as well as connectivity to other cloud data stores from Microsoft (Windows Azure Blobs, the Windows Azure Data Market) and others (Amazon Web Services S3 buckets).
There are, however, many unknowns in the version of Hadoop on Windows Azure that will be publicly released:
  • The current release is a private beta only; there is little information on a roadmap and planned release features.
  • Pricing hasn’t been announced.
  • During the beta, there’s a limit to the size of files you can upload, and Microsoft included a disclaimer that “the beta is for testing features, not for testing production-level data loads.” So it’s unclear what the release-version performance will be like.
To see video demos (screencasts) of the beta functionality of Hadoop on Windows Azure, see my BigData Playlist on YouTube at bit.ly/LyX7Sj.

Lynn Langit (LynnLangit.com) runs her own technical training and consulting company. She designs and builds data solutions that include both RDBMS and NoSQL systems. She recently returned to private practice after working as a developer evangelist for Microsoft for four years. She is the author of three books on SQL Server Business Intelligence, most recently “Smart Business Intelligence Solutions with SQL Server 2008” (Microsoft Press, 2009). She is also the cofounder of the non-profit TKP (TeachingKidsProgramming.org).
Thanks to the following technical expert for reviewing this article: Denny Lee

Hadoop on Azure: An Introduction



Taken from: http://architects.dzone.com/articles/hadoop-azure-introduction




Hadoop on Azure: An Introduction

11.22.2012
| 4227 views |
  5xeb23px
I am in complete awe on how this technology is resonating with today’s developers. If I invite developers for an evening event, Big Data is always a sellout. This particular post is about getting everyone up to speed about what Hadoop is at a high level. Big data is a technology that manages voluminous amount of unstructured and semi-structured data. Due to its size and semi-structured nature, it is inappropriate for relational databases for analysis. Big data is generally in the petabytes and exabytes of data.
  1. However, it is not just about the total size of data (volume)
  2. It is also about the velocity (how rapidly is the data arriving)
  3. What is the structure? Does it have variations?
Sources of Big Data
Science Scientists are regularly challenged by large data sets in many areas, including meteorology, genomics, connectomics, complex physics simulations, and biological and environmental research.
Sensors Data sets grow in size in part because they are increasingly being gathered by ubiquitous information-sensing mobile devices, aerial sensory technologies (remote sensing), software logs, cameras, microphones, radio-frequency identification readers, and wireless sensor networks.
Social networks I am thinking of Facebook, LinkedIn, Yahoo, Google
Social influencers Blog comments, YELP likes, Twitter, Facebook likes, Apple's app store, Amazon, ZDNet, etc
Log files Computer and mobile device log files, web site tracking information, application logs, and sensor data. But there are also sensors from vehicles, video games, cable boxes or, soon, household appliances
Public Data Stores Microsoft Azure MarketPlace/DataMarket, The World Bank, SEC/Edgar, Wikipedia, IMDb
Data warehouse appliances Teradata, IBM Netezza, EMC Greenplum, which includes internal, transactional data that is already prepared for analysis
Network and in-stream monitoring technologies Packets in TCP/IP, email, etc
Legacy documents Archives of statements, insurance forms, medical record and customer correspondence

Two problems to solve
Storage Problem How do I store a petabyte of data reliably? Afterall, a petabyte is over 333 three TB drives.
Money Problem 1 petabyte costs a lot. For just 70 TB you will pay over $100,000. (eBay ad Dell/EMC CLARiiON CX3-40 -70TB- FAST 4G 15K SAN Storage is only 70 TB for $112,000)
Two seminal papers
The Google File System It is about a scalable distributed file system for large distributed data-intensive applications. It provides fault tolerance while running on inexpensive commodity hardware, and it delivers high aggregate performance to a large number of clients http://research.google.com/archive/gfs.html
MapReduce: Simplified Data Processing on Large Clusters MapReduce is a programming model and an associated implementation for processing and generating large data sets. Users specify a map function that processes a key/value pair to generate a set of intermediate key/value pairs, and a reduce function that merges all intermediate values associated with the same intermediate key http://research.google.com/archive/mapreduce.html
Hadoop : What is it?
  1. Hadoop is an open-source software framework that supports data-intensive distributed applications. Hadoop is written in Java.
  2. I met its creator, Doug Cutting, who was working at Yahoo at the time. Hadoop is named after his son's toy elephant. I was hosting a booth at the time, and I remember Doug was curious about finding some cool stuff to bring home from the booth to give to his son. Another great idea, Doug!
  3. One of the goals of Hadoop is to run applications on large clusters of commodity hardware. The cluster is composed of a single master and multiple worker nodes.
  4. Hadoop leverages the the programming model of map/reduce. It is optimized for processing large data sets.
  5. MapReduce is typically used to do distributed computing on clusters of computer. A cluster had many “nodes,” where each node is a computer in a cluster.
  6. The goal of map reduce is to break huge data sets into smaller pieces, distribute those pieces to various slave or worker nodes in the cluster, and process the the data in parallel. Hadoop leverages a distributed file system to store the data on various nodes.
It is about two functions Hadoop comes down to two functions. As long as you can write the map() and reduce() function, your data type is supported, whether we are talking abuot (1) text files (2) xml files (3) json files (4) even graphics, sound or video files.
The core is map() and reduce()
Understanding these methods is the key to mastering Hadoop
Map Step The map step is all about dividing the problem into smaller sub-problems. A master node has the job of distributing the work to worker nodes. The worker node just does one thing and returns the work back to the master node.
Reduce Step Once the master gets the work from the worker nodes, the reduce() step takes over and combines all the work. By combining the work you can form some answer and ultimately output.

01.public class WordCount {
02. 
03.public static class Map extends MapReduceBase
04.implements Mapper<LongWritable, Text, Text, IntWritable>
05.{
06. 
07.private final static IntWritable one = new IntWritable(1);
08.private Text word = new Text();
09. 
10.public void map(LongWritable key, Text value,
11.OutputCollector<Text, IntWritable> output,
12.Reporter reporter) throws IOException
13.{
14. 
15.String line = value.toString();
16.StringTokenizer tokenizer = new StringTokenizer(line);
17.while (tokenizer.hasMoreTokens())
18.{
19.word.set(tokenizer.nextToken());
20.output.collect(word, one);
21.}
22.}
23.}
24. 
25.public static class Reduce
26.extends MapReduceBase
27.implements Reducer<Text, IntWritable, Text, IntWritable>
28.{
29. 
30.public void reduce(Text key, Iterator<IntWritable> values,
31.OutputCollector<Text, IntWritable> output,
32.Reporter reporter) throws IOException
33.{
34.int sum = 0;
35.while (values.hasNext())
36.{
37.sum += values.next().get();
38.}
39.output.collect(key, new IntWritable(sum));
40.}
41.}
42. 
43.public static void main(String[] args) throws Exception
44.{
45.JobConf conf = new JobConf(WordCount.class);
46.conf.setJobName("wordcount");
47. 
48.conf.setOutputKeyClass(Text.class);
49.conf.setOutputValueClass(IntWritable.class);
50. 
51.conf.setMapperClass(Map.class);
52.conf.setCombinerClass(Reduce.class);
53.conf.setReducerClass(Reduce.class);
54. 
55.conf.setInputFormat(TextInputFormat.class);
56.conf.setOutputFormat(TextOutputFormat.class);
57. 
58.FileInputFormat.setInputPaths(conf, new Path(args[0]));
59.FileOutputFormat.setOutputPath(conf, new Path(args[1]));
60. 
61.JobClient.runJob(conf);
62.}
63.}

The “map” in MapReduce
  1. There is a master node and many slave nodes.
  2. The master node takes the input, divides it into smaller sub-problems, and distributes the input to worker or slave nodes. worker node may do this again in turn, leading to a multi-level tree structure.
  3. The worker/slave nodes processes the data into a smaller problem, and passes the answer back to its master node.
  4. Each mapping operation is independent of the others, all maps can be performed in parallel.
The “reduce” in MapReduce
  1. The master node then collects the answers from the worker or slave nodes. It then aggregates the answers and creates the needed output, which is the answer to the problem it was originally trying to solve.
  2. Reducers can also preform the reduction phase in parallel. That is how the system can process petabytes in a matter of hours.
Their are 3 key methods
The map() function will generate a list of key/value pairs based on the data
The shuffle() phase will bring things together for the reduce() phase
The reduce() phase will take the list of key/value pairs and hand that to you to do something with.

  1. The Hello World sample for Hadoop is a word count example.
  2. Let's assume our quote is this:
    • It is time for all good men to come to the aid of their country.
  3. map() function (see the "to" part)
    finds "to" twice
    (It, 1) (is, 1) (time, 1) (for, 1) (all, 1) (good, 1) (to, 1) (men, 1) (to, 1) (come, 1) (the, 1) (aid, 1) (of, 1) (their, 1) (country, 1)
    shuffle() function
    (see the "to" part)
    creates (to, 1, 1)
    (It, 1) (is, 1) (time, 1) (for, 1) (all, 1) (good, 1) (to, 1, 1) (men, 1) (come, 1) (the, 1) (aid, 1) (of, 1) (their, 1) (country, 1)
    reduce() function
    (see the "to" part) creates (to, 2)
    (It, 1) (is, 1) (time, 1) (for, 1) (all, 1) (good, 1) (men, 1) (to, 2) (come, 1) (the, 1) (aid, 1) (of, 1) (their, 1) (country, 1)

gsnjbqtb
High-level Architecture zlkw4wsp
  1. There are two main layers to both the master node and the slave nodes – the MapReduce layer and the Distributed File System Layer. The master node is responsible for mapping the data to slave or worker nodes.
Hadoop is a platform
Hadoop Common The common utilities that support the other Hadoop modules.
Hadoop Distributed File System (HDFS) A distributed file system that provides high-throughput access to application data.
Hadoop YARN A framework for job scheduling and cluster resource management.
Hadoop MapReduce A YARN-based system for parallel processing of large data sets.
There also related modules that are commonly associated with Hadoop.
Apache Pig A platform for analyzing large data sets.It includes a high-level language for expressing data analysis programs
A key point of Pig programs is that they support substantial parallelization
Pig consists of a compiler that produces sequences of Map-Reduce programs
Pig's language layer currently consists of a textual language called Pig Latin
Hive Hive is a data warehouse system for Hadoop.

It provides a SQL-like language.

It helps with data summarization and ad-hoc queries.

I am not sure yet whether this is required with Hadoop on Azure.

If not, just a few command line tasks to do:
Install Hive tar -xzvf hive-x.y.z.tar.gz
Set the environment variable HIVE_HOME $ cd hive-x.y.z/ $ export HIVE_HOME={{pwd}}
Add $HIVE_HOME/bin to your PATH: $ export PATH=$HIVE_HOME/bin:$PATH

But from what I saw here, looks like there is an ODBC Hive Setup Module.
WehnMing Ye has this video:
http://channel9.msdn.com/Events/windowsazure/learn/Hadoop-on-Windows-Azure
I signed up I recently signed up for the Windows Azure HDInsight Service here https://www.hadooponazure.com/. Logging in cezw1k2b After logging in, you will be presented with this screen: qf2laifm Next post : Calculate PI with Hadoop
  1. We will create a job name called “Pi Example.” This very simple sample will calculate PI using a cluster of comptuers.
  2. This is not necessarily the best example of big data, it is more of a compute problem.
  3. The final command line will look like this:
    1. Hadoop jar hadoop-examples-0.20.203.1-SNAPSHOT.jar pi 16 10000000
  4. More details on this sample coming soon.
Published at DZone with permission of Bruno Terkaly, author and DZone MVB. (source)
(Note: Opinions expressed in this article and its replies are the opinions of their respective authors and not those of DZone, Inc.)

Tuesday, October 9, 2012

ZooKeeper Overview


Taken from - https://cwiki.apache.org/confluence/display/ZOOKEEPER/ProjectDescription

ZooKeeper Overview
ZooKeeper allows distributed processes to coordinate with each other through a shared hierarchical name space of data registers (we call these registers znodes), much like a file system. Unlike normal file systems ZooKeeper provides its clients with high throughput, low latency, highly available, strictly ordered access to the znodes. The performance aspects of ZooKeeper allows it to be used in large distributed systems. The reliability aspects prevent it from becoming the single point of failure in big systems. Its strict ordering allows sophisticated synchronization primitives to be implemented at the client.
The name space provided by ZooKeeper is much like that of a standard file system. A name is a sequence of path elements separated by a slash ("/"). Every znode in ZooKeeper's name space is identified by a path. And every znode has a parent whose path is a prefix of the znode with one less element; the exception to this rule is root ("/") which has no parent. Also, exactly like standard file systems, a znode cannot be deleted if it has any children.
The main differences between ZooKeeper and standard file systems are that every znode can have data associated with it (every file can also be a directory and vice-versa) and znodes are limited to the amount of data that they can have. ZooKeeper was designed to store coordination data: status information, configuration, location information, etc. This kind of meta-information is usually measured in kilobytes, if not bytes. ZooKeeper has a built-in sanity check of 1M, to prevent it from being used as a large data store, but in general it is used to store much smaller pieces of data.
The service itself is replicated over a set of machines that comprise the service. These machines maintain an in-memory image of the data tree along with a transaction logs and snapshots in a persistent store. Because the data is kept in-memory ZooKeeper is able to get very high throughput and low latency numbers. The downside to an in memory database is that the size of the database that ZooKeeper can manage is limited by memory. This limitation is further reason to keep the amount of data stored in znodes small.
The servers that make up the ZooKeeper service must all know about each other. As long as a majority of the servers are available the ZooKeeper service will be available. Clients must also know the list of servers. The clients create a handle to the ZooKeeper service using this list of servers.
Clients only connect to a single ZooKeeper server. The client maintains a TCP connection through which it sends requests, gets responses, gets watch events, and sends heart beats. If the TCP connection to the server breaks, the client will connect to a different server. When a client first connects to the ZooKeeper service, the first ZooKeeper server will setup a session for the client. If the client needs to connect to another server, this session will get reestablished with the new server.
Read requests sent by a ZooKeeper client are processed locally at the ZooKeeper server to which the client is connected. If the read request registers a watch on a znode, that watch is also tracked locally at the ZooKeeper server. Write requests are forwarded to other ZooKeeper servers and go through consensus before a response is generated. Sync requests are also forwarded to another server, but does not actually go through consensus. Thus, the throughput of read requests scales with the number of servers and the throughput of write requests decreases with the number of servers.
Order is very important to ZooKeeper. (They tend to be a bit obsessive compulsive.) All updates are totally ordered. ZooKeeper actually stamps each update with a number that reflects this order. We call this number the zxid (ZooKeeper Transaction Id). Each update will have a unique zxid. Reads (and watches) are ordered with respect to updates. Read responses will be stamped with the last zxid processed by the server that services the read.


Wednesday, September 19, 2012

hadoop - Incompatible namespaceIDs in hadoop/dfs/data


Are you seeing - org.apache.hadoop.hdfs.server.datanode.DataNode: java.io.IOException: Incompatible namespaceIDs in hadoop/dfs/data: namenode namespaceID = X; datanode namespaceID = Y


The fix (in the order)

  1. Delete VERSION files in the dfs directory (#find /hadoop/dfs -name VERSION -exec rm -rf "{}" \;)
  2. Format the namenode (#hadoop namenode -format)

Cause

A new namespaceIDs (present in VERSION file) is generated each time the HDFS is formatted. It should be same for both VERSION files of data and name node.

Caution

Backup all your data. I am not responsible for your loss of data.



Tuesday, September 18, 2012

Hadoop Default Ports Quick Reference


Taken from - http://www.cloudera.com/blog/2009/08/hadoop-default-ports-quick-reference/

Hadoop Default Ports Quick Reference

Is it 50030 or 50300 for that JobTracker UI? I can never remember!
Hadoop’s daemons expose a handful of ports over TCP. Some of these ports are used by Hadoop’s daemons to communicate amongst themselves (to schedule jobs, replicate blocks, etc.). Others ports are listening directly to users, either via an interposed Java client, which communicates via internal protocols, or via plain old HTTP.
This post summarizes the ports that Hadoop uses; it’s intended to be a quick reference guide both for users, who struggle with remembering the correct port number, and systems administrators, who need to configure firewalls accordingly.

Web UIs for the Common User

The default Hadoop ports are as follows:
DaemonDefault PortConfiguration Parameter
HDFSNamenode50070dfs.http.address
Datanodes50075dfs.datanode.http.address
Secondarynamenode50090dfs.secondary.http.address
Backup/Checkpoint node?50105dfs.backup.http.address
MRJobracker50030mapred.job.tracker.http.address
Tasktrackers50060mapred.task.tracker.http.address
? Replaces secondarynamenode in 0.21.
Hadoop daemons expose some information over HTTP. All Hadoop daemons expose the following:
/logs
Exposes, for downloading, log files in the Java system property hadoop.log.dir.
/logLevel
Allows you to dial up or down log4j logging levels. This is similar to hadoop daemonlog on the command line.
/stacks
Stack traces for all threads. Useful for debugging.
/metrics
Metrics for the server. Use /metrics?format=json to retrieve the data in a structured form. Available in 0.21.
Individual daemons expose extra daemon-specific endpoints as well. Note that these are not necessarily part of Hadoop’s public API, so they tend to change over time.
The Namenode exposes:
/
Shows information about the namenode as well as the HDFS. There’s a link from here to browse the filesystem, as well.
/dfsnodelist.jsp?whatNodes=(DEAD|LIVE)
Shows lists of nodes that are disconnected from (DEAD) or connected to (LIVE) the namenode.
/fsck
Runs the “fsck” command. Not recommended on a busy cluster.
/listPaths
Returns an XML-formatted directory listing. This is useful if you wish (for example) to poll HDFS to see if a file exists. The URL can include a path (e.g., /listPaths/user/philip) and can take optional GET arguments: /listPaths?recursive=yes will return all files on the file system;/listPaths/user/philip?filter=s.* will return all files in the home directory that start with s; and/listPaths/user/philip?exclude=.txt will return all files except text files in the home directory. Beware that filter and exclude operate on the directory listed in the URL, and they ignore the recursive flag.
/data and /fileChecksum
These forward your HTTP request to an appropriate datanode, which in turn returns the data or the checksum.
Datanodes expose the following:
/browseBlock.jsp, /browseDirectory.jsp, tail.jsp, /streamFile, /getFileChecksum
These are the endpoints that the namenode redirects to when you are browsing filesystem content. You probably wouldn’t use these directly, but this is what’s going on underneath.
/blockScannerReport
Every datanode verifies its blocks at configurable intervals. This endpoint provides a listing of that check.
The secondarynamenode exposes a simple status page with information including which namenode it’s talking to, when the last checkpoint was, how big it was, and which directories it’s using.
The jobtracker‘s UI is commonly used to look at running jobs, and, especially, to find the causes of failed jobs. The UI is best browsed starting at /jobtracker.jsp. There are over a dozen related pages providing details on tasks, history, scheduling queues, jobs, etc.
Tasktrackers have a simple page (/tasktracker.jsp), which shows running tasks. They also expose/taskLog?taskid=<id> to query logs for a specific task. They use /mapOutput to serve the output of map tasks to reducers, but this is an internal API.

Under the Covers for the Developer and the System Administrator

Internally, Hadoop mostly uses Hadoop IPC to communicate amongst servers. (Part of the goal of theApache Avro project is to replace Hadoop IPC with something that is easier to evolve and more language-agnostic; HADOOP-6170 is the relevant ticket.) Hadoop also uses HTTP (for the secondarynamenode communicating with the namenode and for the tasktrackers serving map outputs to the reducers) and a raw network socket protocol (for datanodes copying around data).
The following table presents the ports and protocols (including the relevant Java class) that Hadoop uses. This table does not include the HTTP ports mentioned above.
DaemonDefault PortConfiguration ParameterProtocolUsed for
Namenode8020fs.default.name?IPC:ClientProtocolFilesystem metadata operations.
Datanode50010dfs.datanode.addressCustom Hadoop Xceiver: DataNodeand DFSClientDFS data transfer
Datanode50020dfs.datanode.ipc.addressIPC:InterDatanodeProtocol,ClientDatanodeProtocol
ClientProtocol
Block metadata operations and recovery
Backupnode50100dfs.backup.addressSame as namenodeHDFS Metadata Operations
JobtrackerIll-defined.?mapred.job.trackerIPC:JobSubmissionProtocol,InterTrackerProtocolJob submission, task tracker heartbeats.
Tasktracker127.0.0.1:0¤mapred.task.tracker.report.addressIPC:TaskUmbilicalProtocolCommunicating with child jobs
? This is the port part of hdfs://host:8020/.
? Default is not well-defined. Common values are 8021, 9001, or 8012. See MAPREDUCE-566.
¤ Binds to an unused local port.
That’s quite a few ports! I hope this quick overview has been helpful.