Showing posts with label SQOOP. Show all posts
Showing posts with label SQOOP. Show all posts

Monday, January 16, 2017

HADOOP (PROOF OF CONCEPTS) WEATHER REPORT ANALYSIS

Hello Friends,


Welcome back... This blog is for analysis of Weather Report POC which was given to me and asked by one of my friends to complete it. While searching for the same I came across a very good website which I just can't wait to share with you all.. In this POC I have modified and used both for convenience of making you all understand the concept.

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 city.

Ex: California-r-00000
Newjersy-r-00000
Newyork-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:- .txt


This POC Input file and Problem statement was shared to me by Mr. Amol Wani which contains temperature statistics with time for multiple cities.Schema of record set:-

CA_25-Jan-2014 00:12:345 15.7 01:19:345 23.1 02:34:542 12.3 03:12:187 16 04:00:093 -14 05:12:345 35.7 06:19:345 23.1 07:34:542 12.3 08:12:187 16 09:00:093 -7 10:12:345 15.7 11:19:345 23.1 12:34:542 -22.3 13:12:187 16 14:00:093 -7 15:12:345 15.7 16:19:345 23.1 19:34:542 12.3 20:12:187 16 22:00:093 -7

CA is city code, here it stands for California followed by date. After that each pair of values represent time and temperature.



DOWNLOAD MY INPUT FILE FROM BELOW LINK:



1. TO TAKE INPUT DATA ON HDFS


hadoop fs -mkdir /InputData
hadoop fs -put weather_report.txt /InputData
jar xvf WeatherReportPoc.jar 




2.     MAP REDUCE CODES:-


WEATHER REPORT PROCESSOR 
(DRIVER CLASS)

package com.poc;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.FloatWritable;
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.MultipleOutputs;
import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat;

