Thursday, November 1, 2012

Pentaho solution for Sqoop call with dynamic partitions

Pentaho big data release doesnt have a step for SQOOP as of this writing..
Simple solution is to use a "Shell script" stage to call sqoop..

But our requirement has a twist .. and most of sqoop users might have as well..
- capture incremental data for certain tables and keep it in a new hive partition
- Run the sqoop extract for a window of dates where each day's data goes to a partition (DAILY , HOURLY, MONTHLY etc)

Design


The solution is to use a DB for configuration and pentaho to frame sqoop calls

JOB


We are calling the pensqoop transformation to frame the list of sqoop calls to run
shell stage actually runs the sqoop script.. here check the "execute for each input row " so that this script will be called for each sqoop call framed in previous step.

TRANSFORMATION



1. get table configurations (like db connection, incremental column to use for incremental data capture etc)
2. switch the flow for DAILY partitions or FULL refresh
3. frame sqoop call using string manipulation at javascript stage
4. "copy rows to result" to send the output to calling job

Wednesday, June 20, 2012

Hive Maintenance

Hive log filling up tmp space

/tmp is usually the neglected folder in any unix environment, but there is where Hive is going to place all its log files. If you didnt have enough space allocated for this then your scheduled sqoop or hive queries are going to fail because of silly reason - no temp space. Not worth it..


1.       Create a cron to remove the /tmp/<username>/*.txt files frequenty
2.       Change the hive log directory to a different location which is monitored by ops ..
I I prefer the 2nd since this change will not remove any files in the process. 
t   To make the configuration change got to /etc/hive/conf/hive-site.xml and add the following property

<property>
  <name>hive.querylog.location</name>
  <value>/var/<username>/tmp</value>
  <description>Directory where structured hive query logs are created</description>
</property>



Freeing up unused HDFS memory

Two directories we can concentrate to free-up a lot of HDFS memory

1. /tmp/hive-<username>/
2.  /local/hadoop/mapred/staging/user/.staging/job*

If you have Hue you can cleanup some of long unused saved reports
   1. /tmp/hive-beeswax-<username>/

If you can limit the files by date less than the current month, you are safe..(since this is delete operation)
Hadoop, Hive and Hue does have demon or code to clean-up these folders (Im not sure) but 
Usually all these files are there because of orphaned mapreduce jobs. Like when you CTRL+C instead of formally  "hadoop -job kill <job_id>"

And ofcourse scanning the hdfs directory and hive directory for junk files will also help..

the following are default configuration directories

for beeswax cleanup -- hadoop fs -rmr /tmp/hive-beeswax-*/hive*   

for hive tmp --   hadoop fs -rmr /tmp/hive-bhchandr/hive*
older job files --  hadoop fs -rmr /local/hadoop/mapred/staging/<application_user>/.staging/job_2012*

Free Sqoop temp space

remove the compile folders for sqoop job
/tmp/sqoop-<username>/compile/

