UNIT 5: BIG DATA ANALYTICS & HADOOP ECOSYSTEM
1. Big Data Fundamentals
Definition: Extremely large, complex datasets that traditional data processing tools cannot handle effectively. Evolved from relational databases to distributed systems due to internet-scale data growth.
Characteristics (4 V’s):
-
Volume: Massive scale (TB/PB/EB). Example: Social media logs.
-
Velocity: High-speed generation/processing. Example: Real-time sensor data.
-
Variety: Structured, semi-structured, unstructured. Example: Text, images, videos.
-
Veracity: Uncertainty/quality issues (inconsistency, missing data).
[!TIP] Exam often asks for examples for each V. Link to smart cities: traffic cameras (velocity), social media feeds (variety), energy meter readings (volume).
Applications in Smart Cities:
-
Intelligent Traffic Management: Real-time analysis of GPS/traffic camera data to optimize signals.
-
Energy Grid Optimization: Predictive analytics on smart meter data for load balancing.
-
Public Safety: Crime pattern analysis using historical police records + social media.
-
Waste Management: Sensor-based bin fill-level monitoring for efficient collection routes.
2. Hadoop Distributed File System (HDFS)
Goals: Store huge files reliably across clusters of commodity hardware; provide high throughput access.
Architecture & Components:
-
NameNode (Master): Manages file system namespace, metadata (file permissions, block locations). Single point of failure.
-
DataNode (Slave): Stores actual data blocks (default 128MB), handles read/write requests, reports block status.
-
Secondary NameNode: Not a backup! Periodically merges fsimage (metadata snapshot) and edit logs to prevent edit log overgrowth. Does not provide failover.
Data Replication & Fault Tolerance:
-
Default replication factor = 3 (3 copies).
-
Blocks stored on different DataNodes; rack-aware placement (2 on same rack, 1 on different rack).
-
If DataNode fails, NameNode replicates blocks from remaining replicas to maintain replication factor.
[!TIP] Common Pitfall: Secondary NameNode ≠ Backup NameNode. It only aids checkpoints.
3. MapReduce Programming Model
Core Concepts:
-
Map Function: Processes input key-value pairs → intermediate key-value pairs.
map(key, value): for each word w in value: emit(w, 1) -
Reduce Function: Aggregates values for each intermediate key.
reduce(key, list<values>): sum = 0 for v in values: sum += v emit(key, sum)
Data Flow & Execution Phases:
-
Input Splitting: Input data divided into logical splits (InputFormat).
-
Map Phase: Parallel execution of map tasks; output stored in memory buffer.
-
Shuffle & Sort: Transfer map outputs to reducers; sort by key.
-
Reduce Phase: Aggregate sorted data; write to HDFS.
-
Commit Phase: Final output committed.
Capacity Scheduler:
-
Purpose: Multi-tenant resource allocation (queues with capacity guarantees).
-
Features:
-
Queues allocated minimum capacity (guaranteed resources).
-
Uses FIFO within queues, but queues can have priorities.
-
Supports resource-based scheduling (memory, CPU).
-
Allows elastic usage: free queue capacity can be borrowed temporarily.
-
[!TIP] Distinguish from Fair Scheduler: Fair aims for equal share over time; Capacity enforces minimum guarantees.
4. Hadoop Ecosystem Architecture
Core Components:
-
HDFS: Storage layer.
-
MapReduce/YARN: Processing/resource management.
- YARN (Yet Another Resource Negotiator): Separates resource management (ResourceManager) from task scheduling (ApplicationMaster).
-
Common Utilities: Hadoop Common, Avro, etc.
Cluster Setup:
-
Master Node: NameNode, ResourceManager, Secondary NameNode.
-
Worker/Slave Nodes: DataNode, NodeManager.
-
Configuration via
core-site.xml,hdfs-site.xml,yarn-site.xml.
Data Ingestion Tools:
-
Sqoop: Relational DB ↔ Hadoop (structured data).
-
Flume: Streaming data ingestion (logs, events) into HDFS.
5. Apache Hive: Data Warehousing on Hadoop
Architecture:
User → CLI/Beeline → HiveServer2 (JDBC/ODBC) → Driver → Metastore (metadata) & Execution Engine → MapReduce/Tez/Spark → HDFS
-
HiveServer2: Supports concurrent client connections (multi-user).
-
Metastore: Central repository for schema info (table definitions, partitions, column types). Stores in RDBMS (MySQL/PostgreSQL) or embedded Derby. Critical for query planning.
-
Driver: Compiles, optimizes, executes queries; manages lifecycle.
HiveQL vs SQL:
| Feature | HiveQL | SQL (RDBMS) |
|---|---|---|
| Schema | Schema-on-read (flexible) | Schema-on-write (rigid) |
| Latency | High (batch processing) | Low (OLTP) |
| Updates | Limited (ACID in newer versions) | Full ACID |
| Scale | Petabyte-scale (distributed) | Terabyte-scale (single) |
Metastore Deep Dive:
-
Role: Stores metadata only (not data). Tables map to HDFS directories.
-
Storage: External RDBMS recommended for production.
-
Integration: Hive, Pig, Spark can share metastore.
HiveQL DDL Commands:
-
CREATE TABLE:
CREATE TABLE employees (id INT, name STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY ',' STORED AS TEXTFILE; -
ALTER TABLE (add column):
ALTER TABLE employees ADD COLUMNS (dept STRING); -
DROP TABLE:
DROP TABLE employees; -- Metadata + data deleted (if not EXTERNAL)
Indexing & Partitioning:
-
Partitioning: Divide tables by column values (e.g.,
PARTITIONED BY (year INT)). Improves query performance by pruning partitions. -
Indexing: Accelerate selective queries. Types: Compact (concise key storage), Bitmap (for low-cardinality columns). Note: Hive indexes less mature than RDBMS.
6. Apache Pig: Data Flow Scripting
Pig Latin: High-level scripting language for data transformation.
-
Data Model:
-
Atom: Single value (int, string).
-
Tuple: Ordered set of fields (like row).
-
Bag: Collection of tuples (like table).
-
-
Syntax Example:
A = LOAD 'data.txt' USING PigStorage(',') AS (id:int, name:chararray); B = FILTER A BY id > 100; C = GROUP B BY name; D = FOREACH C GENERATE group, COUNT(B); STORE D INTO 'output';
Why Pig? (vs Java MapReduce):
-
Abstraction: No Java coding; fewer lines.
-
Automatic Optimization: Pig’s optimizer reorders operations.
-
Flexibility: Schema-less processing; handles nested data.
-
Extensibility: User-defined functions (UDFs) in Java/Python.
Execution Modes:
-
Local Mode: Single JVM; for development/testing (uses local filesystem).
-
Hadoop (MapReduce) Mode: Distributed execution on cluster (default).
7. NoSQL Databases: MongoDB
Indexing Strategies:
-
Single Field:
db.collection.createIndex({field: 1}) -
Compound:
db.collection.createIndex({field1: 1, field2: -1})(order matters for sort). -
Text: For text search;
db.collection.createIndex({field: "text"}). -
Multikey: Automatically created for array fields.
Aggregation Framework:
Pipeline of stages:
-
$match: Filter documents (like WHERE).
-
$group: Group by field, compute aggregates (
$sum,$avg). -
$project: Reshape documents (include/exclude fields, computed fields).
-
$sort, $limit, $skip, etc.
db.sales.aggregate([
{ $match: { amount: { $gt: 100 } } },
{ $$\displaystyle group: { _id: " $$item", total: { $sum: "$amount" } } },
{ $sort: { total: -1 } }
])
Comparison with Relational Databases:
| Feature | MongoDB (NoSQL) | Relational (RDBMS) |
|---|---|---|
| Schema | Dynamic (flexible) | Fixed (rigid) |
| Scaling | Horizontal (sharding) | Vertical (scale-up) |
| Joins | Limited ($lookup stage) |
Native (SQL JOINs) |
| Transactions | Multi-document (since 4.0) | Full ACID |
| Use Case | Unstructured/semi-structured | Structured, complex joins |
8. Supporting Tools & Coordination: Apache ZooKeeper
Purpose: Centralized service for coordination, configuration, naming, synchronization in distributed systems.
Advantages:
-
Simple API: Znodes (data nodes) with
getData,setData,create,delete. -
High Availability: Replicated across ensemble (odd number of servers).
-
Ordered Updates: Sequential znodes for coordination (e.g., leader election).
-
Watch Mechanism: Clients watch znodes for changes (event-driven).
Role in Hadoop Ecosystem:
-
HBase: Manages region assignment, master election.
-
Hive: Stores metastore configuration.
-
Kafka: Manages broker metadata, topic configuration.
-
YARN: Stores application state (in some setups).
9. Analytics Concepts & Techniques
Term Frequency (TF) & Inverse Document Frequency (IDF):
- TF: Frequency of term t in document d.
$$TF(t,d) = \frac{\text{number of times } t \text{ appears in } d}{\text{total terms in } d}$$
- IDF: Measures term rarity across corpus D.
$$IDF(t,D) = \log \left( \frac{\text{total documents in } D}{\text{documents containing } t} \right)$$
- TF-IDF:
$$TF\text{-}IDF(t,d,D) = TF(t,d) \times IDF(t,D)$$
-
Higher value = term is frequent in doc but rare in corpus → discriminative.
-
Application: Text mining, information retrieval (search ranking), document clustering.
Regression vs ANOVA:
| Aspect | Regression | ANOVA |
|---|---|---|
| Purpose | Model relationship between dependent & independent variables (continuous). | Compare means across multiple groups (categorical independent variable). |
| Output | Predictive equation (line/curve). | F-statistic; significance of group differences. |
| Variables | Dependent (continuous), independent(s) (continuous/categorical). | Dependent (continuous), independent (categorical, ≥2 groups). |
| Use Case | Predict house price from size, location. | Test if crop yield differs across fertilizer types. |
| Assumptions | Linearity, normality, homoscedasticity. | Normality, homogeneity of variances, independence. |
[!TIP] Regression can handle multiple predictors; ANOVA typically one categorical factor (but MANOVA extends).
10. Information Management in Big Data
Data Lifecycle Management:
-
Ingestion: Sqoop/Flume → HDFS.
-
Storage: HDFS (raw), Hive (structured), HBase (NoSQL).
-
Processing: MapReduce/Spark → transformed data.
-
Archival: Move cold data to cheaper storage (e.g., HDFS archival, tape).
Metadata Management & Data Governance:
-
Metadata: Data about data (schema, lineage, ownership). Stored in Hive Metastore, Apache Atlas.
-
Governance: Policies for data quality, security, compliance (e.g., GDPR). Tools: Apache Ranger (access control), Atlas (lineage).
Security & Access Control:
-
Authentication: Kerberos (tickets).
-
Authorization: HDFS permissions (POSIX-like), Apache Ranger (fine-grained, policy-based).
-
Encryption: HDFS Transparent Encryption (zones), TLS for data in transit.
11. Applications & Case Studies
Smart Cities Development:
-
Traffic Flow Optimization: Real-time GPS data → predictive signal timing.
-
Energy Management: Smart grid analytics for demand forecasting.
-
Waste Management: IoT sensor data → dynamic collection routes.
-
Public Health: Disease outbreak prediction from social media + hospital records.
Sector-Specific Analytics:
-
Water Management: Sensor networks for leakage detection, quality monitoring (cross-topic with IWRM).
-
Structural Health Monitoring: Vibration/strain data from bridges → predictive maintenance.
[!TIP] Link Big Data 4 V’s to smart cities: Velocity (real-time traffic), Volume (city-wide sensor data), Variety (images, audio, text), Veracity (noisy sensor readings).
Quick Reference: High-Frequency Exam Formulas & Definitions
-
TF-IDF: $$\displaystyle \boxed{TF\text{-}IDF(t,d,D) = \frac{f_{t,d}}{\sum_{t' \in d} f_{t',d}} \times \log \frac{|D|}{|\{d \in D: t \in d\}|}} $$
-
HDFS Replication: Default block size = 128 MB (Hadoop 2.x+), replication factor = 3.
-
Hive Metastore: $$\displaystyle \boxed{\text{Central metadata repository storing schema, partitions, table locations}} $$.
-
Pig Latin Data Model: Atom < Tuple < Bag.
-
ZooKeeper Ensemble: Minimum 3 servers (odd number for quorum).
-
MapReduce Shuffle: Transfer of map output to reduce inputs (sorted by key).