public class WeatherReportProcessor {

public static String caOutputName = "California";
public static String nyOutputName = "Newyork";
public static String njOutputName = "Newjersy";
public static String ausOutputName = "Austin";
public static String bosOutputName = "Boston";
public static String balOutputName = "Baltimore";

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.setMapOutputKeyClass(Text.class);
job.setMapOutputValueClass(FloatWritable.class);

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

MultipleOutputs.addNamedOutput(job, caOutputName, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, nyOutputName, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, njOutputName, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, bosOutputName, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, ausOutputName, TextOutputFormat.class, Text.class, Text.class);
MultipleOutputs.addNamedOutput(job, balOutputName, 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)

package com.poc;

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

import org.apache.hadoop.io.FloatWritable;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;

public class WeatherMapper extends Mapper<Object, Text, Text, FloatWritable> {

private final static IntWritable one = new IntWritable(1);
private Text word = new Text();

public void map(Object key, Text dayReport, Context context)
throws IOException, InterruptedException {
StringTokenizer st2 = new StringTokenizer(dayReport.toString(), "\t");

int counter = 0;
String cityDateString = "";
String maxTempTime = "";
String minTempTime = "";
String curTime = "";
float curTemp = 0;
float minTemp = Float.MAX_VALUE;
float maxTemp = Float.MIN_VALUE;

while (st2.hasMoreElements()) {
if (counter == 0) {
cityDateString = st2.nextToken();
} else {
if (counter % 2 == 1) {
curTime = st2.nextToken();
} else if (counter % 2 == 0) {
curTemp = Float.parseFloat(st2.nextToken());
if (minTemp > curTemp) {
minTemp = curTemp;
minTempTime = curTime;
} else if (maxTemp < curTemp) {
maxTemp = curTemp;
maxTempTime = curTime;
}
}
}
counter++;
}

FloatWritable fValue = new FloatWritable();
Text cityDate = new Text();

fValue.set(maxTemp);
cityDate.set(cityDateString);
context.write(cityDate, fValue);

fValue.set(minTemp);
cityDate.set(cityDateString);
context.write(cityDate, fValue);
}
}


WEATHER REDUCER 
(HAVING REDUCER LOGIC)

package com.poc;

import java.io.IOException;

import org.apache.hadoop.io.FloatWritable;
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, FloatWritable, Text, Text> {

MultipleOutputs<Text, Text> mos;

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

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

for (FloatWritable value : values) {
if (counter == 0)
f1 = value.get();
else
f2 = value.get();

counter = counter + 1;
}
if (f1 > f2) {
result = new Text(Float.toString(f2) + "\t" + Float.toString(f1));
} else {
result = new Text(Float.toString(f1) + "\t" + Float.toString(f2));
}
String fileName = "";
if (key.toString().contains("CA")) {
fileName = WeatherReportProcessor.caOutputName;
} else if (key.toString().contains("NY")) {
fileName = WeatherReportProcessor.nyOutputName;
} else if (key.toString().contains("NJ")) {
fileName = WeatherReportProcessor.njOutputName;
} else if (key.toString().substring(0, 3).equals("AUS")) {
fileName = WeatherReportProcessor.ausOutputName;
} else if (key.toString().substring(0, 3).equals("BOS")) {
fileName = WeatherReportProcessor.bosOutputName;
} else if (key.toString().substring(0, 3).equals("BAL")) {
fileName = WeatherReportProcessor.balOutputName;
}

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 WeatherReportPoc.jar com.poc.WeatherReportProcessor /InputData/weather_report.txt /WeatherOutput






Explanation:- 


In map method, we are parsing each input line and maintains a counter for extracting date and each temperature & time information.For a given input line, first extract date(counter ==0) and followed by alternatively extract time(counter%2==1) since time is on odd number position like (1,3,5....) and get temperature otherwise. Compare for max & min temperature and store it accordingly. Once while loop terminates for a given input line, write maxTempTime and minTempTime with date.

In reduce method, for each reducer task, setup method is executed and create MultipleOutput object. For a given key, we have two entry (maxtempANDTime and mintempANDTime). Iterate values list , split value and get temperature & time value. Compare temperature value and create actual value sting which reducer write in appropriate file.

In main method,a instance of Job is created with Configuration object. Job is configured with mapper, reducer class and along with input and output format. MultipleOutputs information added to Job to indicate file name to be used with input format.

4.     PIG SCRIPT

PigScript1.pig

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

B = DISTINCT A;
DUMP B; 



(In the output we can clearly see that it is reading all files and as we have given DISTINCT command it is removing duplicate entries)

PigScript2.pig

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

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


NOTE:- Don't give DISTINCT command if you want to export all entries.






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 WEATHER;";




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



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

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



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


sqoop export --connect jdbc:mysql://localhost/WEATHER--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...

















Saturday, January 7, 2017

HADOOP (PROOF OF CONCEPT) RETAIL DATA BY MAHESH CHANDRA MAHARANA

INDUSTRY: RETAIL

Data Input Format :- .xls (My Input Data is in excel 2007-2003 Format)


Kindly check my blog to read any kind of Excel sheet and use the Excel Input format, record reader and excel parser given in that blog. Please find link to my blog below:


This POC Input file and Problem statement was shared to me by Mr. Sunil Pashikanti like this below was created 3000 records:-

ATTRIBUTES are like:-
1. RETAIL_ID
2. RETAIL_NAME
3. TYPE_OF_CRAWLING
4. PRODUCT_URL
5. TITTLE
6. SALE_PRICE
7. REG_PRICE
8. REBATE_PERCENTAGE
9. STOCK_INFO

Example:

12 Amazon BS http://www.amazon.com/dell/lp Amazon.com:Dell Laptop 100.00
150.00 33 InStock


DOWNLOAD MY INPUT FILE FROM BELOW LINK:

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

PROBLEM STATEMENT: -

1. take the complete Excel Input data on HDFS
  
2. Develop a Map Reduce Use Case to get the below filtered results from the HDFS Input data(Excel data)

     IF Type_Of_Crawling is -->'BS'
          -salePrice < 100.00  & RebatePercent>50  --> store "HighBuzzProducts"
          -RegPrice<150.00 & RebatePercent in 25-50 --> store "NormalProducts"
          -lengthOf(title)>100 ---> 'rare products'
          
     IF Type_Of_Crawling is -->'ODC'
          - salePrice < 150.00 --> store "OnDemandCrawlProducts"
          - StockInfo --> "InStock"  -->store "AvailableProducts"
    ELSE
          store in "OtherProducts"

  NOTE: In the mentioned file names only 5 outputs have to be generated

3. 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

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

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


NOTE:- For this POC I have used custom input format to read EXCEL files using external jar. So the corresponding jar files to be added during coding and to the lib directory of hadoop for successful execution. You can use poi-xml jar for the reading .xlsx file (2010 onwards excel format ).

Below is the steps to make it work... 

1. Download and Install ant from below link.

http://muug.ca/mirror/apache-dist//ant/binaries/apache-ant-1.9.8-bin.tar.gz


2. To install give following command in terminal:

tar -xzvf <apache ant Path>

3. Update bashrc:-

nano ~/.bashrc

Add below two lines:-

export ANT_HOME=${ant_dir}

export PATH=${ANT_HOME}/bin

Now Source bashrc by command:
source ~/.bashrc

4. Then restart the system. (Very Important for the effect to take place)

5. Download the required Jar files from below link:



Place both jar files during Eclipse compilation and only SNAPSHOT.jar in hadoop lib directory.
6. If still not working try to add CLASSPATH:
export CLASSPATH=.:$CLASSPATH:<Path to the jar file 1>:<Path to jar file 2>

Hope it will work now.

POC Processing Details



MAP REDUCE PROCESS IN DETAILS:-


1.     TO TAKE XLS INPUT DATA ON HDFS

hadoop fs -mkdir /Input
hadoop fs -put POC.xls /Input
jar xvf poc.jar 



2.     MAP REDUCE CODES:-

EXCEL INPUT DRIVER 
(DRIVER CLASS)

package com.poc;

import java.io.IOException;
import org.apache.hadoop.conf.Configuration;
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.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 org.apache.hadoop.util.GenericOptionsParser;

public class PocDriver {

static public int count = 0;

public static void main(String[] args) throws IOException, InterruptedException, ClassNotFoundException {
Configuration conf = new Configuration();

GenericOptionsParser parser = new GenericOptionsParser(conf, args);
args = parser.getRemainingArgs();

Job job = new Job(conf, "Retail_Poc");
job.setJarByClass(PocDriver.class);

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

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

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

job.setMapperClass(PocMapper.class);
job.setReducerClass(PocReducer.class);

MultipleOutputs.addNamedOutput(job, "HighBuzzProducts", TextOutputFormat.class, IntWritable.class, Text.class);
MultipleOutputs.addNamedOutput(job, "NormalProducts", TextOutputFormat.class, IntWritable.class, Text.class);
MultipleOutputs.addNamedOutput(job, "RareProducts", TextOutputFormat.class, IntWritable.class, Text.class);
MultipleOutputs.addNamedOutput(job, "OnDemandCrawlProducts", TextOutputFormat.class, IntWritable.class,
Text.class);
MultipleOutputs.addNamedOutput(job, "AvailableProducts", TextOutputFormat.class, IntWritable.class, Text.class);
MultipleOutputs.addNamedOutput(job, "OtherProducts", TextOutputFormat.class, IntWritable.class, Text.class);

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

}
}


EXCEL INPUT FORMAT 
(CUSTOM INPUT FORMAT TO READ EXCEL FILES)

package com.poc;
import java.io.IOException;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.InputSplit;
import org.apache.hadoop.mapreduce.RecordReader;
import org.apache.hadoop.mapreduce.TaskAttemptContext;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;

public class ExcelInputFormat extends FileInputFormat<LongWritable, Text> {
            @Override
            public RecordReader<LongWritable, Text> createRecordReader(InputSplit split, TaskAttemptContext context)
                                    throws IOException, InterruptedException {
                        return new ExcelRecordReader();
            }
}

EXCEL RECORD READER 
(TO READ EXCEL FILE AND SEND AS KEY, VALUE FORMAT)

package com.poc;
import java.io.IOException;
import java.io.InputStream;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FSDataInputStream;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.InputSplit;
import org.apache.hadoop.mapreduce.RecordReader;
import org.apache.hadoop.mapreduce.TaskAttemptContext;
import org.apache.hadoop.mapreduce.lib.input.FileSplit;
import com.sreejithpillai.excel.parser.ExcelParser;

public class ExcelRecordReader extends RecordReader<LongWritable, Text> {
            private LongWritable key;
            private Text value;
            private InputStream is;
            private String[] strArrayofLines;
            @Override
            public void initialize(InputSplit genericSplit, TaskAttemptContext context)
                                    throws IOException, InterruptedException {
                        FileSplit split = (FileSplit) genericSplit;
                        Configuration job = context.getConfiguration();
                        final Path file = split.getPath();
                        FileSystem fs = file.getFileSystem(job);
                        FSDataInputStream fileIn = fs.open(split.getPath());
                        is = fileIn;
                        String line = new ExcelParser().parseExcelData(is);
                        this.strArrayofLines = line.split("\n");
            }
            @Override
            public boolean nextKeyValue() throws IOException, InterruptedException {
                        if (key == null) {
                                    key = new LongWritable(0);
                                    value = new Text(strArrayofLines[0]);
                        } else {
                                    if (key.get() < (this.strArrayofLines.length - 1)) {
                                                long pos = (int) key.get();
                                                key.set(pos + 1);
                                                value.set(this.strArrayofLines[(int) (pos + 1)]);
                                                pos++;
                                    } else {
                                                return false;
                                    }
                        }
                        if (key == null || value == null) {
                                    return false;
                        } else {
                                    return true;
                        }
            }
            @Override
            public LongWritable getCurrentKey() throws IOException, InterruptedException {
                        return key;
            }
            @Override
            public Text getCurrentValue() throws IOException, InterruptedException {
                        return value;
            }
            @Override
            public float getProgress() throws IOException, InterruptedException {
                        return 0;
            }
            @Override
            public void close() throws IOException {
                        if (is != null) {
                                    is.close();
                        }
            }
}


EXCEL PARSER 
(TO PARSE EXCEL SHEET)

package com.poc;
import java.io.IOException;
import java.io.InputStream;
import java.util.Iterator;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.apache.poi.hssf.usermodel.HSSFSheet;
import org.apache.poi.hssf.usermodel.HSSFWorkbook;
import org.apache.poi.ss.usermodel.Cell;
import org.apache.poi.ss.usermodel.Row;

public class ExcelParser {
            private static final Log LOG = LogFactory.getLog(ExcelParser.class);
            private StringBuilder currentString = null;
            private long bytesRead = 0;
            public String parseExcelData(InputStream is) {
                        try {
                                    HSSFWorkbook workbook = new HSSFWorkbook(is);
                                    HSSFSheet sheet = workbook.getSheetAt(0);
                                    Iterator<Row> rowIterator = sheet.iterator();
                                    currentString = new StringBuilder();
                                    while (rowIterator.hasNext()) {
                                                Row row = rowIterator.next();
                                                Iterator<Cell> cellIterator = row.cellIterator();
                                                while (cellIterator.hasNext()) {
                                                            Cell cell = cellIterator.next();
                                                            switch (cell.getCellType()) {
                                                            case Cell.CELL_TYPE_BOOLEAN:
                                                                        bytesRead++;
                                                                        currentString.append(cell.getBooleanCellValue() + "\t");
                                                                        break;
                                                case Cell.CELL_TYPE_NUMERIC:
                                                                        bytesRead++;
                                                                        currentString.append(cell.getNumericCellValue() + "\t");
                                                                        break;
                                                            case Cell.CELL_TYPE_STRING:
                                                                        bytesRead++;
                                                                        currentString.append(cell.getStringCellValue() + "\t");
                                                                        break;
                                                            }
                                                }
                                                currentString.append("\n");
                                    }
                                    is.close();
                        } catch (IOException e) {
                                    LOG.error("IO Exception : File not found " + e);
                        }
                        return currentString.toString();
            }
            public long getBytesRead() {
                        return bytesRead;
            }
}

EXCEL MAPPER 
(HAVING MAPPER LOGIC)

package com.poc;

import java.io.IOException;

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

public class PocMapper extends Mapper<LongWritable, Text, Text, Text> {
public void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
try {
if (value.toString().contains("RTL_NAME") && value.toString().contains("TYPE_OF_CRAWLING"))
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 dr1 = data.trim().replaceAll("\\s+", "\t");
String[] str1 = dr1.split("\t");

int id = (int) Double.parseDouble(str1[0]);
int regprice = (int) Double.parseDouble(str1[6]);
int rebate = (int) Double.parseDouble(str1[7]);
int saleprice = (int) Double.parseDouble(str1[5]);
String dr = Integer.toString(id) + "\t" + str1[1] + "\t" + str1[2] + "\t" + str1[3] + "\t" + str1[4]+ "\t" + Integer.toString(saleprice) + "\t" + Integer.toString(regprice) + "\t"+Integer.toString(rebate) + "\t" + str1[8];

context.write(new Text(""), new Text(dr));
}
} catch (Exception e) {
e.printStackTrace();
}
}
}

