Showing posts with label Hadoop. Show all posts
Showing posts with label Hadoop. Show all posts

Tuesday, December 29, 2015

How to install hadoop on Linux - Single node setup (Pseudo-Distributed Mode )

This blog explains how to install and configure hadoop in Pseudo-Distributed Mode on Linux (Oracle Enterprise Linux). And common errors encountered during setup.

Install sun jdk

download jdk 6 (latest update) from http://www.oracle.com/technetwork/java/javase/overview/index.html
 I have downloaded update 37 from below url.
 http://www.oracle.com/technetwork/java/javase/downloads/jdk6u37-downloads-1859587.html

cd /softwares/hadoop
./jdk-6u37-linux-x64.bin
jdk will be installed under same folder ( jdk1.6.0_37)


Download latest stable hadoop version 

Download hadoop-1.0.4-bin.tar.gz from apache mirror http://www.motorlogy.com/apache/hadoop/common/stable/

Extract hadoop:
 tar -xvf hadoop-1.0.4-bin.tar.gz

Now hadoop is extracted under /scratch/rajiv/softwares/hadoop/hadoop-1.0.4


set JAVA_HOME for hadoop

 vi hadoop-1.0.4/conf/hadoop-env.sh

uncomment below line and update path to jdk 1.6

# export JAVA_HOME=/usr/lib/j2sdk1.5-sun

I have updated it to  
 export JAVA_HOME=/softwares/hadoop/jdk1.6.0_37


Configure SSH

You would need ssh server(demon) and client on your box. 


$yum -y install openssh-server openssh-clients

If yum is repository is not configured, refer this post



Ensure if the SSH server daemon sshd is running.
$ /sbin/service sshd status
If the SSH server daemon sshd is not running, start this daemon using below command:

$ /sbin/service sshd start


Alternativly to determine if SSH is running, enter the following command:
$ pgrep sshd

If SSH is running, this command returns one or more process ids


Run "which ssh" to check if you have ssh client installed.

Generate ssh key (passphraseless )
$ ssh-keygen -t rsa -P ""

This creates rsa keypair under home folder with empty password.
Verify that id_rsa  id_rsa.pub files are present under .ssh folder.

Copy keypair to authorizes keys

copy id_rsa.pub to authorized_keys
hadoop-1.0.4]$ cat ~/.ssh/id_rsa.pub >> ~/.ssh/authorized_keys

To Test ssh configuration run  "ssh localhost". Below output confirms that ssh is working.
    The authenticity of host 'localhost (127.0.0.1)' can't be established.

    RSA key fingerprint is 89:20:58:b4:06:c6:c6:5a:08:1e:43:eb:cc:e5:45:49.

    Are you sure you want to continue connecting (yes/no)? yes

    Warning: Permanently added 'localhost' (RSA) to the list of known hosts.
Modify hadoop config file - core-site.xml 

1) modify core-site.xml to add two properties under configuration tag.

a) hadoop.tmp.dir - if this property is not configured, hdfs will be configured under "/tmp/hadoop-/dfs". 
b) fs.default.name - if this property is not configured, startup of secondarynamenode will fail

Sample conf/core-site.xml

2) modify hdfs-site.xml
add dfs.replication property and set value to 1


Sample hdfs-site.xml

c) modify mapred-site.xml 
add mapred.job.tracker property and set value to "localhost:9001"
Format hdfs file system
Run below command
$./hadoop-1.0.4/bin/hadoop namenode -format

start hadoop

/scratch/rajiv/softwares/hadoop/hadoop-1.0.4/bin/start-all.sh


 make sure all hadoop processes are running

run "$JAVA_HOME/bin/jps" and make sure below processes are running TaskTracker JobTracker DataNode SecondaryNameNode NameNode


stop hadoop

./bin/stop-all.sh  

stopping jobtracker localhost: stopping tasktracker stopping namenode localhost: stopping datanode localhost: stopping secondarynamenode


Common errors and troubleshooting steps

1) If ssh is not configured, following errors will be displayed while starting hadoop

$ ./start-all.sh

starting namenode, logging to /scratch/rajiv/softwares/hadoop/hadoop-1.0.4/libexec/../logs/hadoop-rajiv-namenode-myhostname.out rajiv@localhost's password: rajiv@localhost's password: localhost: Permission denied, please try again. localhost: Permission denied, please try again. rajiv@localhost's password: localhost: starting datanode, logging to /scratch/rajiv/softwares/hadoop/hadoop-1.0.4/bin/../logs/hadoop-rajiv-datanode-myhostname.out

