Skip to content
CE-803 (B) · Data Analytics/Quick Revision Short Notes

Data Analytics (CE-803 (B)) - Unit 5 Short Notes

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:

  1. Input Splitting: Input data divided into logical splits (InputFormat).

  2. Map Phase: Parallel execution of map tasks; output stored in memory buffer.

  3. Shuffle & Sort: Transfer map outputs to reducers; sort by key.

  4. Reduce Phase: Aggregate sorted data; write to HDFS.

  5. 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:

  1. CREATE TABLE:

    
    CREATE TABLE employees (id INT, name STRING) 
    
    ROW FORMAT DELIMITED 
    
    FIELDS TERMINATED BY ',' 
    
    STORED AS TEXTFILE;
    
    
  2. ALTER TABLE (add column):

    
    ALTER TABLE employees ADD COLUMNS (dept STRING);
    
    
  3. 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:

  1. $match: Filter documents (like WHERE).

  2. $group: Group by field, compute aggregates ($sum, $avg).

  3. $project: Reshape documents (include/exclude fields, computed fields).

  4. $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:

  1. Ingestion: Sqoop/Flume → HDFS.

  2. Storage: HDFS (raw), Hive (structured), HBase (NoSQL).

  3. Processing: MapReduce/Spark → transformed data.

  4. 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).

DiagramCANVAS: Hadoop Cluster Architecture showing Master Node (NameNode, ResourceManager, Secondary NameNode) and Worker Nodes (DataNode, NodeManager) with HDFS and YARN layers
DiagramCANVAS: Hive Architecture with CLI/Beeline, HiveServer2, Driver, Metastore (RDBMS), Execution Engine (MapReduce/Tez), HDFS
Go to where you left off?

Quick Add to Notes

Save questions, your own notes and screenshots into notes filed by unit. It takes a free account.

Create free account

Have an account? Log in