Showing posts with label map reduce. Show all posts
Showing posts with label map reduce. Show all posts

Thursday, March 28, 2013

Benchmarking a Hadoop Cluster




Benchmarks make good tests because you also get numbers that you can compare with other clusters as a sanity check on whether your new cluster is performing roughly as expected. And you can tune a cluster using benchmark results to squeeze the best performance out of it.


To get the best results, you should run benchmarks on a cluster that is not being used by others. In practice, this is just before it is put into service and users start relying on it. Once users have scheduled periodic jobs on a cluster, it is generally impossible to find a time when the cluster is not being used.


Hadoop comes with several benchmarks that you can run very easily with minimal setup cost. Benchmarks are packaged in the test JAR file, and you can get a list of them, with descriptions, by invoking the JAR file with no arguments:

% hadoop jar $HADOOP_INSTALL/hadoop-*-test.jar

Most of the benchmarks show usage instructions when invoked with no arguments.

For example:

% hadoop jar $HADOOP_INSTALL/hadoop-*-test.jar TestDFSIO 

TestFDSIO.0.0.4
Usage: TestFDSIO -read | -write | -clean [-nrFiles N] [-fileSize MB] [-resFile
resultFileName] [-bufferSize Bytes]


Benchmarking HDFS with TestDFSIO

TestDFSIO tests the I/O performance of HDFS. It does this by using a MapReduce job as a convenient way to read or write files in parallel. Each file is read or written in a separate map task, and the output of the map is used for collecting statistics related to the file just processed. The statistics are accumulated in the reduce to produce a summary.

The following command writes 10 files of 1,000 MB each:

% hadoop jar $HADOOP_INSTALL/hadoop-*-test.jar TestDFSIO -write -nrFiles 10
-fileSize 1000


At the end of the run, the results are written to the console and also recorded in a local
file (which is appended to, so you can rerun the benchmark and not lose old results):
% cat TestDFSIO_results.log
----- TestDFSIO ----- : write
Date & time: Sun Apr 12 07:14:09 EDT 2009
Number of files: 10
Total MBytes processed: 10000
Throughput mb/sec: 7.796340865378244
Average IO rate mb/sec: 7.8862199783325195
IO rate std deviation: 0.9101254683525547
Test exec time sec: 163.387
The files are written under the /benchmarks/TestDFSIO directory


To run a read benchmark, use the -read argument. Note that these files must already
exist (having been written by TestDFSIO -write):

% hadoop jar $HADOOP_INSTALL/hadoop-*-test.jar TestDFSIO -read -nrFiles 10
-fileSize 1000

Here are the results for a real run:

----- TestDFSIO ----- : read
Date & time: Sun Apr 12 07:24:28 EDT 2009
Number of files: 10
Total MBytes processed: 10000
Throughput mb/sec: 80.25553361904304
Average IO rate mb/sec: 98.6801528930664
IO rate std deviation: 36.63507598174921
Test exec time sec: 47.624
When you’ve finished benchmarking, you can delete all the generated files from HDFS
using the -clean argument:

% hadoop jar $HADOOP_INSTALL/hadoop-*-test.jar TestDFSIO -clean


Benchmarking MapReduce with Sort

Hadoop comes with a MapReduce program that does a partial sort of its input. It is very useful for benchmarking the whole MapReduce system, as the full input dataset is transferred through the shuffle. The three steps are: generate some random data, perform the sort, then validate the results.

First, we generate some random data using RandomWriter. It runs a MapReduce job with 10 maps per node, and each map generates (approximately) 1 GB of random binary data, with keys and values of various sizes.



Here’s how to invoke RandomWriter (found in the example JAR file, not the test one) to write its output to a directory called random-data:

% hadoop jar $HADOOP_INSTALL/hadoop-*-examples.jar randomwriter random-data

Next, we can run the Sort program:


% hadoop jar $HADOOP_INSTALL/hadoop-*-examples.jar sort random-data sorted-data

