Sunday, July 19, 2015

Name Node


NameNode acts as the data manager between the requested client and content holder - DataNodes.

On start up, DataNodes send the block report to NameNode on every hour, along with the heartbeats on every 3 seconds.  NameNode keeps track of every data change.  At any given time, the NameNode has a complete view of all DataNodes in the cluster, their current health, and what blocks they have available.  File to block mapping on the NameNode is stored on disk.

NameNode does not directly send requests to DataNodes. It uses replies to heartbeats to send instructions to the DataNodes.  The instructions include commands to replicate blocks to other nodes, remove local block replicas, re-register and send an immediate block report, and shut down the node.  NameNode stores its filesystem metadata on local filesystem disks. 2 Key files are (1)FsiImage (2)EditsLog

FsiImage contains the complete snapshot of filesystem at iNode level.  iNode is an internal representation of a file or directory's metadata and contains such information as the file's replication level, modification and access times, access permissions, block size, and the blocks a file is made up of.  This design makes not to worry about the changing DataNodes' hostname or IP address.

Edits file (journal) contains only incremental modifications made to the metadata. It uses a write ahead log which reduces I/O operations to sequential, append-only operations, which avoids costly seek operations and yields better overall performance.

On NameNode startup, the fsimage file is loaded into RAM and any changes in the edits file are replayed, bringing the in-memory view of the filesystem up to date.  NameNode filesystem metadata is served entirely from RAM. This makes it fast, but limits the amount of metadata a box can handle. Roughly 1 million blocks occupies roughly 1 GB of heap.

Thus, NameNode performs the filesystem operations in the highly distributed methodology for the Client request(s).

Sunday, July 12, 2015

HDFS Design


In principle, HDFS has a block size higher than most other file systems. The default is 128M and some go as high as 1G. Files in HDFS are write once.

There are three daemons that make up a standard HDFS cluster.
  1. NameNode - 1 per cluster. Meta data's centralized server to provide a global picture of the filesystem's state.
  2. Secondary NameNode - 1 per cluster.  Performs internal NameNode transaction log check pointing.
  3. DataNode - Many per cluster.  Stores block data (contents of files).

NameNode stores its filesystem metadata on local filesystem disks in a few different files, but the two most important of which are fsimage and edits.  Fsimage contains a complete snapshot of the filesystem metadata including a serialized form of all the directory and file inodes in the filesystem.  Edits file (journal) contains only incremental modifications made to the metadata, which acts as write ahead log.

Secondary NameNode is not only backup of NameNode but also shares the workload via checkpointing process.  In which secondary NameNode applies the updates from the edits file to the fsimage file and sends it back to the primary.  Checkpointing is controlled by duration (default 60 mins) and/or file size and/or transaction count of Edits file.

Daemon responsible for storing and retrieving block (chunks of a file) data is called the DataNode .  Datanodes regularly report their status to the NameNode in a heartbeat mode; default 3 mins.  It sends Block Report (list of all usable blocks of DataNode disks) to NameNode; default 60 mins.

3 key daemons of HDFS Architecture, is represented in the attached diagram.

Saturday, July 11, 2015

HDFS Goals


In last tip, we saw Google’s Whitepapers on Big Data.  Apache Hadoop has been originated from Google’s Whitepapers:

  1. Apache HDFS is derived from GFS  (Google File System).
  2. Apache MapReduce is derived from Google MapReduce
  3. Apache HBase is derived from Google BigTable.


This Tip#2, is on HDFS goals & roles in Big Data.  What is HDFS?

HDFS is a distributed and scalable file system designed for storing very large files with streaming data access patterns, running clusters on COMMODITY hardware. HDFS is not a POSIX-compliant filesystem.

In HDFS, each machine in a cluster stores a subset of the data (blocks) that makes up the complete filesystem. Itz metadata is stored on a centralized server, acting as a directory of block data and providing a global picture of the filesystem's state.

Top 5 Goals of HDFS

  1. Store millions of large files, each greater than tens of gigabytes, and filesystem sizes reaching tens of petabytes.
  2. Use a scale-out model based on inexpensive commodity servers with internal JBOD  ("Just a bunch of disks") rather than RAID to achieve large-scale storage. 
  3. Accomplish availability and high throughput through application-level replication of data.
  4. Optimize for large, streaming reads and writes rather than low-latency access to many small files. 
  5. Support the functionality and scale requirements of MapReduce processing.

