Showing posts with label MapReduce. Show all posts
Showing posts with label MapReduce. Show all posts

Thursday, February 16, 2017

WAYS TO BULK LOAD DATA IN HBASE

Dear Friends,


Going ahead with my post, this one was asked by one of my friend about HBase, for which I am sharing my thoughts and working procedure for the Loading of Bulk Data in HBase.

HBase is an open-source, NoSQL, distributed, column-oriented data store which has been implemented from Google BigTable that runs on top of HDFS. It was developed as part of Apache’s Hadoop project and runs on top of HDFS (Hadoop Distributed File System). HBase provides all the features of Google BigTable. We can call HBase a “Data Store” than a “Data Base” as it lacks many of the features available in traditional database, such as typed columns, secondary indexes, triggers, and advanced query languages, etc.
The Data model consists of Table name, row key, column family, columns, time stamp. While creating tables in HBase, the rows will be uniquely identified with the help of row keys and time stamp. In this data model the column family are static whereas columns are dynamic. Now let us look into the HBase Architecture. Hbase is a column oriented database where one has to specify what data belongs to which column family name.. So a Hbase table comprises of this minimum thing ie; A table Name and atleast 1 Column family name.

Apache HBase is all about giving you random, real-time, read/write access to your Big Data, but how do you efficiently get that data into HBase in the first place? Intuitively, a new user will try to do that via the client APIs or by using a MapReduce job with TableOutputFormat, but those approaches are problematic, as you will learn below. Instead, the HBase bulk loading feature is much easier to use and can insert the same amount of data more quickly.

In this blog I will take you through the number of ways to achieve Bulk Loading of Data in HBase.

There are basically 3 ways to Bulk load the data in HBase:-

1. Using ImportTsv Class to load txt Data to HBase.

2. Using Hive's HCatalog & Pig command.

3. Using MapReduce API.

You can download the Sample.txt file used in this blog HERE.

(NOTE:- While driving these examples, Please be sure to have your Hadoop Daemons & Hbase Daemons are up and running.)

1. Using ImportTsv Class to load txt Data to HBase:-


A) Uploading Sample.txt file to HDFS:-


Upload the sample file into HDFS by the following command:

Command > hadoop fs -put Sample.txt /Input

B) Create Table in HBase:-


For using this method we have to first create a table in HBase with number of column family according to the data. Here I am using 2 Column family in my data.

First go to HBase shell by giving below command and create a table with column family names:

Command > hbase shell      (To enter into HBase shell)

Command > create ‘Test′,’cf1’,'cf2'                 (To create a table with column family)


 C) Using ImportTsv Class  LOAD the Sample.txt file to HBase:-


Now we are set and ready to load the file in HBase. To load the file we will be using ImportTsv class from Hadoop/HBase jar file using the below command (goto hbase folder and give command):-

Command > ./bin/hbase org.apache.hadoop.hbase.mapreduce.ImportTsv -Dimporttsv.separator=”,” -Dimporttsv.columns=HBASE_ROW_KEY,cf1,cf2 Test /Input/Sample.txt


Check the data is loaded in Hbase Table:-

Command >scan 'Test'



Here’s a explanation of the different configuration elements:

-Dimporttsv.separator= " ," specifies that the separator is a comma separated value.
-Dimporttsv.bulk.output=output is a relative path to where the HFiles will be written. Since your user on the VM is “cloudera” by default, it means the files will be in /user/cloudera/output. Skipping this option will make the job write directly to HBase. (We have not used but is useful).
-Dimporttsv.columns=HBASE_ROW_KEY,f:count is a list of all the columns contained in this file. The row key needs to be identified using the all-caps HBASE_ROW_KEY string; otherwise it won’t start the job. (I decided to use the qualifier “count” but it could be anything else.)

2. Using Hive's HCatalog & Pig command:-


In this method different jar files from PIG, HIVE and HCatalog is required which can be exported using HADOOP_CLASSPATH Command, else error: ClassNotFoundException will come with respective class details.
(For the safe-side and since my classpath command didn't worked, I copied all jar file from pig/lib, hive/lib & hive/hcatalog/lib to hadoop/lib, After which it worked fine without any error.)

A) Create a Script using HIVE SerDe & Table Properties:-