The overall execution time of the sort is the metric we are interested in, but it’s instructive to watch the job’s progress via the web UI (http://jobtracker-host:50030/), where you can get a feel for how long each phase of the job takes.


sanity check, we validate that the data in sorted-data is, in fact, correctly sorted:

% hadoop jar $HADOOP_INSTALL/hadoop-*-test.jar testmapredsort -sortInput random-data \

-sortOutput sorted-data

This command runs the SortValidator program, which performs a series of checks on the unsorted and sorted data to check whether the sort is accurate. It reports the outcome to the console at the end of its run:
SUCCESS! Validated the MapReduce framework's 'sort' successfully.


MRBench (invoked with mrbench) runs a small job a number of times. It acts as a good
counterpoint to sort, as it checks whether small job runs are responsive.

NNBench (invoked with nnbench) is useful for load-testing namenode hardware.

Gridmix is a suite of benchmarks designed to model a realistic cluster workload by
mimicking a variety of data-access patterns seen in practice. See the documentation
in the distribution for how to run Gridmix








Thursday, June 21, 2012

HADOOP INSTALLATION ON LINUX

HADOOP INSTALLATION ON LINUX


In pioneer days they used oxen for heavy pulling, and when one ox couldn’t budge a log,they didn’t try to grow a larger ox. We shouldn’t be trying for bigger computers, but for more systems of computers.
  
                                                                             —Grace Hopper


We live in the data age. It’s not easy to measure the total volume of data stored electronically,but an IDC estimate put the size of the “digital universe” at 0.18 zettabytes in 2006 and is forecasting a tenfold growth by 2011 to 1.8 zettabytes.1 A zettabyte is1021 bytes, or equivalently one thousand exabytes, one million petabytes, or one billion terabytes. That’s roughly the same order of magnitude as one disk drive for every person in the world.


Hadoop was created by Doug Cutting, the creator of Apache Lucene, the widely used text search library. Hadoop has its origins in Apache Nutch, an open source web search engine, itself a part of the Lucene project.




The Hadoop projects that are covered in this book are described briefly here:

Common
A set of components and interfaces for distributed filesystems and general I/O
(serialization, Java RPC, persistent data structures).


Avro
A serialization system for efficient, cross-language RPC and persistent data
storage.

MapReduce
A distributed data processing model and execution environment that runs on large clusters of commodity machines.

HDFS
A distributed filesystem that runs on large clusters of commodity machines.

Pig
A data flow language and execution environment for exploring very large datasets. Pig runs on HDFS and MapReduce clusters.

Hive
A distributed data warehouse. Hive manages data stored in HDFS and provides a query language based on SQL (and which is translated by the runtime engine to MapReduce jobs) for querying the data.

HBase
A distributed, column-oriented database. HBase uses HDFS for its underlying
storage, and supports both batch-style computations using MapReduce and point queries (random reads).

ZooKeeper
A distributed, highly available coordination service. ZooKeeper provides primitives such as distributed locks that can be used for building distributed applications.


Sqoop
A tool for efficient bulk transfer of data between structured data stores (such as relational databases) and HDFS.

Oozie
A service for running and scheduling workflows of Hadoop jobs (including Map-
Reduce, Pig, Hive, and Sqoop jobs).

Download a stable release from one of the apache download mirror (http://www.apache.org/dyn/closer.cgi/hadoop/common/)
, which is packaged as a gzipped tar file.

Hadoop 2.0.0 is the latest version (hadoop-2.0.0-alpha.tar.gz) download from (http://apache.mirrors.lucidnetworks.net/hadoop/common/hadoop-2.0.0-alpha/

Unpack this file 

% tar xzf hadoop-2.0.0-alpha.tar.gz

set JAVA_HOME and HADOOP_INSTALL variables , as hadoop is written in java it requires java installation location.


% export HADOOP_INSTALL = /usr/rjuluri/HADOOP/hadoop-2.0.0-alpha

% export JAVA_HOME = /usr/rjuluri/middleware/Jdev11.1.3/jdk160_18

% export PATH=$PATH:$HADOOP_INSTALL/bin:$HADOOP_INSTALL/sbin

To verify the installation, run the following command

% hadoop version

Hadoop 2.0.0-alpha
Subversion http://svn.apache.org/repos/asf/hadoop/common/branches/branch-2.0.0-alpha/hadoop-common-project/hadoop-common -r 1338348
Compiled by hortonmu on Wed May 16 01:28:50 UTC 2012
From source with checksum 954e3f6c91d058b06b1e81a02813303f

Hadoop can be run in one of three modes:

Standalone (or local) mode
There are no daemons running and everything runs in a single JVM. Standalone
mode is suitable for running MapReduce programs during development, since it
is easy to test and debug them.

Common     fs.default.name                              file:/// (default)
HDFS          dfs.replication                                N/A

YARN     yarn.resourcemanager.address              N/A

In standalone mode, there is no further action to take, since the default properties are set for standalone mode and there are no daemons to run.

Pseudodistributed mode
The Hadoop daemons run on the local machine, thus simulating a cluster on a
small scale.

Common     fs.default.name                            hdfs://localhost/
HDFS          dfs.replication                              1

YARN     yarn.resourcemanager.address             localhost:8032

Modify config files under etc/hadoop directory of hadoop installation

Common:
fs.default.name
hdfs://localhost/

HDFS:
dfs.replication
1

MAP-REDUCE:
yarn.resourcemanager.address
localhost:8032
yarn.nodemanager.aux-services
mapreduce.shuffle







fs.default.name
hdfs://localhost/






dfs.replication
1






mapred.job.tracker
localhost:8021




If you are running YARN, use the yarn-site.xml file:





yarn.resourcemanager.address
localhost:8032


yarn.nodemanager.aux-services
mapreduce.shuffle






make sure that SSH is installed and a server is running
Then, to enable password-less login, generate a new SSH key with an empty passphrase:

% ssh-keygen -t rsa -P '' -f ~/.ssh/id_rsa
% cat ~/.ssh/id_rsa.pub >> ~/.ssh/authorized_keys

Test this with:

% ssh localhost

To start the HDFS and YARN daemons, type:

% start-dfs.sh
% start-yarn.sh

These commands will start the HDFS daemons, and for YARN, a resource manager and a node manager. The resource manager web UI is at http://localhost:8088/

You can stop the daemons with:

% stop-dfs.sh
% stop-yarn.sh



Fully distributed mode
The Hadoop daemons run on a cluster of machines.

Common     fs.default.name                            hdfs://namenode/
HDFS          dfs.replication                              3 (default)
YARN     yarn.resourcemanager.address    resourcemanager:8032




Popular Posts