2) If hdfs is not configured, starting hadoop will fail  with following message


 To fix this make sure " core-site.xml" has below snippet

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



 $./start-all.sh


starting namenode, logging to /scratch/rajiv/softwares/hadoop/hadoop-1.0.4/libexec/../logs/hadoop-rajiv-namenode-myhostname.out localhost: starting datanode, logging to /scratch/rajiv/softwares/hadoop/hadoop-1.0.4/bin/../logs/hadoop-rajiv-datanode-myhostname.out localhost: starting secondarynamenode, logging to /scratch/rajiv/softwares/hadoop/hadoop-1.0.4/bin/../logs/hadoop-rajiv-secondarynamenode-myhostname.out localhost: Exception in thread "main" java.lang.IllegalArgumentException: Does not contain a valid host:port authority: file:/// localhost: at org.apache.hadoop.net.NetUtils.createSocketAddr(NetUtils.java:162) localhost: at org.apache.hadoop.hdfs.server.namenode.NameNode.getAddress(NameNode.java:198) localhost: at org.apache.hadoop.hdfs.server.namenode.NameNode.getAddress(NameNode.java:228) localhost: at org.apache.hadoop.hdfs.server.namenode.NameNode.getServiceAddress(NameNode.java:222) localhost: at org.apache.hadoop.hdfs.server.namenode.SecondaryNameNode.initialize(SecondaryNameNode.java:161) localhost: at org.apache.hadoop.hdfs.server.namenode.SecondaryNameNode.(SecondaryNameNode.java:129) localhost: at org.apache.hadoop.hdfs.server.namenode.SecondaryNameNode.main(SecondaryNameNode.java:567) starting jobtracker, logging to /scratch/rajiv/softwares/hadoop/hadoop-1.0.4/libexec/../logs/hadoop-rajiv-jobtracker-myhostname.out localhost: starting tasktracker, logging to /scratch/rajiv/softwares/hadoop/hadoop-1.0.4/bin/../logs/hadoop-rajiv-tasktracker-myhostname.out

3) Cause: invalid tmp folder specified:
$ ./bin/hadoop namenode -format
To fix this check core-site.xml and ensure hadoop.tmp.dir element has correct value.


13/03/26 21:49:33 INFO namenode.NameNode: STARTUP_MSG: /************************************************************ STARTUP_MSG: Starting NameNode STARTUP_MSG: host = myhostname/myipaddress STARTUP_MSG: args = [-format] STARTUP_MSG: version = 1.0.4 STARTUP_MSG: build = https://svn.apache.org/repos/asf/hadoop/common/branches/branch-1.0 -r 1393290; compiled by 'hortonfo' on Wed Oct 3 05:13:58 UTC 2012 ************************************************************/ 13/03/26 21:49:33 INFO util.GSet: VM type = 64-bit 13/03/26 21:49:33 INFO util.GSet: 2% max memory = 17.77875 MB 13/03/26 21:49:33 INFO util.GSet: capacity = 2^21 = 2097152 entries 13/03/26 21:49:33 INFO util.GSet: recommended=2097152, actual=2097152 13/03/26 21:49:34 INFO namenode.FSNamesystem: fsOwner=rajiv 13/03/26 21:49:34 INFO namenode.FSNamesystem: supergroup=supergroup 13/03/26 21:49:34 INFO namenode.FSNamesystem: isPermissionEnabled=true 13/03/26 21:49:34 INFO namenode.FSNamesystem: dfs.block.invalidate.limit=100 13/03/26 21:49:34 INFO namenode.FSNamesystem: isAccessTokenEnabled=false accessKeyUpdateInterval=0 min(s), acc essTokenLifetime=0 min(s) 13/03/26 21:49:34 INFO namenode.NameNode: Caching file names occuring more than 10 times 13/03/26 21:49:34 ERROR namenode.NameNode: java.io.IOException: Cannot create directory /u01/hadoop/tmp/dfs/na me/current at org.apache.hadoop.hdfs.server.common.Storage$StorageDirectory.clearDirectory(Storage.java:297) at org.apache.hadoop.hdfs.server.namenode.FSImage.format(FSImage.java:1320) at org.apache.hadoop.hdfs.server.namenode.FSImage.format(FSImage.java:1339) at org.apache.hadoop.hdfs.server.namenode.NameNode.format(NameNode.java:1164) at org.apache.hadoop.hdfs.server.namenode.NameNode.createNameNode(NameNode.java:1271) at org.apache.hadoop.hdfs.server.namenode.NameNode.main(NameNode.java:1288)