After loading the data in HDFS define the HBase schema for the data in HIVE shell. Continuing with the Sample example, create a script file called sample.ddl, which will contain the HBase schema for data used by HIVE. To do so write the below code in a file and name it as Sample.ddl:

Script Sample.ddl :- 

CREATE TABLE sample_hcat_load_table (id STRING, cf1 STRING, cf2 STRING) STORED BY 'org.apache.hadoop.hive.hbase.HBaseStorageHandler' 
WITH SERDEPROPERTIES ( 'hbase.columns.mapping' = 'd:cf1,d:cf2' ) 
TBLPROPERTIES ( 'hbase.table.name' = 'sample_load_table'); 


B) Now Create and register the HBase table in HCatalog.


To register the ddl file use HCatalog. (Hcatalog will be inside HIVE folder (/usr/local/hadoop/hive/hcatalog) export the Hcatalog home and path in ~/.bashrc file (Like you did in installing hive)). After that source ~/.bashrc file to update it by giving below command:

Command > source ~/.bashrc

Now register the ddl file using syntax :- hcat -f $HBase_Table_Name.

The following HCatalog command-line command runs the DDL script Sample.ddl:

Command > hcat -f sample.ddl



Goto HBase shell by giving below command to check whether the table is created or not:-

Command > hbase shell


C) Create the import file using PIG Script:-.


The following command/script instructs Pig to load data from Sample. and store it in sample_load_table.

Script Hbase-bulk-load.pig:-

A = LOAD '/Input/Sample.txt' USING PigStorage(',') AS (id:chararray, c1:chararray, c2:chararray);

STORE A INTO 'simple_hcat_load_table' USING org.apache.hive.hcatalog.pig.HCatStorer();



Use Pig command to populate the HBase table via HCatalog bulkload:-

Continuing with the example, execute the following command:

Command > pig -useHCatalog Hbase-bulk-load.pig

Command > pig Hbase-bulk-load.pig 

(Since in my system it failed to read the Sample.txt data from HDFS I used local storage for my ease of usage by giving command pig -x local Hbase-bulk-load.pig or pig -x local -useHCatalog Hbase-bulk-load.pig )



Goto HBase shell and give scan command to check the result:-



Below is another example for achieving the same (I have not tried it.).

A = LOAD '/hbasetest.txt' USING PigStorage(',') as (id:chararray, c1:chararray, c2:chararray);
STORE A INTO 'hbase://mydata'  USING
org.apache.pig.backend.hadoop.hbase.HBaseStorage('mycf:intdata');


3. Using MapReduce API.


HBase's Put API can be used to insert the data into HDFS, but the data has to go through the complete HBase path as explained here. So, for inserting the data in bulk into HBase using the Put API is lot slower than the bulk loading option. There are some references to bulk loading (1, 2), but either they are incomplete or a bit too complicated.

1. Extract data from source(in our case from Text File).
2. Transform data into HFiles.
3. Loading the files into HBase by telling RegionServers where to find them.

Below is the coding I used for the same for my Sample.txt data file. You can modify it according to your requirement.

NOTE:- This code doesn't create a tablein HBase so, before ruuning this code in Hadoop environment, make sure to create a table in HBase using create command with coulmn families.


HBaseBulkLoadDriver

DRIVER CLASS

package com.poc.hbase;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.conf.Configured;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.hbase.HBaseConfiguration;
import org.apache.hadoop.hbase.client.HTable;
import org.apache.hadoop.hbase.client.Put;
import org.apache.hadoop.hbase.io.ImmutableBytesWritable;
import org.apache.hadoop.hbase.mapreduce.HFileOutputFormat;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.input.TextInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.hadoop.util.Tool;
import org.apache.hadoop.util.ToolRunner;

public class HBaseBulkLoadDriver extends Configured implements Tool {
    private static final String DATA_SEPERATOR = ",";
    private static final String TABLE_NAME = "sample-data";
    private static final String COLUMN_FAMILY_1="cf1";
    private static final String COLUMN_FAMILY_2="cf2";
  
    public static void main(String[] args) {
        try {
            int response = ToolRunner.run(HBaseConfiguration.create(), new HBaseBulkLoadDriver(), args);
            if(response == 0) {
                System.out.println("Job is successfully completed...");
            } else {
                System.out.println("Job failed...");
            }
        } catch(Exception exception) {
            exception.printStackTrace();
        }
    }

