Showing posts with label Mapreduce. Show all posts
Showing posts with label Mapreduce. Show all posts

Saturday, February 20, 2016

Compression Settings in Mapreduce

1. To read compressed input files, and the output is also compressed
   
     In the driver class, add the below:

 jobConf.setBoolean("mapred.output.compress",true)      
jobConf.setClass("mapred.output.compression.codec","GzipCodec.class","CompressionCodec.class")

    Running the program over compressed input:

    % hadoop jar MaxTempWithCompression input/tanu/input.txt.gz output

    % gunzip -c output/part-r-00000.gz
         1949    111
         1950    20

2. To compress the mapper output
 
 jonconf.setCompressMapOutput(true);
 jobConf.setMapOutputCompressorClass(GzipCodec.class);

Thursday, February 18, 2016

Designing MR Jobs to Improve Debugging

When dealing with tons of jobs each day, with tera,peta,zeta bytes of data in xmls and other files, designing our MR jobs efficiently becomes crucial to Big Data handling.

I will discuss certain ways that can help us design and develop an efficient MR job.

1. Implement Exception Handling

Most basic is to ensure implementing exception handling with all possible exceptions in the catch block, in order of most special to generic ones. Also, do not leave the catch block empty. Write proper code, that will help debugging issues. Few helpful System.err.println can be given to be viewed in Yarn logs (How to view sysouts in yarn logs) for debugging.

If you want the job to run successfully, then do not throw the exception. Rather catch it.

2. Use Counters

Consider a real world scenario where we have thousands of xml each day, and due to 1 or 2 invalid XMLs, the complete job fails.
A solution for such problem is using Counters. That can ensure how many xmls failed and on which line no, and still the job can continue successfully. The Invalid xmls can be later processed after doing the correction if needed.

Steps:
1. Create an enum in your mapper class ( outside map method)

enum InvalidXML{
     Counter;
}

2. Write you XML processing code and in catch block

catch(XMLStreamException e){
      context.setStatus("Detected Invalid Record in Xml. Check Logs);
      context.getCounter(InvalidXML.Counter).increment(1);
}

We can also print the input file name using
      String fileName = ((FileSplit) context.getInputSplit()).getPath().getName();
      context.setStatus("Detected Invalid Record in Xml " + fileName+ ". Check Logs");

3. When this jobs Runs, it will complete successfully, but in Job Logs it will display the status as below



On clicking logs, we can see


Another use of Counters is in displaying the Total no of Files Processed or Records Processed. For that, counters need to be incremented in the try block.

Counters will display files and issues for the current date (how many records failed etc), but if we need to debug issues for an older date. For that we can do the following

3. Use MultiOutputFormat in catch block to write errors to separate text files with date etc.

The MultiOutputFormat class simplifies writing output data to multiple outputs.
write(String alias, K key, V value, org.apache.hadoop.mapreduce.TaskInputOutputContext context)

In the catch block, we can write something like below
catch(SomeException e){
     MultiOutputFormat.write("XMLErrorLogs"+date, storeId,  businessDate+e.getMessage(), context);
}

The name of the text files can be ErrorLogs etc along with date etc.
The key can be your Id on which you would like to search in hive
The value can be the content you wish to search along with the errorMessage.

Once these logs are loaded into Hive, we can easily query to see, which files from which Stores are giving more issues. Or on an earlier business date, what errors we got in which files and for which stores.
Lot of relevant information can be stored and queried for valuable analysis. This can really help to debug and support huge big data application.

Hope this article will help many.

Thursday, October 15, 2015

HIVE - To get DDL of an existing table

HIVE - to get DDL of an existing table, use the below command in hive shell.

SHOW CREATE TABLE <tablename>;

Wednesday, October 14, 2015

Recovering Deleted files/dir from Hadoop Trash

Command
hadoop fs -ls hdfs://ABCDE/user/test/.Trash/Current

where
ABCDE - is the hdfs name or url (<ABCDE> or <127.0.0.1:9000>)
test - is the username


This is possible if the TRASH feature is enabled in hadoop in core-site.xml
<property>
<name>fs.trash.interval</name>
<value>120</value> 
</property>

<property>
<name>fs.trash.checkpoint.interval</name>
<value>45</value> 
</property> 

Tuesday, October 13, 2015

Getting input file name in Mapper

Get InputSplit on Context. Cast it to FileSplit and getPath on that 
Path filePath = ((FileSplit) context.getInputSplit()).getPath();
String filePathString = ((FileSplit) context.getInputSplit()).getPath().toString();
String fileName = ((FileSplit) context.getInputSplit()).getPath().getName();

Sunday, October 11, 2015

MapReduce - Viewing sysouts in yarn

To view your mapreduce system.out.printlns in YARN

After the job completes, type the below command.

$yarn logs -applicationId <application_1442077641322_19221>

It will display all sysouts in yarn logs.

Use grep along with the above command to view data for a particular text.

$yarn logs -applicationId <application_1442077641322_19221> | grep <mysysout>

Hope it helps.

Thursday, October 8, 2015

Unsupported major minor version error

Problem: On running the mapreduce jar, getting below exception
Exception in thread "main" java.lang.UnsupportedClassVersionError: 
 jobname : Unsupported major.minor version 51.0

 at java.lang.ClassLoader.defineClass1(Native Method)
 at java.lang.ClassLoader.defineClass(ClassLoader.java:643)
 at java.security.SecureClassLoader.defineClass(SecureClassLoader.java:142)
 at java.net.URLClassLoader.defineClass(URLClassLoader.java:277)
 at java.net.URLClassLoader.access$000(URLClassLoader.java:73)
 at java.net.URLClassLoader$1.run(URLClassLoader.java:212)
 at java.security.AccessController.doPrivileged(Native Method)
 at java.net.URLClassLoader.findClass(URLClassLoader.java:205)
 at java.lang.ClassLoader.loadClass(ClassLoader.java:323)
 at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:294)
 at java.lang.ClassLoader.loadClass(ClassLoader.java:268)
Reason
This error occurs if the jar is compiled with a version higher than the version where it isdeployed
Solution
1. Check current version of java by using the below command:  $ java -version
2. Get that java version installed in your machine.
3. Go to Project > Properties > Project Build Path > Edit Jre > Select Alternate Jre
4.Go to Project > Properties > Java Compiler > Select Enable Project Specific Settings > Change the Compiler Compliance Level
5. Create the jar again and deploy. It must work fine. :)