4) while running the basic example,

$ ./bin/hadoop jar hadoop-examples-*.jar grep input output 'dfs[a-z.]+'

13/05/07 23:04:22 INFO util.NativeCodeLoader: Loaded the native-hadoop library 13/05/07 23:04:22 WARN snappy.LoadSnappy: Snappy native library not loaded 13/05/07 23:04:22 INFO mapred.JobClient: Cleaning up the staging area file:/scratch/softwares/hadoop/tmp/mapred/staging/rajiv762105165/.staging/job_local_0001 13/05/07 23:04:22 ERROR security.UserGroupInformation: PriviledgedActionException as:rajiv cause:org.apache.hadoop.mapred.InvalidInputException: Input path does not exist: hdfs://localhost:54310/user/rajiv/input org.apache.hadoop.mapred.InvalidInputException: Input path does not exist: hdfs://localhost:54310/user/rajiv/input at org.apache.hadoop.mapred.FileInputFormat.listStatus(FileInputFormat.java:197) at org.apache.hadoop.mapred.FileInputFormat.getSplits(FileInputFormat.java:208) at org.apache.hadoop.mapred.JobClient.writeOldSplits(JobClient.java:989) at org.apache.hadoop.mapred.JobClient.writeSplits(JobClient.java:981) at org.apache.hadoop.mapred.JobClient.access$600(JobClient.java:174) at org.apache.hadoop.mapred.JobClient$2.run(JobClient.java:897) at org.apache.hadoop.mapred.JobClient$2.run(JobClient.java:850) at java.security.AccessController.doPrivileged(Native Method) at javax.security.auth.Subject.doAs(Subject.java:396) at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1121) at org.apache.hadoop.mapred.JobClient.submitJobInternal(JobClient.java:850) at org.apache.hadoop.mapred.JobClient.submitJob(JobClient.java:824) at org.apache.hadoop.mapred.JobClient.runJob(JobClient.java:1261) at org.apache.hadoop.examples.Grep.run(Grep.java:69) at org.apache.hadoop.util.ToolRunner.run(ToolRunner.java:65) at org.apache.hadoop.examples.Grep.main(Grep.java:93) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:39) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:25) at java.lang.reflect.Method.invoke(Method.java:597) at org.apache.hadoop.util.ProgramDriver$ProgramDescription.invoke(ProgramDriver.java:68) at org.apache.hadoop.util.ProgramDriver.driver(ProgramDriver.java:139) at org.apache.hadoop.examples.ExampleDriver.main(ExampleDriver.java:64) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:39) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:25) at java.lang.reflect.Method.invoke(Method.java:597) at org.apache.hadoop.util.RunJar.main(RunJar.java:156)

Cause:

the input folder is not present in hdfs

org.apache.hadoop.mapred.InvalidInputException: Input path does not exist: hdfs://localhost:54310/user/rajiv/input


Solution:


./bin/hadoop dfs -ls hdfs:/user/rajiv
if input folder is not present then copy it from local file system

./bin/hadoop dfs -copyFromLocal ./input hdfs:/user/rajiv/input

above command assume that input folder is present under current dir


5) if hadoop is not started, while running example below error is displayed

bin/hadoop jar hadoop-examples-*.jar grep input output 'dfs[a-z.]+' 13/05/07 23:09:05 INFO ipc.Client: Retrying connect to server: localhost/127.0.0.1:54310. Already tried 0 time(s). 13/05/07 23:09:06 INFO ipc.Client: Retrying connect to server: localhost/127.0.0.1:54310. Already tried 1 time(s).



http://hadoop.apache.org/docs/r1.0.4/mapred_tutorial.html#Example%3A+WordCount+v1.0



Thursday, August 22, 2013

How to build hadoop source code using ant & ivy (hadoop release-1.0.4)


Download hadoop source code

Use svn anonymous checkout to get hadoop source code

create folder to checkout source code, for example:
mkdir -p /scratch/rajiv/softwares/hadoop/source


Now checkout from svn repository.