    @Override
    public int run(String[] args) throws Exception {
        int result=0;
        String outputPath = args[1];
        Configuration configuration = getConf();
        configuration.set("data.seperator", DATA_SEPERATOR);
        configuration.set("hbase.table.name",TABLE_NAME);
        configuration.set("COLUMN_FAMILY_1",COLUMN_FAMILY_1);
        configuration.set("COLUMN_FAMILY_2",COLUMN_FAMILY_2);
        Job job = new Job(configuration);
        job.setJarByClass(HBaseBulkLoadDriver.class);
        job.setJobName("Bulk Loading HBase Table::"+TABLE_NAME);
        job.setInputFormatClass(TextInputFormat.class);
        job.setMapOutputKeyClass(ImmutableBytesWritable.class);
        job.setMapperClass(HBaseBulkLoadMapper.class);
        FileInputFormat.addInputPaths(job, args[0]);
        FileSystem.getLocal(getConf()).delete(new Path(outputPath), true);
        FileOutputFormat.setOutputPath(job, new Path(outputPath));
        job.setMapOutputValueClass(Put.class);
        HFileOutputFormat.configureIncrementalLoad(job, new HTable(configuration,TABLE_NAME));
        job.waitForCompletion(true);
        if (job.isSuccessful()) {
            HBaseBulkLoad.doBulkLoad(outputPath, TABLE_NAME);
        } else {
            result = -1;
        }
        return result;
    }
}

HBaseBulkLoadMapper

MAPPER CLASS

package com.poc.hbase;

import org.apache.hadoop.hbase.client.Put;
import org.apache.hadoop.hbase.util.Bytes;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.hbase.io.ImmutableBytesWritable;
import org.apache.hadoop.mapreduce.Mapper;

public class HBaseBulkLoadMapper extends Mapper<LongWritable, Text, ImmutableBytesWritable, Put> {
    private String hbaseTable;
    private String dataSeperator;
    private String columnFamily1;
    private String columnFamily2;
    private ImmutableBytesWritable hbaseTableName;

    public void setup(Context context) {
        Configuration configuration = context.getConfiguration();
        hbaseTable = configuration.get("hbase.table.name");
        dataSeperator = configuration.get("data.seperator");
        columnFamily1 = configuration.get("COLUMN_FAMILY_1");
        columnFamily2 = configuration.get("COLUMN_FAMILY_2");
        hbaseTableName = new ImmutableBytesWritable(Bytes.toBytes(hbaseTable));
    }

    public void map(LongWritable key, Text value, Context context) {
        try {
            String[] values = value.toString().split(dataSeperator);
            String rowKey = values[0];
            Put put = new Put(Bytes.toBytes(rowKey));
            put.add(Bytes.toBytes(columnFamily1), Bytes.toBytes("cf1"), Bytes.toBytes(values[1]));
            put.add(Bytes.toBytes(columnFamily2), Bytes.toBytes("cf2"), Bytes.toBytes(values[2]));
            context.write(hbaseTableName, put);
        } catch(Exception exception) {
            exception.printStackTrace();
        }
    }
}

HBaseBulkLoad

HBASE CONFIGURATION CLASS

package com.poc.hbase;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.hbase.HBaseConfiguration;
import org.apache.hadoop.hbase.client.HTable;
import org.apache.hadoop.hbase.mapreduce.LoadIncrementalHFiles;

public class HBaseBulkLoad {

public static void doBulkLoad(String pathToHFile, String tableName) {
try {
Configuration configuration = new Configuration();
configuration.set("mapreduce.child.java.opts", "-Xmx1g");
HBaseConfiguration.addHbaseResources(configuration);
LoadIncrementalHFiles loadFfiles = new LoadIncrementalHFiles(configuration);
HTable hTable = new HTable(configuration, tableName);
loadFfiles.doBulkLoad(new Path(pathToHFile), hTable);
System.out.println("Bulk Load Completed..");
} catch (Exception exception) {
exception.printStackTrace();
}
}
}

NOTE:- To create a table you can tweak and use below coding.

You have to create the table first using Java API. You can do it with the below code

//Create table and do pre-split
HTableDescriptor descriptor = new HTableDescriptor(
Bytes.toBytes(tableName)
);

