UNIT 1: Big Data Fundamentals and Distributed Processing
I. Big Data Core Concepts
What is Big Data?
Big Data refers to extremely large, complex, and rapidly growing datasets that traditional data processing tools (like relational databases) cannot effectively capture, store, manage, and analyze. It is characterized by the 3Vs (and often more):
| Dimension | Explanation | Key Technologies/Challenges |
|---|---|---|
| Volume | The sheer scale of data (Terabytes, Petabytes, Exabytes). | Scalable storage (HDFS), distributed processing. |
| Variety | The diverse formats and sources: structured (SQL), semi-structured (XML, JSON), unstructured (text, images, video, logs). | Schema-on-read, NoSQL databases, complex parsing. |
| Velocity | The speed at which data is generated, collected, and needs to be processed (real-time/streaming). | Stream processing engines (Spark Streaming, Flink), low-latency ingestion. |
[!TIP] Exam Focus: Always define Big Data and elaborate on all three Vs (Volume, Variety, Velocity) with concrete examples for each (e.g., Volume: social media posts; Variety: sensor data + images; Velocity: stock ticker data).
Challenges with Big Data
-
Storage: Storing massive, diverse datasets cost-effectively.
-
Processing: Analyzing data in reasonable time; requires parallel/distributed computing.
-
Analysis: Extracting meaningful insights from noisy, incomplete data.
-
Security & Privacy: Protecting sensitive information across distributed systems.
-
Scalability: Systems must scale horizontally (add more machines) seamlessly.
-
Visualization: Representing complex, high-dimensional results understandably.
Big Data Analytics
The process of examining large and varied datasets to uncover hidden patterns, correlations, market trends, and customer preferences.
-
Overview: Uses advanced analytics (predictive modeling, machine learning) on big data platforms.
-
Real-World Applications:
-
Healthcare: Predictive diagnostics, drug discovery.
-
Retail: Recommendation systems, inventory optimization, sentiment analysis.
-
Finance: Fraud detection, algorithmic trading, risk assessment.
-
Telecom: Network optimization, customer churn prediction.
-
Manufacturing: Predictive maintenance, supply chain analytics.
-
II. Data Preprocessing and Mining Techniques
Data Cleaning and Sampling
-
Data Cleaning: Identifying and correcting (or removing) inaccurate, incomplete, duplicate, or inconsistent data. Crucial as "garbage in, garbage out."
- Methods: Handling missing values (imputation, deletion), smoothing noise, resolving inconsistencies.
-
Sampling: Selecting a representative subset of data for analysis when processing the entire dataset is infeasible.
- Methods: Simple random sampling, stratified sampling (preserves class distribution), reservoir sampling (for streams).
[!TIP] Common Pitfall: Don't confuse sampling (for analysis) with data cleaning (for quality). Both are preprocessing steps.
Classification Techniques
Decision Trees
A flowchart-like tree structure where each internal node tests an attribute, each branch represents an outcome, and each leaf node holds a class label.
-
Construction: Uses recursive partitioning (splitting data based on attribute values) to create homogeneous subsets.
-
Key Algorithms: ID3 (uses Information Gain), C4.5 (uses Gain Ratio, handles continuous attributes), CART (uses Gini Index, for classification & regression).
-
Splitting Criteria:
- Information Gain (ID3): Measures reduction in entropy (impurity).
$$IG(S, A) = Entropy(S) - \sum_{v \in Values(A)} \frac{|S_v|}{|S|} Entropy(S_v)$$
* **Gini Index (CART):** Measures impurity of a split.
$$Gini(S) = 1 - \sum_{i=1}^{c} (p_i)^2$$
-
Advantages: Easy to understand, interpretable, handles both numerical & categorical data.
-
Disadvantages: Prone to overfitting, unstable (small data changes can change tree).
Naive Bayes Classification
Based on Bayes' Theorem with a "naive" assumption of conditional independence among features given the class.
- Bayes' Theorem:
$$P(C|X) = \frac{P(X|C) \cdot P(C)}{P(X)}$$
Where:
* $C$ = class, $$\displaystyle X = (x_1, x_2, ..., x_n) $$ = feature vector.
* $P(C|X)$ = posterior probability.
* $P(C)$ = prior probability of class.
* $P(X|C)$ = likelihood.
-
Naive Assumption: $$\displaystyle P(X|C) = P(x_1|C) \cdot P(x_2|C) \cdot ... \cdot P(x_n|C) $$
-
Classification Rule: Assign class $C$ that maximizes:
$$C_{NB} = \arg\max_C P(C) \prod_{i=1}^{n} P(x_i|C)$$
-
Advantages: Simple, fast, works well with high-dimensional data, requires small training data.
-
Disadvantages: "Naive" independence assumption rarely holds true in reality.
Association Rules
Discovers interesting relationships (correlations) between variables in large databases. Commonly used in market basket analysis.
-
Key Metrics:
- Support ($supp(X)$): Frequency of itemset $X$.
$$supp(X) = \frac{\text{Number of transactions containing } X}{\text{Total number of transactions}}$$
2. **Confidence ($conf(X \Rightarrow Y)$):** Likelihood of $Y$ given $X$.
$$conf(X \Rightarrow Y) = \frac{supp(X \cup Y)}{supp(X)}$$
3. **Lift ($lift(X \Rightarrow Y)$):** Measures strength of rule over random occurrence.
$$lift(X \Rightarrow Y) = \frac{conf(X \Rightarrow Y)}{supp(Y)}$$
* $$\displaystyle lift > 1 $$: Positive correlation (rule is useful).
* $$\displaystyle lift = 1 $$: Independence.
* $$\displaystyle lift < 1 $$: Negative correlation.
-
Algorithm: Apriori (uses candidate generation & pruning based on support threshold).
-
Applications: Cross-selling, product placement, recommendation systems, medical diagnosis.
III. Hadoop Ecosystem and MapReduce
Hadoop Architecture
A framework for distributed storage and processing of very large datasets on clusters of commodity hardware.
-
Core Building Blocks:
-
HDFS (Hadoop Distributed File System): Distributed, scalable, fault-tolerant file system.
-
NameNode: Master server (manages metadata, file system namespace).
-
DataNode: Slave server (stores actual data blocks, typically 128MB/256MB).
-
-
YARN (Yet Another Resource Negotiator): Cluster resource management & job scheduling.
-
ResourceManager: Master (allocates resources globally).
-
NodeManager: Per-node agent (manages containers, monitors resources).
-
ApplicationMaster: Per-application negotiator (requests containers from RM, coordinates execution).
-
-
MapReduce: Programming model for parallel processing.
-
JobTracker (pre-YARN) / ApplicationMaster (YARN): Manages job lifecycle.
-
TaskTracker (pre-YARN) / Container (YARN): Executes map/reduce tasks.
-
-
[!TIP] Exam Focus: Be able to draw and label the Hadoop 2.x (YARN) architecture diagram, clearly differentiating HDFS, YARN, and MapReduce components.
MapReduce Programming Model
A programming paradigm for processing massive datasets in parallel across a cluster.
-
Core Idea:
mapfunction processes key-value pairs to produce intermediate key-value pairs.reducefunction aggregates all values associated with the same intermediate key. -
Data Flow:
Input → Split → Map → Shuffle & Sort → Reduce → Output -
Key Types & Formats:
-
InputFormat: Defines how input files are split and read. (e.g.,
TextInputFormat- line by line). -
Mapper Output Key/Value Types: Intermediate types (
K2, V2). -
Reducer Output Key/Value Types: Final output types (
K3, V3).
-
-
Combiner: A local, mini-reducer that runs after the map task to reduce network traffic by aggregating map output locally. Optional.
-
Partitioner: Controls which reducer a particular intermediate key is sent to. Default is HashPartitioner. Custom partitioners enable load balancing.
-
Examples:
-
Word Count (Classic):
-
Map: (docID, text) → (word, 1) for each word.
-
Reduce: (word, [1,1,...]) → (word, sum).
-
-
Inverted Index:
-
Map: (docID, text) → (word, docID) for each word.
-
Reduce: (word, [docID1, docID2,...]) → (word, sorted list of docIDs).
-
-
IV. High-Level Big Data Tools
Apache Hive
A data warehouse infrastructure built on Hadoop. Provides SQL-like querying (HiveQL) for data summarization, analysis, and transformation.
-
Architecture:
-
Hive CLI / JDBC/ODBC Driver: User interfaces.
-
Driver: Receives HiveQL queries, creates session handles, manages execution plan.
-
Compiler: Parses, validates, transforms query into MapReduce/Tez/Spark execution plan.
-
Metastore: Central repository for metadata (table schemas, partitions, locations). Typically uses MySQL/PostgreSQL.
-
Execution Engine: Executes the plan (via MapReduce/Tez/Spark).
-
HDFS / HBase: Underlying storage.
-
-
User-Defined Functions (UDFs): Custom functions written in Java to extend Hive's built-in functionality.
-
Procedure:
-
Extend
org.apache.hadoop.hive.ql.exec.UDFclass. -
Implement
evaluate()method (can be overloaded). -
Package as a JAR file.
-
Add JAR to Hive session:
ADD JAR <path_to_jar>; -
Create temporary/permanent function:
CREATE TEMPORARY FUNCTION <alias> AS '<fully.qualified.ClassName>';
-
-
Apache Pig
A high-level platform for creating MapReduce programs using a scripting language called Pig Latin.
-
Architecture:
-
Pig Latin Script: User's high-level script.
-
Parser: Checks syntax, validates, builds logical plan (DAG of operators).
-
Optimizer: Performs logical optimizations (e.g., projection pushdown, filter pushdown).
-
Compiler: Converts logical plan into series of MapReduce jobs.
-
Execution Engine: Submits jobs to Hadoop cluster.
-
HDFS: Storage.
DiagramCANVAS: Flowchart: Pig Latin Script → Parser → Logical Plan → Optimizer → Optimized Logical Plan → Compiler → Series of MapReduce Jobs → Execution Engine → Hadoop Cluster (YARN/MapReduce) → HDFS. -
-
Pig Latin Application Flow:
-
Write Script: Create
.pigfile with Pig Latin statements (LOAD,FILTER,GROUP,FOREACH,STORE). -
Execute: Run via
pig -x mapreduce script.pigor Grunt shell. -
Parse & Type Check: Pig parses script, checks semantics.
-
Logical Plan Creation: Builds DAG of operators (e.g.,
LOAD → FILTER → GROUP → FOREACH). -
Logical Optimization: Applies rules (e.g., combine consecutive
FILTERs). -
Physical Plan Creation: Translates logical operators into physical operators (MapReduce stages).
-
MapReduce Job Generation: Compiles physical plan into one or more MapReduce jobs.
-
Execution: Jobs are submitted to Hadoop cluster and executed sequentially.
-
R Programming Language
An open-source language and environment for statistical computing and graphics.
-
Features & Advantages:
-
Powerful data handling & storage capabilities.
-
Vast collection of packages (CRAN) for statistical techniques, machine learning, visualization.
-
Excellent graphics capabilities (base, ggplot2).
-
Platform-independent.
-
Strong community support.
-
Integrates with Hadoop (RHadoop, SparkR).
-
-
Major Components of R Environment:
-
R Console: Interactive command-line interface.
-
R Script: Text file with R commands.
-
R Workspace: Current objects (variables, functions) in memory.
-
R Packages: Bundles of code, data, documentation (e.g.,
dplyr,ggplot2,caret). -
R Help System:
?function,help(package="pkg"). -
RStudio: Popular IDE with script editor, console, workspace viewer, plots pane.
-
-
Operations on Vectors (R's fundamental data structure):
-
Creation:
v <- c(1, 2, 3)orv <- 1:5. -
Indexing:
v[1](first),v[-2](exclude second),v[c(1,3)](specific),v[v > 2](logical). -
Arithmetic: Element-wise operations:
v * 2,v1 + v2(recycling rule applies if lengths differ). -
Recycling Rule: When operating on vectors of unequal length, the shorter vector is repeated until it matches the longer one. Warning if longer length isn't a multiple of shorter.
-
Functions:
length(v),sum(v),mean(v),sort(v),v[order(v)]. -
Vectorized Operations: Most functions operate on entire vectors without explicit loops (
sqrt(v),log(v)).
-
[!TIP] Exam Focus: For R, be prepared to write simple code snippets for vector creation, indexing (especially logical and negative indexing), and arithmetic operations demonstrating recycling. Know the purpose of key packages (
dplyrfor manipulation,ggplot2for graphics).