$svn checkout http://svn.apache.org/repos/asf/hadoop/common/tags/release-1.0.4/ release-1.0.4
svn: OPTIONS of 'http://svn.apache.org/repos/asf/hadoop/common/tags/release-1.0.4': could not connect to server (http://svn.apache.org)
Got above error initially as I am behind proxy server.

Edit ~/.subversion/servers file and un-comment below lines and provide your proxy server hostname(preferably fully qualified hostname) and port. If you don't know the proxy serve name, just run "wget google.com" from terminal, it will print proxy server hostname.

Make sure you edit these properties under "[global]" section. Same properties are present under "group".
But modifying them will not have any effect. You will get same exception as above.
#http-proxy-host=my.proxy.server.name.here#http-proxy-port=80#http-compression = no

Now run svn checkout again
$svn checkout http://svn.apache.org/repos/asf/hadoop/common/tags/release-1.0.4/ release-1.0.4

Now the hadoop source code is checkout out under current folder.
$ls hadoop-common-1.0.4
 


Build hadoop source code
Hadoop 1.0.4 source doesn't have maven project defined.
There is no pom.xml present and trying to build using maven will fail with below exception.
[ERROR] The goal you specified requires a project to execute but there is no POM in this directory (/scratch/rajiv/softwares/hadoop/source/hadoop-common-1.0.4). Please verify you invoked Maven from the correct directory. -> [Help 1]
Maven project was available till hadoop release-0.23.7.
Then versions from release-0.3.0  to release-0.9.2 ant is used.
Then version from release-1.0.0 to release-1.2.0-rc1 ant and ivy are used(noticed that during build, ivy downloads maven2 artifacts).
 
Now from hadoop 2.0 on-wards ant and ivy are removed and only maven project is present.

 
Refer 
http://svn.apache.org/repos/asf/hadoop/common/tags/release-0.23.7/
 
http://svn.apache.org/repos/asf/hadoop/common/tags/release-1.0.4/

http://svn.apache.org/repos/asf/hadoop/common/tags/release-2.0.1-alpha/


Install ant and Ivy

Download and extract ivy
 wget http://apache.osuosl.org//ant/ivy/2.3.0/apache-ivy-2.3.0-bin.tar.gz

download and extract ant
 
wget http://archive.apache.org/dist/ant/binaries/apache-ant-1.8.4-bin.tar.gz
$ant jar 

Exception in thread "main" java.lang.NoClassDefFoundError: org/apache/tools/ant/launch/Launcher
Caused by: java.lang.ClassNotFoundException: org.apache.tools.ant.launch.Launcher
        at java.net.URLClassLoader$1.run(URLClassLoader.java:202)
        at java.security.AccessController.doPrivileged(Native Method)
        at java.net.URLClassLoader.findClass(URLClassLoader.java:190)
        at java.lang.ClassLoader.loadClass(ClassLoader.java:306)
        at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:301)
        at java.lang.ClassLoader.loadClass(ClassLoader.java:247)
Could not find the main class: org.apache.tools.ant.launch.Launcher.  Program will exit

set JAVA_HOME and ANT_HOME to fix above error.

$ setenv JAVA_HOME
/scratch/rajiv/softwares/hadoop/jdk
$ setenv ANT_HOME /scratch/rajiv/softwares/hadoop/ant

$/scratch/rajiv/softwares/hadoop/ant/bin/ant jar 
Buildfile: /scratch/rajiv/softwares/hadoop/source/hadoop-common-1.0.4/build.xml

clover.setup:

clover.info:
     [echo]
     [echo]      Clover not found. Code coverage reports disabled.
     [echo]

clover:

ivy-download:
      [get] Getting: http://repo2.maven.org/maven2/org/apache/ivy/ivy/2.1.0/ivy-2.1.0.jar
      [get] To: /scratch/rajiv/softwares/hadoop/source/hadoop-common-1.0.4/ivy/ivy-2.1.0.jar
      [get] Error getting http://repo2.maven.org/maven2/org/apache/ivy/ivy/2.1.0/ivy-2.1.0.jar to /scratch/rajiv/softwares/hadoop/source/hadoop-common-1.0.4/ivy/ivy-2.1.0.jar