descriptor.addFamily(
new HColumnDescriptor(Constants.COLUMN_FAMILY_NAME)
);

HBaseAdmin admin = new HBaseAdmin(config);

byte[] startKey = new byte[16];
Arrays.fill(startKey, (byte) 0);

byte[] endKey = new byte[16];
Arrays.fill(endKey, (byte)255);

admin.createTable(descriptor, startKey, endKey, REGIONS_COUNT);
admin.close();

Run the Jar File

Compile the above coding in eclipse with including HBase jars while compilation and export the jar file and run.

Command > hadoop jar hbase.jar com/poc/hbase/HBaseBulkLoadDriver /Input/Sample.txt /Out



Now goto HBase terminal to check data is loaded.

Hbase shell > scan 'sample-table'



That's all friends...

Now go ahead and tweak the coding to learn more about  HBase working Mechanism.



References:-



Hope you all understood the procedures... 
Please do notify me for any corrections...
Kindly leave a comment for any queries/clarification...
(Detailed Description of each phase to be added soon).
ALL D BEST...


Friday, February 10, 2017

MULTIPLE OUTPUT WITH MULTIPLE INPUT FILE NAME

Dear Friends,


I was being asked to solve how to process different files at a time and store the same under each file name. Its a real-time problem where say for example, you have log files from different places and you have to process the  same logic on all but have to store it in different file name. How to do this????

In this Blog, I will take you through how to do the same using simple multiple output method in  MapReduce program. Here I am using wordcount program logic.

Problem Statement is as below.
1. N no.of input files will be in HDFS. Each input file is having list of sentences/words.
2. Write a Mapreduce program which will give wordcount of each input file in corresponding part-r file. Where part-r filename has to be <input file name> -r-0000.

The problem statement though looks difficult yet very easy to understand and implement. (Just think simple and logically).

Solution:-

The simple logical solution is:-
1. Extract the name of each file using FileSplit method.
2. Give output of the each file after processing as the name extracted by FileSplit using multiple output method.



DOWNLOAD MY INPUT FILE FROM BELOW LINK:

https://drive.google.com/file/d/0BzYUKIo7aWL_M0s2UFRKS2xoMVE/view?usp=sharing




1. TO TAKE INPUT DATA ON HDFS


hadoop fs -mkdir /Input
hadoop fs -put Input* /Input
jar xvf mulout.jar 






2. MAP REDUCE CODES:-


DRIVER CLASS


package com.mulout.wordcount;

import java.io.IOException;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.input.TextInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.hadoop.mapreduce.lib.output.LazyOutputFormat;
import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat;
import org.apache.hadoop.util.GenericOptionsParser;

public class Multiwordcnt {

public static void main(String[] args) throws IOException, InterruptedException, ClassNotFoundException {

Configuration conf = new Configuration();
Job myJob = new Job(conf, "Multiwordcnt");
args = new GenericOptionsParser(conf, args).getRemainingArgs();
FileSystem fs = FileSystem.get(new Configuration());
fs.delete(new Path("/NewOut/"), true);

myJob.setJarByClass(Multiwordcnt.class);
myJob.setMapperClass(MyMapper.class);
myJob.setReducerClass(MyReducer.class);
myJob.setMapOutputKeyClass(Text.class);
myJob.setMapOutputValueClass(IntWritable.class);
// myJob.setNumReduceTasks(0);
myJob.setOutputKeyClass(Text.class);
myJob.setOutputValueClass(IntWritable.class);
LazyOutputFormat.setOutputFormatClass(myJob, TextOutputFormat.class);

myJob.setInputFormatClass(TextInputFormat.class);
myJob.setOutputFormatClass(TextOutputFormat.class);

FileInputFormat.addInputPath(myJob, new Path(args[0]));
FileOutputFormat.setOutputPath(myJob, new Path(args[1]));

System.exit(myJob.waitForCompletion(true) ? 0 : 1);
}

}


EXPLANATION:- In driver class LazyOutputFormat is used to store the file in -r-0000 format, without using the same we will not get output.
(Here I have used delete syntax to delete if the existing folder is there in HDFS.)

MAPPER CLASS


package com.mulout.wordcount;

import java.io.IOException;
import java.util.StringTokenizer;

