Posted in

How to Perform RDBMS CRUD Operations with Hadoop MapReduce Integration

This article introduces the way to perform RDBMS operations with Hadoop integration. Hadoop is a trending technology these days and to understand the subject, you need to clear some basic facts about this technology. In this post, experts will explain how to read the RDBMS data and manipulate it with Hadoop MapReduce and write it back to RDBMS.

We are introducing a way to perform simple RDBMS read and write operations using Hadoop. Now a days Hadoop is a trending technology stack in Big Data industry.

It is always a common case where giant company or organization which is using RDBMS, now want to migrate their stack to Hadoop.

For using Hadoop as well as other Big Data components, Data transformation is an important task.

In this article I am going to explain how to read the data from RDBMS, manipulate the data using Hadoop Map Reduce paradigm and them write the data back to RDBMS.

Though it is a simple case in which source and sink for the data is same RDBMS, This technique or concepts can be used to perform operation using any type of source and sink. In this article I have used MySql as RDBMS.

Problem:

Consider a simple case, where there is a large chain retail store in which millions of transactions are performed during a day. At the end of the day/Month/year the data is available for analysis so that the Business developers can analyze the data and can take the decisions.

In our example the database of a transaction table is used and we are interested to count the number of unique products sold during the day/Month/Year.

Solution:

In order to solve the problem we are using Hadoop MapReduce example for processing and executing the massive count operation.

We have wrote Custom datatypes for reading and writing the data from RDBMS.

Environment:

Java : 1.6.0_31
Hadoop : 2.2.0
MySQL : 5.5

Here is my solution code,

DBInputWritable.java

importjava.io.DataInput;

importjava.io.DataOutput;

importjava.sql.PreparedStatement;

importjava.io.IOException;

importjava.sql.SQLException;

importjava.sql.ResultSet;

importorg.apache.hadoop.io.Writable;

importorg.apache.hadoop.mapreduce.lib.db.DBWritable;

publicclassDBInputWritableimplementsWritable,DBWritable

{

     privateinttransactionID;

     private String productsBought;

    

     @Override

     publicvoid write(PreparedStatementstatement) throwsSQLException {

          statement.setInt(1, transactionID);

          statement.setString(2, productsBought);

     }


     @Override

     publicvoidreadFields(ResultSetresultSet) throwsSQLException {

          transactionID = resultSet.getInt(1);

          productsBought = resultSet.getString(2);        

     }


     publicintgetTransactionID()

     {

          returntransactionID;

     }


     publicvoidsetTransactionID(inttransactionID)

     {

          this.transactionID = transactionID;

     }


     public String getProductsBought()

     {

          returnproductsBought;

     }


     publicvoidsetProductsBought(String productsBought)

     {

          this.productsBought = productsBought;

     }


     @Override

     publicvoid write(DataOutputout) throwsIOException {

     }


     @Override

     publicvoidreadFields(DataInputin) throwsIOException {

     }

}

 

DBOutputWritable.java

importjava.io.IOException;

importjava.io.DataInput;

importjava.io.DataOutput;

importjava.sql.SQLException;

importjava.sql.ResultSet;

importjava.sql.PreparedStatement;

importorg.apache.hadoop.io.Writable;

importorg.apache.hadoop.mapreduce.lib.db.DBWritable;


publicclassDBOutputWritableimplementsWritable,DBWritable{   

     intcount;

     String productName;

   

     publicDBOutputWritable(){}

     publicDBOutputWritable(String productName, intcount)

     {

          this.productName=productName;

          this.count = count;

     }

    

     @Override

     publicvoid write(DataOutputout) throwsIOException {

     }

     @Override

     publicvoidreadFields(DataInputin) throwsIOException {

     }


     @Override

     publicvoid write(PreparedStatementstatement) throwsSQLException {

          statement.setString(1, productName);

          statement.setInt(2, count);        

     }

     @Override

     publicvoidreadFields(ResultSetresultSet) throwsSQLException {

          productName = resultSet.getString(1);

          count = resultSet.getInt(2);

     }

}

DBDriverTest.java

importjava.io.IOException;

importorg.apache.hadoop.conf.Configuration;

importorg.apache.hadoop.conf.Configured;

importorg.apache.hadoop.io.NullWritable;

importorg.apache.hadoop.io.LongWritable;

importorg.apache.hadoop.io.IntWritable;

importorg.apache.hadoop.io.Text;

importorg.apache.hadoop.mapreduce.Job;

importorg.apache.hadoop.mapreduce.Reducer;

importorg.apache.hadoop.mapreduce.Mapper;

importorg.apache.hadoop.mapreduce.lib.db.DBConfiguration;

importorg.apache.hadoop.mapreduce.lib.db.DBInputFormat;

importorg.apache.hadoop.util.Tool;

importorg.apache.hadoop.mapreduce.lib.db.DBOutputFormat;

importorg.apache.hadoop.util.ToolRunner;

importcom.work.hadoop.mapreduce.datatypes.DBInputWritable;

importcom.work.hadoop.mapreduce.datatypes.DBOutputWritable;