remove log files older than 1/2 day at jobtracker and data-nodes

 /var/log/hadoop-*-*/userlogs/*            

recover from safemode

Hadoop can go into safemode when the local directory mapped to HDFS is full.. this usually happens when your HDFS files usedup all the space.. but the trick is you cannot remove files unless you recover from safemode.

So first remove some local files from
/opt/local/hadoop/mapred/local/taskTracker
/opt/local/hadoop/mapred/local/taskTracker/distcache
/opt/local/hadoop/mapred/local/taskTracker/<username>

then
hadoop dfsadmin -safemode leave

since Hadoop is out of safemode.. you can cleanup some HDFS using the steps in the first sections


/var/log/hadoop-0.20/history/done

Thursday, May 24, 2012

Hive - digging deeper into metastore

We usually use information_schema or database metadata tables to query the tables, columns, indexes ..in the case of traditional databases. How about the same in Hive?

In the case of Hive there is a metastore which acts as a metadata for the databases, Hive uses this database to store the tables, partitions, databases, serde in this database. Say you want to know the tables in a database or physical(HDFS) location of the tables or similar column names across the tables..

I have no idea about the default derby DB.. so letme talk about metastore in a traditional database like Mysql..( If you dont have your metastore in an external DB you can follow this link to do so.. https://ccp.cloudera.com/display/CDHDOC/Hive+Installation#HiveInstallation-ConfiguringtheHiveMetastore)

Hive metastore has minimal tables when compared to the metadata layer of a traditional database but Im sure the metadata schema will get bigger and complex in future.

Lets go through the most desirable tables ...

COLUMNS
DBS  -- the list of schemas
PARTITIONS
TBLS - tables

Querying these tables will give a basic idea and the datamodel of Hive metastore ( I couldnt find any over the internet), understanding of these tables are essential if you are seriously into Hive and have some production data.

for eg.. In a weird case we found some partitions of one database is being loaded into HDS location of another database (may be a code issue ) but using the following query Im able to narrow down to the affected tables..

select distinct a.tbl_id,a.tbl_name
    from TBLS a
        join PARTITIONS c on (a.TBL_ID=c.TBL_ID)
        join SDS b on (c.SD_ID=b.SD_ID)
 where b.location like '%sbx.db%' and a.DB_ID=1;


Wednesday, April 25, 2012

Secondary Sort using Python and MRJob

MRjob is an excellent module developed as open source project by Yelp.. I chosed mrjob because of the following features


  1. It provides a seamless JSON reader and writer (i.e the mapper can read json lines and convert them into   lists)
  2. We can test hadoop job locally (in windows or unix) on a small dataset without actually using huge hdfs files (quick !!)
  3. Can orchestrate many mappers and reducers in the same code 

My task is to parse json formatted web log files and parse them, say the columns are sessionid,stepno and data. so the psuedo-code
  1. Read the json files using mrjob protocol
        DEFAULT_INPUT_PROTOCOL = 'json_value'
        DEFAULT_OUTPUT_PROTOCOL = 'repr_value'  #  output is delimited

  2.  yield sessionid, (sessionid,stepno,data) from mapper
                      the Mapreduce will make sure that all sessionids(key) goes to same mapper.. with the remaining values sent as a dictionary (value) to make it easier for us to srt in reducer

 3.  Use sorted from itertools of python module to sort by stepnumber in reducer

 def reducer(self, sessionId, details):
                sdetail = sorted(details, key=lambda x: x[1])  # sorting by stepno for each session
                for d in sdetail:
                        line_data='\t'.join(str(n) for n in d)

We are doing the secondary sort to scan through each events as the sequence is very important to do the funnel analysis of logs..

Complete code :

import sys,time
#sys.path.append('/usr/lib/python2.4/site-packages/')
from mrjob.job import MRJob
from mrjob.protocol import JSONValueProtocol
from itertools import groupby
from operator import itemgetter, attrgetter

class uet(MRJob):
        DEFAULT_INPUT_PROTOCOL = 'json_value'
        DEFAULT_OUTPUT_PROTOCOL = 'repr_value'

        def mapper(self, _, line):
                        sessionId = line['sessionId']
                        data = line['data']
                        if len(sessionId) < 13:
                                for i in range(len(data)):
                                        no = data[i]['no']
                                        yield sessionId,(sessionId, no,data)

        def reducer(self, sessionId, details):
                sdetail = sorted(details, key=lambda x: x[1])  # sorting by stepno for each session
                for d in sdetail:
                        line_data='\t'.join(str(n) for n in d)
                        print str(line_data)


if __name__ == '__main__':
    uet.run()



Tuesday, April 17, 2012

Top NBA Players - by twitter followers

I really dont have any idea of how advertising companies decide on the price for certain celebrities. Because its very hard to measure the direct relationship to the sales and to derive an ROI. Also Im not sure if any celebrity with enormous following can make a big impact. I did data collection for fun to see who is actually more popular in Twitter and have more following. I used infochimps rest api calls to get the aggregated information and formatted files to make it readable by Tableau. (used Python for ETL).. seems like SHAQ even after out of the league has a greater fan following than active player.. I havent included Kobe and Rose.. and I tried my best to get the official ids of each player..


Tuesday, April 10, 2012

DW in Hive - handling big dimensions

This is a an issue bothering our reporting data quality. How to handle big dimension tables in Hive data warehouse. How to balance performance and data-quality..

Problem statement:
Currently dimension hive table (dim_customer) is partitioned by date. The daily incremental creates the new partition so that we can improve the performance at the reporting side. This poses 2 critical issues
1. slowly changing dimension is lost
2. Compromise in data-quality
3. Have to filter by the dimension table in the reporting

Solution:
The only way to solve this problem is to have the dimension table as one big Hive table instead of partitions. But this creates issues with the refresh strategy and overhead on reporting Hive query..

The following is a recipe to solve this block.. step.1 is certainly the priority

1. Increase processing power
Hadoop is not only about mega storage it is also about mega processing ..so if we process big files then we got to have good number of nodes. say we have 30 nodes to process partitioned dimension table..we have to move to 120 nodes for single dimension file strategy. Processing power is tripled - it is directly proportional !.

2. Use SQOOP merge
We cannot extract the whole table from transactional system every time..source transactional systems might not allow.. we can only capture the change data. SQOOP merge comes handy for this purpose. We can overwrite only the incremental records in the hive table (type 1 SCD). Again this is a Map reduce program.. we need processing power..

3. Use Bucketed Hive tables
Hive performance block comes while joining. We can create a table with bucketing.. like hashing index on the customer_id.


If the tables being joined are bucketized, and the buckets are a multiple of each other, the buckets can be joined with each other. If table A has 8 buckets are table B has 4 buckets, the following join

SELECT /*+ MAPJOIN(b) */ a.KEY, a.value
FROM a JOIN b ON a.KEY = b.KEY
can be done on the mapper only. Instead of fetching B completely for each mapper of A, only the required buckets are fetched. For the query above, the mapper processing bucket 1 for A will only fetch bucket 1 of B. It is not the default behavior, and is governed by the following parameter