import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.mapreduce.lib.input.FileSplit;

public class MyMapper extends Mapper<LongWritable, Text, Text, IntWritable> {

Text emitkey = new Text();
IntWritable emitvalue = new IntWritable(1);

public void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {

String filePathString = ((FileSplit) context.getInputSplit()).getPath().getName().toString();
String line = value.toString();
StringTokenizer tokenizer = new StringTokenizer(line);
while (tokenizer.hasMoreTokens()) {

String filepathword = filePathString + "*" + tokenizer.nextToken();
emitkey.set(filepathword);
context.write(emitkey, emitvalue);
}
}
}

EXPLANATION:- In Mapper class we took the File Input Name using FileSplit menthod and combined that with the individual word and kept as output key. Then we assinged 1 for each word as output value for futher processing in reducer

REDUCER CLASS


package com.mulout.wordcount;

import java.io.IOException;

import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.lib.output.MultipleOutputs;

public class MyReducer extends Reducer<Text, IntWritable, Text, IntWritable> {
Text emitkey = new Text();
IntWritable emitvalue = new IntWritable();
private MultipleOutputs<Text, IntWritable> multipleoutputs;

public void setup(Context context) throws IOException, InterruptedException {
multipleoutputs = new MultipleOutputs<Text, IntWritable>(context);
}

public void reduce(Text key, Iterable<IntWritable> values, Context context)
throws IOException, InterruptedException {
int sum = 0;

for (IntWritable value : values) {
sum = sum + value.get();
}
String pathandword = key.toString();
String[] splitted = pathandword.split("\\*");
String path = splitted[0];
String word = splitted[1];
emitkey.set(word);
emitvalue.set(sum);
System.out.println("word:" + word + "\t" + "sum:" + sum + "\t" + "path:  " + path);
multipleoutputs.write(emitkey, emitvalue, ("/NewOut/"+path));
}

public void cleanup(Context context) throws IOException, InterruptedException {
multipleoutputs.close();
}
}

EXPLANATION:- In reducer class we splitted the key containing Input File Name and added all 1 to get sum of number of times the word occurred and then used multiple output method with 3 parameters <Key,Value,Path> to display our result in individual File Name.
(Here I have used additional output folder "/NewOut/" for storing my results.)



3. EXECUTING THE MAP REDUCE CODE


Command > hadoop jar mulout.jar com/mulout/wordcount/Multiwordcnt /Input /Out1







That's all....

Now you can take N number of Input files and process it and store it in same File name.



Hope you all understood the procedures... 
Please do notify me for any corrections...
Kindly leave a comment for any queries/clarification...
(Detailed Description of each phase to be added soon).
ALL D BEST...


Tuesday, January 17, 2017

HADOOP POC ON EXCEL DATA WEATHER REPORT ANALYSIS

Hello Friends,


Glad to present this blog which is for analysis of Weather Report POC, which is in Excel Format. This POC  was given to me and asked by one of my friends to complete it.

Most of the time we get data in Excel Format and according to that we have to make changes in our coding. So, in this POC I have modified my previous code to accept the excel data, for the convenience of making you all understand the concept.

NOTE:- Though this POC is to read EXCEL data, I have not used the same in my coding but still it worked. (I have no idea how & why it happened. Kindly share if you know anything on the same.)
I worked out this POC on my previous POC's processed system.  So all required jar files for excel reading were already there in hadoop lib folder.
If you face any problem in reading the input file kindly use EXCEL INPUT FORMAT from my previous blog to read the data. 

UPDATE:- CORRECTION:- In this blog the Input file is not in Excel format. so it works directly without using Excel Input Format Class. (Please find the Excel Input File HERE and Compiled Coding Jar file HERE)

Problem Statement:


1. The system receives temperatures of various cities captured at regular intervals of time on each day in an input file.

2. All cities weather information for a week will be inputted to the system in a single input file.

3. System will process the input data file and generates a report with Maximum and Minimum temperatures of each day.

4. Generates a separate output report for each Month.

Ex: January-r-00000
February-r-00000
March-r-00000

5. Develop a PIG Script to filter the Map Reduce Output in the below fashion
- Provide the Unique data
- Sort the Unique data based on RETAIL_ID in DESC order

6. EXPORT the same PIG Output from HDFS to MySQL using SQOOP