Thursday, July 9, 2015

ThankYou Card


A lot of people don’t realize this, but the Pulse app was built as a class project at Stanford University in 2010. Pulse today powers a lot of content you see on LinkedIn’s homepage feed. They reached an incredible milestone for the LinkedIn publishing platform: 1 million professionals have now written a post on LinkedIn.

Over 1 million unique writers publish more than 130,000 posts a week on LinkedIn. About 45% of readers are in the upper ranks of their industries: managers, VPs, CEOs, etc. The top content-demanding industries are tech, financial services and higher education. The average post now reaches professionals in 21 industries and 9 countries.

In this big metric, I'm also taking tiny contribution @ Pulse.  In conjunction with 1 million post celebration, I just received 'Thank You' card from LinkedIn.

If you’re writing on LinkedIn or interested in starting, join the Writing on LinkedIn Group. Not sure “What’s stopping you?”.  Continuous Learning & Continuous Sharing is mantra for our industry.

Sunday, July 5, 2015

Google WhitePaper


Dear readers, I recently got few requests to share few 'Useful Tips' on Big Data ecosystem.

Tip is a small piece or part fitted to the end of an object.  Let me fill Big Data Tips during this quarter - Q3 2015.

For those interested in the history, the super base class of Big Data ecosystem is Google's whitepaper.

The first, presented in 2003, describes a pragmatic, scalable, distributed file system optimized for storing enormous datasets, called "Google File system", or GFS by Sanjay Ghemawat, Howard Gobioff, and Shun-Tak Leung. In addition to simple storage, GFS was built to support large-scale, data-intensive, distributed processing applications.

The following year-2004, another paper, titled "Map-Reduce: Simplified Data Processing on Large Clusters" was presented by Jeffrey Dean and Sanjay Ghemawat, defining a programming model and accompanying framework that provided automatic parallelization, fault tolerance, and the scale to process hundreds of terabytes of data in a single job over thousands of machines.

When paired, these two systems could be used to build large data processing clusters on relatively inexpensive, commodity machines. Google White papers are industry break through and directly inspired the development of HDFS and Hadoop MapReduce, respectively.

Google's WhitePaper References are available at:
http://static.googleusercontent.com/media/research.google.com/en//archive/gfs-sosp2003.pdf
http://static.googleusercontent.com/media/research.google.com/en//archive/mapreduce-osdi04.pdf

Stay tuned for continuous tip & trick to ()l)earn more.

Wednesday, July 1, 2015

Netflix Migration


Last weekend, I read an interesting business use case on Big Data - Cassandra migration. We know about the streaming media leader Netflix; Itz about their Big Data migration from Traditional storage.
About Netflix:
Netflix is the world’s leading Internet television network with more than 48 million streaming members in more than 40 countries. It has successfully shifted its business model from DVDs by mail to online media and today leads the streaming media industry with more than $1.5 billion digital revenue. Netflix dominates the all peak-time Internet usage and its shares continue to skyrocket with soaring numbers of subscribers.
Trigger Point:
In 2010, Netflix began moving its data to Amazon Web Services (AWS) to offer subscribers more flexibility across devices with 'Cloud-First/Mobile-First' Strategy. At the time, Netflix was using Oracle as the back-end database and was approaching limits on traffic and capacity with the ballooning workloads managed in the cloud.

Interestingly, the entire migration of more than 80 clusters and 2500+ nodes was completed with only two engineers
Use Case: 
Systems that understand each person’s unique habits and preferences and bring to light products and items that a user may be unaware of and not looking for. In a nutshell, personalizes viewing for over 50 Million Customers.
Challenges:
Challenges related to the given use case:
  • Affordable capacity to store and process immense amounts of data 
  • Data Volume more than 2.1 billion reads and 4.3 billion writes per day
  • Single point of failure with Oracle’s legacy relational architecture
  • Achieving business agility for international expansion