publicclassDBDriverTestextends Configured implements Tool {


     enumDBCounte

r     {

          TotalProductsFound,

          UniqueProducts

     }

             

     staticclassDBMapperextends

              Mapper<LongWritable, DBInputWritable, Text, IntWritable> {

          privateIntWritableone = newIntWritable(1);

          protectedvoid map(LongWritableid, DBInputWritablevalue, Context context){

              try {

                   String[] keys = value.getProductsBought().split(",");


                   for (String key : keys) {

                        context.write(new Text(key), one);                      

                        context.getCounter(DBCounter.TotalProductsFound).increment(1);

                       }

              } catch (IOExceptione) {

                   e.printStackTrace();

              } catch (InterruptedExceptione) {

                   e.printStackTrace();

              }

          }

     }


     staticclassDBReducerextends

              Reducer<Text, IntWritable, DBOutputWritable, NullWritable> {

          @Override

          protectedvoid reduce(

                   Text key,

                   Iterable<IntWritable>values,

                   Reducer<Text, IntWritable, DBOutputWritable, NullWritable>.Context context)

                   throwsIOException, InterruptedException {

              inttotal = 0;

              for (IntWritablevalue : values) {

                   total += value.get();

              }


              context.write(newDBOutputWritable(key.toString(), total),

                        NullWritable.get());

               context.getCounter(DBCounter.UniqueProducts).increment(1);            

          }

     }

     publicint run(String[] args) throws Exception {

          Configuration conf = newConfiguration();

          DBConfiguration.configureDB(conf, "com.mysql.jdbc.Driver", // drive

r                                                                                 // class

                   "jdbc:mysql://localhost:3306/retail_db", // dburl

                   "root", // user name

                   "root"); // password

          Job job = Job.getInstance(conf);

          job.setJarByClass(DBDriverTest.class);

          job.setMapperClass(DBMapper.class);

          job.setReducerClass(DBReducer.class);

          job.setMapOutputValueClass(IntWritable.class);

          job.setMapOutputKeyClass(Text.class);

          job.setOutputKeyClass(DBOutputWritable.class);

          job.setOutputValueClass(NullWritable.class);

          job.setOutputFormatClass(DBOutputFormat.class);

          job.setInputFormatClass(DBInputFormat.class);


          DBInputFormat.setInput(job, DBInputWritable.class,

                   "retail_master", // input table name

                   null, null,

                   new String[] { "transaction_id", "bought_products"}// Input table columns

                   );


          DBOutputFormat.setOutput(job, "product_analysis", // output table name

                   new String[] { "product_name", "count" } // table columns

                   );

          returnjob.waitForCompletion(true) ? 0 : 1;

     }

     publicstaticvoid main(String[] args) {


          try {

              intresult = ToolRunner.run(new Configuration(),

                        newDBDriverTest(), args);

              System.out.println("job status ::" + result);

          } catch (Exception exception) {

               exception.printStackTrace();

          }

     }

}

Input Data and Database Spool:

I have created 2 different tables using the below sql lines.

 CREATE TABLE `retail_db`.`retail_master` (

  `transaction_id` INT NOT NULL AUTO_INCREMENT,

  `customer_name` VARCHAR(45) NOT NULL,

  `bought_products` VARCHAR(255) NOT NULL,

  `amount` INT NULL,

  PRIMARY KEY (`transaction_id`));

CREATE TABLE `retail_db`.`product_analysis` (

  `product_name` VARCHAR(40) NOT NULL,

  `count` INT NULL,

     PRIMARY KEY (`product_name`));

There are 2 tables retail_master which contains the Transaction related data and product_analysis which is the destination table for our analysis.

Currntly my retail_master table contais below data.

Retail Master Table

At the end of the Hadoop job run you will find the product_analysis will contain the data like below image.

product analysis

Code Walk Through:

  • We are using 2 different custom hadoop data types for reading and writing the data from and to the table.
  • DBInputWritable is a custom Writable datatype for reading the data from RDBMS table.
  • DBOutputWritable is a custom Writable datatype for writing the data to RDBMS table.
  • In the above code I used the necessary columns only for simple analysis.
  • This code can be used for tweaking and further analysis.

I hope this program will help you to learn hadoopmapreduce with RDBMS operation.

Hope this article will help you in understanding the way to perform RDBMS operations with Hadoop integration. You can share your experience with professionals by commenting below. Thanks for your time.

Author Bio:

This Article has been written by Ethan Millar, working with Aegis SoftTech as a Technical writer since last 5 years. He especially write articles in Java, Hadoop, CRM, .Net and other major programing languages. The main object to write this article is to Perform RDBMS CRUD Operations with Hadoop MapReduce Integration. The conclusion has been drawn after partial research and implementation with Hadoop integration by Hadoop developers at Aegis SoftTech.

I am Ethan work with Aegis Soft Tech as a software developer. I have vast experience in Hadoop, CRM as well as Java application development.

Privacy Overview

This website uses cookies so that we can provide you with the best user experience possible. Cookie information is stored in your browser and performs functions such as recognising you when you return to our website and helping our team to understand which sections of the website you find most interesting and useful.