7. Store the same PIG Output in a HIVE External Table.

Input File Format:- .xls (EXCEL Format)


This POC Input file and Problem statement was shared to me by Mr. Amol Wani which contains temperature statistics with time for multiple Months. Schema of record set is as shown in picture below :-





DOWNLOAD MY INPUT FILE FROM BELOW LINK:


https://drive.google.com/file/d/0BzYUKIo7aWL_WkFYdWU5QWdJLTA/view?usp=sharing

1. TO TAKE INPUT DATA ON HDFS


hadoop fs -mkdir /InputData
hadoop fs -put WeatherReport.txt /InputData
jar xvf WeatherPoc.jar 

(Please find my jar file HERE)



2.     MAP REDUCE CODES:-


WEATHER REPORT PROCESSOR 
(DRIVER CLASS)

NOTE:- If you face any problem in reading the input file kindly uncomment the following and add necessary class path & jar files.
// job.setInputFormatClass(ExcelInputFormat.class);
// job.setOutputFormatClass(TextOutputFormat.class);
// LazyOutputFormat.setOutputFormatClass(job, TextOutputFormat.class);

Please go through my previous blog on Any Excel Data reading.

package com.poc.weather;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.hadoop.mapreduce.lib.output.LazyOutputFormat;
import org.apache.hadoop.mapreduce.lib.output.MultipleOutputs;
import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat;

import com.poc.ExcelInputFormat;

public class WeatherReportProcessor {

public static String January = "January";
public static String February = "February";
public static String March = "March";
public static String April = "April";
public static String May = "May";
public static String June = "June";
public static String July = "July";
public static String August = "August";
public static String September = "September";
public static String October = "October";
public static String November = "November";
public static String December = "December";

public static void main(String[] args) throws Exception {
Configuration conf = new Configuration();
Job job = new Job(conf, "Weather Report");
job.setJarByClass(WeatherReportProcessor.class);

job.setMapperClass(WeatherMapper.class);
job.setReducerClass(WeatherReducer.class);

// job.setInputFormatClass(ExcelInputFormat.class);
// job.setOutputFormatClass(TextOutputFormat.class);
// LazyOutputFormat.setOutputFormatClass(job, TextOutputFormat.class);

job.setOutputKeyClass(Text.class);
job.setOutputValueClass(Text.class);

MultipleOutputs.addNamedOutput(job, January, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, February, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, March, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, April, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, May, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, June, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, July, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, August, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, September, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, October, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, November, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, December, TextOutputFormat.class, Text.class, Text.class);
// job.setNumReduceTasks(0);

FileInputFormat.addInputPath(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));

System.exit(job.waitForCompletion(true) ? 0 : 1);
}

}

WEATHER MAPPER 
(HAVING MAPPER LOGIC)

In Mapper, after reading input data from excel, I am removing the first two lines which doesn't contain any related data, and then splitting the entire data and taking only Date and Temperatures as my output from Mapper which be be used as input for Reducer.

package com.poc.weather;

import java.io.IOException;

import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;

public class WeatherMapper extends Mapper<LongWritable, Text, Text, Text> {
public void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
try {
if (value.toString().contains("ID") || value.toString().contains("mm"))
return;
else {
String[] str = value.toString().split(" ");
String data = "";
for (int i = 0; i < str.length; i++) {
if (str[i] != null || str[i] != " ") {
data += (str[i] + " ");

}
}
String Trim = data.trim().replaceAll("\\s+", "\t");
String[] Split = Trim.toString().split("\t");
String Date = Split[1] + Split[2] + Split[3] + Split[4] + Split[5];
String Temp = Split[9] + "\t" + Split[10];
context.write(new Text(Date), new Text(Temp));

}
} catch (Exception e) {
e.printStackTrace();
}
}
}

WEATHER REDUCER 
(HAVING REDUCER LOGIC)

In Reducer phase taking the output from Mapper, I am splitting the temperatures to get max and min temp. and comparing them with other data of different hours from a single day to get the max and min temp of that day.
After getting the max and min temp, I am checking the date for sorting them into different months.

package com.poc.weather;

import java.io.IOException;

import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.lib.output.MultipleOutputs;

