UNIT 2: Big Data Analytics and Technologies
Big Data Fundamentals
Definition: Big Data refers to extremely large, complex, and rapidly growing datasets that traditional data processing tools cannot effectively capture, store, manage, and analyze.
Need for Big Data: To uncover hidden patterns, correlations, and insights from massive, diverse data sources for improved decision-making, innovation, and competitive advantage.
Characteristics of Big Data (Extended 7 V's):
| V | Description | Key Aspect |
|---|---|---|
| Volume | Scale of data (Terabytes to Zettabytes) | Size |
| Velocity | Speed of data generation & processing | Real-time/Streaming |
| Variety | Different types & sources (structured, unstructured, semi-structured) | Formats |
| Veracity | Uncertainty, quality, and trustworthiness of data | Accuracy & Reliability |
| Value | Extraction of meaningful insights & business value | Utility |
| Variability | Changing data formats, meanings, or flow rates | Inconsistency |
| Visualization | Graphical representation to understand complex data | Interpretation |
The Core 4 V's: Volume, Velocity, Variety, Veracity are the most fundamental and frequently cited characteristics.
Application: Big Data in Smart Cities
-
Traffic Management: Real-time analysis of traffic camera feeds and GPS data to optimize signal timing and reduce congestion.
-
Energy Grid Optimization: Analyzing consumption patterns from smart meters to predict demand and integrate renewable sources.
-
Public Safety: Predictive policing by analyzing crime reports, social media, and sensor data to deploy resources proactively.
-
Waste Management: Sensor-equipped bins signal when full, optimizing collection routes.
-
Example: Singapore's "Smart Nation" initiative uses integrated data from transport, utilities, and environment to improve urban planning and services.
Hadoop Ecosystem
Hadoop Architecture & Design Principles
-
Core Components: HDFS (Storage), MapReduce (Processing), YARN (Resource Management).
-
Design Principles:
-
Commodity Hardware: Runs on inexpensive, standard hardware.
-
Fault Tolerance: Automatically handles hardware failures by data replication and task re-execution.
-
Scalability: Scales horizontally by adding more nodes (linear scalability).
-
HDFS (Hadoop Distributed File System)
-
Goals/Objectives: Provide a distributed, fault-tolerant, scalable file system for very large files with high throughput access.
-
Architecture:
-
NameNode: Master server. Manages file system namespace, metadata (file permissions, block locations), and regulates access.
-
DataNode: Slave server. Stores actual data blocks, serves read/write requests, and performs block replication/deletion.
-
Secondary NameNode: Not a backup NameNode. Periodically merges the EditLog with FsImage to prevent EditLog from becoming too large. Critical for checkpointing.
-
-
Data Replication & Block Size:
-
Default Block Size: 128 MB (or 256 MB in newer versions). Files are split into blocks.
-
Default Replication Factor: 3. Each block is copied to 3 different DataNodes (1 on same rack, 2 on different racks for rack awareness).
-
Formula: Storage overhead = (Replication Factor - 1) * 100%. For RF=3, overhead is 200%.
-
MapReduce Programming Model
A programming model for processing large datasets in parallel across a cluster.
-
Map Phase:
-
Input:
<Key, Value>pairs (e.g.,<LineOffset, LineText>from HDFS). -
Process: User-defined
map()function processes each input pair and emits zero or more intermediate<Key, Value>pairs. -
Output: Intermediate pairs are partitioned by key and sorted locally.
-
-
Shuffle and Sort Phase:
-
Automatic framework phase.
-
All intermediate values for a given key are grouped together and sorted. This data is transferred from Mapper nodes to Reducer nodes.
-
-
Reduce Phase:
-
Input:
<Key, List<Value>>pairs from the Shuffle phase. -
Process: User-defined
reduce()function aggregates/processes the list of values for each key. -
Output: Final
<Key, Value>pairs written to HDFS.
-
YARN (Yet Another Resource Negotiator) & Capacity Scheduler
-
Purpose: Separples resource management from the MapReduce programming model, enabling multiple data processing engines (Spark, Tez) to run on Hadoop.
-
Architecture:
-
ResourceManager (RM): Global master. Arbitrates resources among all applications.
-
NodeManager (NM): Per-node slave. Manages containers and monitors resource usage.
-
ApplicationMaster (AM): Per-application master. Negotiates resources with RM and works with NM to execute tasks.
-
-
Capacity Scheduler:
-
A pluggable scheduler for YARN designed for multi-tenant clusters.
-
Features:
-
Queues: Organizes applications into queues (e.g.,
production,development). Each queue has a guaranteed capacity (min share) and can use excess capacity. -
Resource Elasticity: Queues can use free capacity from other queues if available.
-
Priority Scheduling: Within a queue, applications are scheduled by priority (FIFO by default).
-
Queue-based Access Control: ACLs control who can submit to which queue.
-
-
Configuration:
capacity-scheduler.xmldefines queues, capacities, max capacities, and ACLs.
-
Higher-Level Data Processing Tools
Apache Hive
-
Purpose: Data warehouse infrastructure built on Hadoop. Provides SQL-like querying (HiveQL) for summarization, ad-hoc querying, and analysis of large datasets.
-
Architecture & Components:
-
HiveServer2 (HS2): Enables JDBC/ODBC clients to execute queries. Supports multi-client concurrency and authentication.
-
Metastore: Central repository for all Hive metadata (table schema, partition info, location, SerDe info). Stores in a relational database (MySQL, PostgreSQL).
- Components: Thrift Server (exposes Metastore API), Database.
-
Driver: Compiles, optimizes, and executes HiveQL queries. Manages the lifecycle of a query.
-
Compiler: Converts HiveQL into a directed acyclic graph (DAG) of MapReduce/Tez/Spark tasks.
-
Execution Engine: Executes the tasks generated by the compiler (typically on YARN).
-
-
HiveQL DDL Commands (Syntax & Examples):
-- 1. CREATE TABLE CREATE TABLE IF NOT EXISTS employees ( emp_id INT, name STRING, dept STRING ) ROW FORMAT DELIMITED FIELDS TERMINATED BY ',' STORED AS TEXTFILE; -- 2. ALTER TABLE (Add Column) ALTER TABLE employees ADD COLUMNS (salary FLOAT); -- 3. DROP TABLE DROP TABLE IF EXISTS employees;
Apache Pig
-
Pig Latin: A high-level data flow language for expressing data transformations.
-
Need & Advantages over Raw MapReduce:
-
Abstraction: Operates at a higher level (data flow) vs. low-level key-value pairs.
-
Automatic Optimization: Pig's logical and physical optimizers combine and reorder operations (e.g.,
FILTERbeforeJOIN). -
Ease of Use: ~20 lines of Pig Latin often replace 200+ lines of Java MapReduce.
-
Flexibility: Handles both structured and semi-structured data naturally.
-
Extensibility: User Defined Functions (UDFs) in Java, Python, etc.
-
Coordination and Management Services
Apache ZooKeeper
-
Purpose: A centralized service for distributed coordination. Solves problems in distributed systems: configuration management, naming, synchronization, group services, and leader election.
-
Data Model: Hierarchical namespace (like a file system) of znodes. Znodes can be ephemeral (deleted on session end) or sequential (name appended with a sequence number).
-
Advantages:
-
Reliability: Runs on a cluster of servers (ensemble) with replication.
-
Simplicity: Simple API (
create,get,set,delete,watch). -
Performance: In-memory data storage for low latency.
-
Ordering Guarantees: Strictly ordered znode updates and watches enable building robust primitives (barriers, queues).
-
Data Storage and Management
NoSQL Databases: MongoDB
-
Indexing:
| Type | Description | Use Case | |----------|----------------|--------------| | Single Field | Index on a single field. | Basic equality/range queries. | | Compound | Index on multiple fields. Order matters. | Queries on multiple fields. | | Multikey | Index on array fields. One index entry per array element. | Querying array contents. | | Text | Index on string content for text search. | Full-text search (supports language-specific stemming). | | Hashed | Index using a hash of the field's value. | Equality matches only (hash collisions possible). |
- Benefits: Dramatically improves query performance, enables efficient sorting, supports unique constraints.
-
Aggregation Framework:
-
A pipeline-based framework for data transformation and analysis.
-
Pipeline Stages: Data flows through stages (
$match,$group,$sort,$project,$limit,$skip,$unwind,$lookup). -
Use Case Example: Find total sales per category for products in a specific region.
db.sales.aggregate([ { $match: { region: "West" } }, // Filter { $$\displaystyle group: { _id: " $$category", total_sales: { $sum: "$amount" } } }, // Group & Sum { $sort: { total_sales: -1 } } // Sort descending ])
-
Information Management in Big Data Context
-
Metadata Management: "Data about data." Crucial for understanding data sources, schemas, lineage, and quality. Tools: Apache Atlas, Hive Metastore.
-
Data Governance: Policies, standards, and processes to ensure high data quality, security, privacy, and compliance.
-
Data Lineage: Tracks the lifecycle of data from origin through transformations to final consumption. Essential for debugging, impact analysis, and compliance (e.g., GDPR).
Text Analytics Concepts
TF-IDF (Term Frequency-Inverse Document Frequency)
-
Definition: A statistical measure to evaluate the importance of a word in a document relative to a collection (corpus).
-
Mathematical Formula:
$$\text{TF-IDF}(t, d, D) = \text{TF}(t, d) \times \text{IDF}(t, D)$$
where:
* $t$ = term (word)
* $d$ = document
* $D$ = corpus (collection of documents)
-
Calculation:
- Term Frequency (TF): Frequency of term $t$ in document $d$.
$$\text{TF}(t, d) = \frac{\text{frequency of } t \text{ in } d}{\text{total number of terms in } d}$$
* **Inverse Document Frequency (IDF):** Measures how common the term is across all documents. Downweights frequent words (like "the", "is").
$$\text{IDF}(t, D) = \log \left( \frac{\text{total number of documents in } D}{\text{number of documents containing } t} \right)$$
- Applications: Information retrieval (search engine ranking), text mining, keyword extraction, document similarity.
Statistical Analysis for Data Analytics
Regression Analysis
-
Purpose: Predictive modeling to estimate the relationship between a dependent variable and one or more independent variables.
-
Types:
-
Linear Regression: Models linear relationship. $$\displaystyle Y = \beta_0 + \beta_1 X + \epsilon $$. (Continuous Y)
-
Logistic Regression: Models probability of a binary outcome. $$\displaystyle p = \frac{1}{1 + e^{-(\beta_0 + \beta_1 X)}} $$. (Categorical Y)
-
-
Key Assumptions (Linear): Linearity, Independence, Homoscedasticity, Normality of residuals, No multicollinearity.
-
Interpretation: Coefficients ($\beta$) indicate the change in Y for a one-unit change in X, holding others constant. $$\displaystyle R^2 $$ measures proportion of variance explained.
ANOVA (Analysis of Variance)
-
Purpose: Compare means across three or more groups to determine if at least one group mean is statistically different.
-
Types:
-
One-way ANOVA: One independent categorical variable (factor), one continuous dependent variable.
-
Two-way ANOVA: Two independent categorical variables. Can test for interaction effects.
-
-
Core Mechanism: Partitions total variance into:
-
Between-group variance (due to factor differences)
-
Within-group variance (random error)
-
-
F-test: $$\displaystyle F = \frac{\text{Mean Square Between (MSB)}}{\text{Mean Square Within (MSW)} $$. Large F-value suggests group means differ significantly.
-
Assumptions: Normality, Homogeneity of variances, Independence of observations.
Comparison: Regression vs. ANOVA
| Aspect | Regression | ANOVA |
|---|---|---|
| Objective | Predict a continuous outcome / Model relationships | Compare group means |
| Dependent Variable | Continuous | Continuous |
| Independent Variables | Continuous or Categorical (can be mixed) | Categorical (Factors) |
| Primary Output | Regression equation, coefficients, $$\displaystyle R^2 $$ | F-statistic, p-value, group means |
| Use Case | "How does salary change with experience and education?" | "Do test scores differ between teaching methods A, B, and C?" |
| Relationship | Can include categorical predictors (then tests similar to ANOVA). ANOVA is a special case of linear regression with only categorical predictors. |