Bucketing and Sqoop merge requires good planning and metadata management..
Managing Hadoop from the scratch is challenging as we bump into the limits sooner, we need to adapt quickly else data might grow beyond limits. One advantage( and complexity) is that the internal processing (mapreduce) is open and its upto the developer to improve. And the biggest advantage of all is scaling out.. you can add nodes easily to really make a difference..






Thursday, March 29, 2012

Data science -- the cool scientist without white gowns

"Scientist" is a cool word during my school days, wanna become one but donno on what. All I see as scientists wore white gown with colorful liquids around. But later during college and working days scientists seemed to be boring people with no personal life, accumulated in educational institutes with college kids helping around.

Recently there is a profile called "Data scientist" all over the BI market and started appearing in every article where Hadoop / Big data is there. It must be some part what a BI/DW person is doing with some specialization. Yes it is the formula,as for me..

BI + big-data + statistics + scripting + visualization = Data scientist

Ok can be scientist and work for a corporate or invent something new for you name? possibly...
But seems like a lot to cover , learn and experience at work. Maybe not really if we are in right job. Im just listing the very higher level outline..(and it is not limited to..) and my intention is not to oversimplify, but certainly to simplify the puzzle..

BI
- DW work like ETL , databases, SQL with exposure to enterprise setup. Collecting data from heterogenous data sources. Log analysis. Dimensional modelling, DW architecture. DB performance.

Big-data
Hadoop is the first thing comes to mind for Big-data.. but good to know about noSQL dbs.
Mapreduce - shared nothing architecture - need for MR - use cases - tools available - pros and cons

Statistics
Basics - application of statistics in real-world - R programming

Scripting
Perl, Python, Java

Visualization
Reporting (I like Tableau), Complex SQLs , Ability to tell a story with data - by whatever way you effectively deliver..

I would like to list some of the coolest learning materials available for above topics..

I believe the thirst for discovery, admiring the hidden secret in the boring pile of data would make a Data Scientist..