Solution
Big Data Cassandra delivers a persistent data-store 
  • 100% up-time and cost effective scale across multiple data centers
  • DataStax expert support the results
  • It delivers a throughput of more than 10 million transactions per second
  • Effortless creation/management of new data clusters across various regions
  • Capture of every detail of customer viewing and log data
The attached Netflix Deployment Diagram is published in Netflix Tech Blog @http://techblog.netflix.com/2012/06/annoucing-archaius-dynamic-properties.html

Tuesday, June 23, 2015

AWS Enterprise Summit 2015


AWS (Amazon Web Service) Enterprise Summit are designed to educate new customers about the AWS platform; offers existing customers deep technical content to be more successful with AWS.
Today(23 Jun), I had the chance to attend AWS Enterprise Summit at Chennai, India. In a nutshell, it covers keynote address, panel discussion, customers use case on technical track, etc. Herez the highlights of Today's sessions:
Development Focus
"Moore's law" is the observation that, over the history of computing hardware, the number of transistors in a dense integrated circuit has doubled approximately every two years.
Itz true and reflected in IT industry. In 1971, Intel's first processor 4004 contained 2,300 embedded transistors to execute. Now, in 2015, Intel's 18-core Xeon Haswell-EP has over 5.5 billion transistors. Amazing growth, right!!!
As hardware is rapidly expanding its band, software's development focus is migrating in the below order:
  • Mainframe - 1970s
  • Personal Computer (PC) - 1980s
  • Data Centre (RDBMS) - 1990s
  • High Performance Computing (HPC) - 2000s
  • Cloud; Horizontal Scaling - 2010s
Amazon Web Service (AWS) platform is the key player in 2010s era.

Business Model


Traditional IT
Emerging IT
Business Model

Large Capital Expenditure (CapEx)
Low variable on demand CapEx
Cost Reduction

CapEx focused
Operation Expenditure (OpEx) focused
Cost Model

Basic Computing
Broad & Deeper Platform
Scalability

Responsible for periodic upgrade
New features arrive daily
Lead to Innovation

Slow to roll new feature
Ready to use rapid feature
Agile

Traditional DR
Automatic DR
Improved Availability

Costing Strategy

Cost drop is achievable by TWO key factors in the business theory.
1. Large customer base
2. Better economy of scale
In alignment with this costing strategy, Amazon had the multiple historical price reductions i.e. 48 price slashing since 2006.

Pricing Philosophy

Like mobile plan/cost, it starts from base to advanced package based on the user's demand. It is upto the customer to select their choice. AWS has the multiple pricing plan as below:

Purchase Model
Description
Usage

Free Tier
With free usage and no commitment
For PoCs and getting started

On Demand
Pay by the hour with no long term commitment
For spiky / seasonal workloads

Reserved
Low one-time payment with significant discount
For committed utilization

Spot
Bid for unused capacity, fluctuates based on demand/supply
Time intensive or transient workloads

Dedicated
Launch instance run on hardware dedicated to a single customer
Highly intensive or compliance loads

Data Growth Trend

It is interesting to observe the data growth @ our industry
  • 7.9 Zetta Bytes of Data Persistence
  • 90% of data growth just in last 2 years
  • 5+ billion devices usage
  • 966 Exa Bytes transfer rate

Big Data Platform

AWS Platform has the end-end solution for Big Data use cases with their own cloud based tools:
  1. Kinesis - allows for large data stream processing and real-time analytics
  2. Elastic MapReduce (EMR) - API styled web service that uses Hadoop ecosystem
  3. Relational Data Store (RDS) - Scalable relational database in the cloud
  4. Simple Storage Service (S3) - Opt for virtually unlimited cloud & internet storage 
  5. RedShift - Fast, fully managed, petabyte-scale data warehouse in the cloud
  6. Dynamo DB - Fully managed NoSQL database service that provides fast and predictable performance with seamless scalability

Closing Note

Statistics indicates that 6% Enterprise into Big Data and 9% Enterprise are into Cloud platform to adapt this data growth.
Thus, AWS (Amazon Web Service) is matured enough in each and every space of Big Data and Cloud platform.  In fact, their own line of business (amazon.com) in world's leading online store, is helped to set the validity of their solution/products, before serving to the industry.  Awesome work, AWS.