EXCEL REDUCER 
(HAVING REDUCER LOGIC)

package com.poc;

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 PocReducer extends Reducer<Text, Text, IntWritable, Text> {
MultipleOutputs<IntWritable, Text> mos;

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

@Override
public void reduce(Text k1, Iterable<Text> k2, Context context) throws IOException, InterruptedException {
while (k2.iterator().hasNext()) {
String sr = k2.iterator().next().toString();
String sr1 = sr.trim().replaceAll("\\s+", "\t");

String[] str1 = sr1.split("\t");

int regprice = Integer.parseInt(str1[6]);
int rebate = Integer.parseInt(str1[7]);
int saleprice = Integer.parseInt(str1[5]);
String dr = str1[0] + "\t" + str1[1] + "\t" + str1[2] + "\t" + str1[3] + "\t" + str1[4] + "\t" + str1[5]
+ "\t" + str1[6] + "\t" + str1[7] + "\t" + str1[8];
if (str1[2].equalsIgnoreCase("BS")) {
if (saleprice < 100 && rebate > 50) {
mos.write("HighBuzzProducts", null, new Text(dr), "/Retail/HighBuzzProducts");
} else if (regprice < 150 && rebate > 25 && rebate < 50) {
mos.write("NormalProducts", null, new Text(dr), "/Retail/NormalProducts");

} else if (str1[4].length() > 100) {
mos.write("RareProducts", null, new Text(dr), "/Retail/RareProducts");
} else {
mos.write("OtherProducts", null, new Text(dr), "/Retail/OtherProducts");
}
} else if (str1[2].equalsIgnoreCase("ODC")) {
if (saleprice < 150) {
mos.write("OnDemandCrawlProducts", null, new Text(dr), "/Retail/OnDemandCrawlProducts");
} else if (str1[8].equalsIgnoreCase("IN_STOCK")) {
mos.write("AvailableProducts", null, new Text(dr), "/Retail/AvailableProducts");
} else {
mos.write("OtherProducts", null, new Text(dr), "/Retail/OtherProducts");
}
} else {
mos.write("OtherProducts", null, new Text(dr), "/Retail/OtherProducts");
}

}
}

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

EXECUTING THE MAP REDUCE CODE

hadoop jar poc.jar com/poc/PocDriver /Input/POC.xls /Poc


Goto Firefox and open name node page by following command:

http://localhost:50070 and browse the file system , then click on HealthCarePOC directory to check the files created.





3.     PIG SCRIPT

A = LOAD '/Retail/' USING PigStorage ('\t') AS (id:int, Name:chararray, crawl:chararray, produrl:chararray, tittle:chararray, sale:int, reg:int, rebate:int, stockinfo:chararray);

B = DISTINCT A;
DUMP B; 





PigScript2.pig

A = LOAD '/Retail/' USING PigStorage ('\t') AS (id:int, Name:chararray, crawl:chararray, produrl:chararray, tittle:chararray, sale:int, reg:int, rebate:int, stockinfo:chararray);

B = DISTINCT A;
C = ORDER B BY id;
STORE C INTO '/RETAILPOC'; 




4.     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 RETAIL;";



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


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

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



sqoop eval --connect jdbc:mysql://localhost/RETAIL --username root --password root --query "create table retailpoc(id int, name varchar(50), crawl varchar(50), produrl varchar(200), tittle varchar(200), sale int, reg int, rebate int, stockinfo varchar(50));";



sqoop export --connect jdbc:mysql://localhost/RETAIL--table retailpoc --export-dir /RETAILPOC --fields-terminated-by '\t';





5.     STORE THE PIG OUTPUT IN A HIVE EXTERNAL TABLE

goto hive shell using command:

hive

show databases;
create database RetailPOC;
use RetailPOC;



create external table retailpoc(id int, Name string, crawl string, produrl string, tittle string, sale int, reg int, rebate int, stockinfo string)
row format delimited
fields terminated by '\t'
stored as textfile location '/RETAILPOC';





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...