public class WeatherReducer extends Reducer<Text, Text, Text, Text> {

MultipleOutputs<Text, Text> mos;

public void setup(Context context) {
mos = new MultipleOutputs<Text, Text>(context);
}

public void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException {
float f1 = 0, f2 = 50;
Text result = new Text();

while (values.iterator().hasNext()) {
String sr = values.iterator().next().toString();
String[] str1 = sr.split("\t");
float max = Float.parseFloat(str1[0]);
float min = Float.parseFloat(str1[1]);

if (max > f1) {
f1 = max;
} else if (min < f2) {
f2 = min;
}

}

result = new Text(Float.toString(f1) + "\t" + Float.toString(f2));

String fileName = "";
if (key.toString().contains("/01/")) {
fileName = WeatherReportProcessor.January;
} else if (key.toString().contains("/02/")) {
fileName = WeatherReportProcessor.February;
} else if (key.toString().contains("/03/")) {
fileName = WeatherReportProcessor.March;
} else if (key.toString().contains("/04/")) {
fileName = WeatherReportProcessor.April;
} else if (key.toString().contains("/05/")) {
fileName = WeatherReportProcessor.May;
} else if (key.toString().contains("/06/")) {
fileName = WeatherReportProcessor.June;
} else if (key.toString().contains("/07/")) {
fileName = WeatherReportProcessor.July;
} else if (key.toString().contains("/08/")) {
fileName = WeatherReportProcessor.August;
} else if (key.toString().contains("/09/")) {
fileName = WeatherReportProcessor.September;
} else if (key.toString().contains("/10/")) {
fileName = WeatherReportProcessor.October;
} else if (key.toString().contains("/11/")) {
fileName = WeatherReportProcessor.November;
} else if (key.toString().contains("/12/")) {
fileName = WeatherReportProcessor.December;
}
// String strArr[] = key.toString().split("_");
// key.set(strArr[1]);
mos.write(fileName, key, result);
}

@Override
public void cleanup(Context context) throws IOException, InterruptedException {
mos.close();
}

}



3. EXECUTING THE MAP REDUCE CODE


hadoop jar WeatherPoc.jar com.poc.weather.WeatherReportProcessor /InputData/WeatherReport.xls /WeatherOutput



We can clearly see that the input records is 8986 but the output is 365. ; It has sorted the data into number of days in a year which has been kept in different months as specified in coding.







4.     PIG SCRIPT

PigScript1.pig

A = LOAD '/WeatherReport/' USING PigStorage ('\t') AS (date:chararray, mintemp:float, maxtemp:float);

B = DISTINCT A;
DUMP B; 





PigScript2.pig

A = LOAD '/WeatherReport/' USING PigStorage ('\t') AS (date:chararray, mintemp:float, maxtemp:float);

B = DISTINCT A;
C = ORDER B BY date DESC;
STORE C INTO '/WeatherPOC'; 







5.     EXPORT the PIG Output from HDFS to MySQL using SQOOP

sqoop eval --connect jdbc:mysql://localhost/ --username root --password root --query "create database if not exists WEATHERPOC;";


sqoop eval --connect jdbc:mysql://localhost/ --username root --password root --query "use WEATHERPOC;";



sqoop eval --connect jdbc:mysql://localhost/ --username root --password root --query "grant all privileges on WEATHERPOC.* to ‘localhost’@’%’;”;

sqoop eval --connect jdbc:mysql://localhost/ --username root --password root --query "grant all privileges on WEATHERPOC.* to ‘’@’localhost’;”;


sqoop eval --connect jdbc:mysql://localhost/WEATHERPOC --username root --password root --query "create table weatherpoc(date varchar(50), mintemp float, maxtemp float);";


sqoop export --connect jdbc:mysql://localhost/WEATHERPOC --table weatherpoc --export-dir /WeatherPOC --fields-terminated-by '\t';



6.     STORE THE PIG OUTPUT IN A HIVE EXTERNAL TABLE

Goto hive shell using command:

hive

show databases;
create database WeatherPOC;
use WeatherPOC;



create external table weatherpoc(Name string, mintemp float, maxtemp float)
row format delimited
fields terminated by '\t'
stored as textfile location '/WeatherPOC';






Hope you all understood the procedures... 
Please do notify me for any corrections...
Kindly leave a comment for any queries/clarification...
(Detailed Description of each phase to be added soon).

ALL D BEST...