BUILD FAILED
/scratch/rajiv/softwares/hadoop/source/hadoop-common-1.0.4/build.xml:2419: java.net.NoRouteToHostException: No route to host
        at java.net.PlainSocketImpl.socketConnect(Native Method)
        at java.net.PlainSocketImpl.doConnect(PlainSocketImpl.java:351)
        at java.net.PlainSocketImpl.connectToAddress(PlainSocketImpl.java:213)
        at java.net.PlainSocketImpl.connect(PlainSocketImpl.java:200)
        at java.net.SocksSocketImpl.connect(SocksSocketImpl.java:366)
        at java.net.Socket.connect(Socket.java:529)
        at java.net.Socket.connect(Socket.java:478)
        at sun.net.NetworkClient.doConnect(NetworkClient.java:163)
        at sun.net.www.http.HttpClient.openServer(HttpClient.java:388)
        at sun.net.www.http.HttpClient.openServer(HttpClient.java:523)
        at sun.net.www.http.HttpClient.(HttpClient.java:227)
        at sun.net.www.http.HttpClient.New(HttpClient.java:300)
        at sun.net.www.http.HttpClient.New(HttpClient.java:317)
        at sun.net.www.protocol.http.HttpURLConnection.getNewHttpClient(HttpURLConnection.java:970)
        at sun.net.www.protocol.http.HttpURLConnection.plainConnect(HttpURLConnection.java:911)
        at sun.net.www.protocol.http.HttpURLConnection.connect(HttpURLConnection.java:836)
        at org.apache.tools.ant.taskdefs.Get$GetThread.openConnection(Get.java:660)
        at org.apache.tools.ant.taskdefs.Get$GetThread.get(Get.java:579)
        at org.apache.tools.ant.taskdefs.Get$GetThread.run(Get.java:569)

Total time: 10 seconds
[rajiv@myhostname hadoop-common-1.0.4]$
ant is not able to get ivy jars. Pass http proxy host and port as -D arguments to ant to fix this.

$ /scratch/rajiv/softwares/hadoop/ant/bin/ant jar -Dhttp.proxyHost=your-proxy-server-name-here -Dhttp.proxyPort=80

This also doesn't work
set  ANT_OPTS
 setenv ANT_OPTS "-Dhttp.proxyHost=your-proxy-server-name-here -Dhttp.proxyPort=80"
$ /scratch/rajiv/softwares/hadoop/ant/bin/ant jar 
cd to hadoop source code folder, make sure build.xml file present under this folder:
 

$cd /scratch/rajiv/softwares/hadoop/source/hadoop-common-1.0.4
 
run ant:
 
$/scratch/rajiv/softwares/hadoop/ant/bin/ant jar



 
sample output (trimmed), refer this link for complete build output log
 

/scratch/rajiv/softwares/hadoop/ant/bin/ant jar Buildfile: /scratch/rajiv/softwares/hadoop/source/hadoop-common-1.0.4/build.xml

clover.setup:

clover.info:
[echo]
[echo] Clover not found. Code coverage reports disabled.
[echo]
 ........................... 
 
BUILD SUCCESSFUL
Total time: 8 minutes 34 seconds
[rajiv@myhostname hadoop-common-1.0.4]$
 






Friday, April 19, 2013

MapReduce - Part I


Mapreduce is a paradigm for distributed computing. It provides a framework for parallel computing. MapReduce enables large scale distributed data processing. MapReduce can be applied to many large scale computing problems. The name MapReduce is inspired from Map and Reduce functions in LISP programming language.But it is not an implementation of this lisp functions.

From a user's perspective, there are two basic operations in MapReduce: Map and Reduce.

Typical problems solved by MapReduce:
Read a lot of data
Map: extract something you care about from each record.
Shuffle and sort - MapReduce framework will read output from Map phase, and perform sort and grouping.

Reduce: This phase aggregate, summarize, filter or transform data and write the results to file.


Why do we need distributed computing like MapReduce

Otherwise some problems are too big to solve

Example: 

20+ billion web pages x 20KB = 400+ terabytes
- One computer can read 30-35 MB/sec from disk
   ~four months to read the web
~1,000 hard drives just to store the web
   Even more to do something with the data
 
Using MapReduce: same problem can be solved with 1000 machines, < 3 hours


The Map function reads a stream of data and parses it into intermediate (key, value) pairs. When that is complete, the Reduce function is called once for each unique key that was generated by Map and is given the key and a list of all values that were generated for that key as a parameter. The keys are presented in sorted order.

Programmers get a simple API and do not have to deal with issues of parallelization, remote execution, data distribution, load balancing, or fault tolerance. The framework makes it easy for one to use thousands of processors to process huge amounts of data (e.g., terabytes and petabytes). 

MapReduce is not a general-purpose framework for all forms of parallel programming. Rather, it is designed specifically for problems that can be broken up into the the map-reduce paradigm.

Credit:
Much of this information is from below articles:
Distributed Systems course - www.cs.rutgers.edu/~pxk/417/notes