CMP422 Data Science and Analytics

Data Science and AnalyticsUnit 918 min read

Big Data Tools: Hadoop, Spark, NoSQL, Cloud & Stream Processing

Unit 9 of Data Science and Analytics explores the core tools and architectures for handling big data—Hadoop (HDFS, MapReduce), Spark (RDDs, DataFrames), NoSQL databases (MongoDB, Cassandra), cloud platforms (AWS, GCP), and stream processing (Kafka, Flink)—with real-world examples from Nepalese and global tech ecosystem

TAKEAWAYS:

  • Hadoop’s HDFS and MapReduce split and parallelize data storage and processing across clusters, making them ideal for batch analytics (e.g., NTC’s network traffic logs).
  • Spark’s in-memory processing (RDDs/DataFrames) accelerates iterative algorithms like machine learning, used by banks for fraud detection (e.g., Nabil Bank’s real-time transaction monitoring).
  • NoSQL databases (MongoDB, Cassandra) store unstructured/semi-structured data (e.g., Daraz’s user reviews or Pathao’s ride GPS coordinates) with flexible schemas.
  • Cloud tools (AWS EMR, GCP Dataproc) provide scalable, pay-as-you-go big data pipelines (e.g., NEPSE’s stock market analytics).
  • Stream processing (Kafka, Flink) handles real-time data (e.g., Khalti’s transaction streams or Ncell’s call detail records) for instant insights.
  • Tool selection depends on data type (batch vs. stream), velocity, and use case (e.g., Hadoop for historical logs, Spark for ML, Kafka for live feeds).

1. Why Big Data Tools? The Scale Challenge

Big data isn’t just "more data"—it’s data that exceeds traditional tools’ capacity in volume (terabytes+), velocity (real-time streams), or variety (text, images, logs). For example:

  • NTC’s network logs: Millions of call records daily → need distributed storage (HDFS) and parallel processing (MapReduce).
  • Daraz’s user reviews: Unstructured text → require NoSQL (MongoDB) for flexible schema.
  • Pathao’s ride requests: Real-time GPS streams → need Kafka/Flink for low-latency processing.

Key metrics defining "big data":

Metric Traditional DB Limit Big Data Tools Handle
Volume GBs PBs/TBs
Velocity Batch (hours/days) Milliseconds (streams)
Variety Structured (SQL) Unstructured (JSON, logs)
Veracity Clean data Noisy/missing data

2. Hadoop: The Batch Processing Workhorse

Hadoop is an open-source framework designed for distributed storage and processing of large datasets across clusters. It consists of two core components:

A. HDFS (Hadoop Distributed File System)

  • Purpose: Stores data across multiple machines (nodes) in a fault-tolerant way.
  • How it works:
    • Data is split into blocks (default: 128MB or 256MB).
    • Blocks are replicated (default: 3 copies) across nodes to prevent loss.
    • NameNode manages metadata (file locations), while DataNodes store actual data.
graph TD
    A["HDFS Cluster"] --> B["NameNode\n(Manages metadata)"]
    A --> C["DataNode 1\n(Stores Block 1, Block 2)"]
    A --> D["DataNode 2\n(Stores Block 1, Block 3)"]
    A --> E["DataNode 3\n(Stores Block 2, Block 3)"]
    B -->|"Metadata"| C
    B -->|"Metadata"| D
    B -->|"Metadata"| E
  • Advantages:
    • Scalable: Add more DataNodes to increase storage.
    • Fault-tolerant: Losing a node doesn’t lose data (replication).
    • Cost-effective: Runs on commodity hardware.
  • Disadvantages:
    • Not for real-time: Batch processing only (minutes/hours).
    • High latency: Not suitable for interactive queries.
    • Complex setup: Requires expertise in cluster management.

Real-world use in Nepal:

  • NTC uses Hadoop to analyze call detail records (CDRs) for network optimization.
  • Nepal Rastra Bank processes historical financial transactions for fraud detection.

B. MapReduce: Parallel Processing Model

  • Purpose: Processes data in parallel across a cluster.
  • How it works:
    1. Input: Data is split into chunks and distributed across nodes.
    2. Map Phase: Each node processes its chunk independently (e.g., count word frequencies in a document).
    3. Shuffle: Intermediate results are sorted and sent to reducers.
    4. Reduce Phase: Aggregates results (e.g., sums all word counts).
sequenceDiagram
    participant User
    participant HDFS
    participant Mapper1
    participant Mapper2
    participant Reducer
    User->>HDFS: Split data into chunks
    HDFS->>Mapper1: Send Chunk 1
    HDFS->>Mapper2: Send Chunk 2
    Mapper1->>Reducer: <key1, value1>
    Mapper2->>Reducer: <key1, value2>
    Reducer->>User: Aggregated result

Worked Example: Word Count in Hadoop Problem: Count how often each word appears in a corpus of Nepali news articles (stored in HDFS). Steps:

  1. Map Phase:
    • Input: "नेपालको अर्थतन्त्र धेरै चुनौतीपूर्ण छ" → Split into words.
    • Output: (नेपालको, 1), (अर्थतन्त्र, 1), (धेरै, 1), (चुनौतीपूर्ण, 1), (छ, 1).
  2. Shuffle: Group by word.
  3. Reduce Phase:
    • Input: (नेपालको, [1]), (अर्थतन्त्र, [1]) → Output: (नेपालको, 1), (अर्थतन्त्र, 1).

Why this matters:

  • Nepal News Agency could use this to analyze trending topics across articles.
  • NEPSE uses similar batch processing to analyze historical stock trends.

3. Apache Spark: Faster In-Memory Processing

Hadoop’s MapReduce is slow for iterative algorithms (e.g., machine learning). Spark improves this by:

  • In-memory processing: Keeps data in RAM for faster access.
  • Lazy evaluation: Only computes when needed (optimizes performance).
  • Rich APIs: Supports SQL, Python (PySpark), R, and Scala.

Key Components of Spark

Component Purpose
RDD Resilient Distributed Dataset: Immutable, partitioned collection of data.
DataFrame Optimized for structured data (like a SQL table).
Spark SQL Runs SQL queries on Spark data.
MLlib Machine learning library (e.g., clustering, classification).
GraphX Graph processing (e.g., social networks).

spark architecture labelled diagram**Spark’s components and their interactions. (Image: lparant, CC BY-SA 4.0, via Wikimedia Commons)

Advantages over Hadoop MapReduce:

Feature Hadoop MapReduce Apache Spark
Processing Speed Disk-based (slow) In-memory (100x faster)
Latency Minutes/hours Seconds/milliseconds
Iterative Algorithms Poor support Optimized (e.g., ML)
Ease of Use Complex (Java) APIs in Python, R, SQL

Real-world use in Nepal:

  • Nabil Bank uses Spark for real-time fraud detection in transactions (e.g., sudden large withdrawals).
  • Daraz processes user clickstreams to personalize recommendations (e.g., "Customers who bought X also bought Y").

Worked Example: Fraud Detection with Spark (Nabil Bank Scenario)

Problem: Detect unusual transactions in real-time (e.g., a customer suddenly transferring ₹50,000 to an unknown account). Steps:

  1. Data Ingestion: Stream transactions via Kafka.
  2. Spark Processing:
    • Use Spark SQL to query:
      SELECT user_id, SUM(amount) as total_spent
      FROM transactions
      WHERE timestamp > NOW() - INTERVAL 1 HOUR
      GROUP BY user_id
      HAVING SUM(amount) > 100000  -- Flag high amounts
      
    • Apply MLlib’s clustering to identify anomalies (e.g., transactions far from the user’s normal pattern).
  3. Alert: Trigger an alert for manual review.

Why Spark?

  • Speed: Processes transactions in milliseconds (vs. Hadoop’s hours).
  • Scalability: Handles thousands of transactions per second.

4. NoSQL Databases: Flexible Storage for Unstructured Data

Traditional SQL databases (e.g., MySQL) struggle with unstructured data (e.g., JSON, logs, GPS coordinates). NoSQL databases solve this with:

  • Schema-less design: No fixed columns (e.g., MongoDB stores documents like { "user": "Alice", "purchases": [...] }).
  • Horizontal scaling: Add more servers to handle growth.
  • Specialized data models: Key-value, document, column-family, or graph.

Types of NoSQL Databases

Type Example Use Case Nepalese Example
Document MongoDB User profiles, product catalogs Daraz’s product listings
Key-Value Redis Caching, session storage eSewa’s user login sessions
Column-Family Cassandra Time-series data (e.g., IoT, logs) NTC’s network traffic logs
Graph Neo4j Social networks, fraud detection Pathao’s ride-sharing network

Worked Example: Daraz’s Product Catalog (MongoDB) Problem: Store product details with flexible attributes (e.g., some products have "size" but not "weight"). Solution: Use MongoDB documents:

{
  "_id": "prod_123",
  "name": "Nepali Woolen Blanket",
  "price": 2500,
  "attributes": {
    "color": ["Red", "Blue"],
    "size": ["S", "M", "L"],
    "weight": null  // Optional field
  },
  "reviews": [
    { "user": "user456", "rating": 5, "comment": "Very warm!" }
  ]
}

Advantages:

  • Flexibility: Add new fields without schema changes.
  • Scalability: Daraz can add more MongoDB servers as product listings grow.

5. Cloud Platforms for Big Data: AWS, GCP, Azure

Managing Hadoop/Spark clusters on-premises is expensive. Cloud providers offer managed services:

Service Provider Purpose Nepalese Use Case
EMR AWS Managed Hadoop/Spark clusters NEPSE’s historical stock analysis
Dataproc GCP Serverless Spark/Hadoop Nepal Rastra Bank’s financial reports
Azure HDInsight Azure Hadoop/Spark on Azure NTC’s big data analytics
BigQuery GCP Serverless SQL for big data Kathmandu Metropolitan City’s traffic analysis

Why Cloud?

  • No infrastructure management: Pay for what you use.
  • Auto-scaling: Handle traffic spikes (e.g., Daraz during sales).
  • Integration: Connects to other cloud services (e.g., AWS S3 for storage).

Worked Example: NEPSE’s Stock Market Analytics (AWS EMR) Problem: Analyze 10+ years of stock price data to predict trends. Steps:

  1. Store data: Upload CSV files to AWS S3.
  2. Process: Use EMR to run Spark jobs on historical data.
  3. Visualize: Export results to QuickSight for dashboards.

Cost: Pay only for the hours the cluster runs (vs. buying physical servers).


Batch processing (Hadoop/Spark) is too slow for real-time needs. Stream processing tools handle data as it arrives:

Key Tools

Tool Purpose Nepalese Example
Apache Kafka Distributed event streaming platform Khalti’s transaction logs
Apache Flink Low-latency stream processing Ncell’s real-time call analytics
Spark Streaming Micro-batch processing (Spark + streams) Pathao’s ride demand forecasting

How Kafka Works:

  1. Producers (e.g., Khalti app) send transactions to a topic (e.g., transactions).
  2. Brokers store data in partitions (scalable storage).
  3. Consumers (e.g., Flink) process data in real-time.
graph TD
    A["Khalti App\n(Producer)"] -->|"Transaction Data"| B["Kafka Broker 1\n(Topic: transactions)"]
    A --> C["Kafka Broker 2"]
    B --> D["Flink Consumer\n(Real-time Fraud Check)"]
    C --> D

Worked Example: Khalti’s Fraud Detection Problem: Flag suspicious transactions (e.g., ₹10,000 sent to a new merchant in 1 second). Steps:

  1. Kafka Topic: transactions receives every payment.
  2. Flink Job:
    • Check if amount > 5000 and merchant_id is new.
    • If yes, send alert to fraud-alerts topic.
  3. Action: Khalti blocks the transaction and notifies the user.

Why Stream Processing?

  • Real-time decisions: Block fraud before it’s too late.
  • Scalability: Handle millions of transactions per second (e.g., during Dashain sales).

7. Choosing the Right Tool: Decision Guide

Not all tools fit every job. Use this table to decide:

Requirement Hadoop (HDFS + MapReduce) Spark NoSQL (MongoDB/Cassandra) Kafka/Flink Cloud (EMR/Dataproc)
Data Type Structured (batch) Structured/semi-structured Unstructured (JSON, logs) Streams (events) All (depends on service)
Processing Speed Slow (hours) Fast (seconds) Fast (read/write) Real-time (ms) Depends on config
Use Case Historical analytics ML, iterative algorithms Flexible schemas, high write throughput Real-time dashboards, fraud detection Managed big data pipelines
Example in Nepal NTC’s network logs Nabil Bank’s fraud detection Daraz’s product catalog Khalti’s transactions NEPSE’s stock analysis
Learning Curve High (Java, cluster setup) Medium (Python/SQL) Low (document-based) Medium (streaming logic) Low (managed service)

8. Big Data Tools in Action: Case Studies

Case 1: Daraz’s Recommendation Engine (Spark + MongoDB)

  • Problem: Personalize product recommendations for 10M+ users.
  • Solution:
    • Spark MLlib: Trains a collaborative filtering model on user purchase history.
    • MongoDB: Stores user profiles and clickstream data flexibly.
    • Result: 20% increase in conversion rates.
  • Problem: Predict surge pricing during traffic jams (e.g., Kathmandu’s Ring Road).
  • Solution:
    • Kafka: Streams real-time GPS data from drivers.
    • Flink: Aggregates demand per 500m grid cell every 30 seconds.
    • Result: Dynamic pricing adjusts in real-time.

Case 3: NTC’s Network Optimization (Hadoop + Cassandra)

  • Problem: Analyze 10TB of call logs to reduce dropped calls.
  • Solution:
    • HDFS: Stores raw CDRs.
    • MapReduce: Identifies peak congestion times.
    • Cassandra: Stores time-series data for fast queries.
    • Result: 15% reduction in call drops during festivals.

In the Real World

  1. Khalti’s Transaction Processing (Kafka + Flink)

    • Tool: Apache Kafka ingests every transaction (₹500M+ daily), while Flink detects fraud in <100ms.
    • How: Flink joins transaction data with user profiles to flag anomalies (e.g., sudden large transfers).
    • Impact: Blocks 90% of fraudulent transactions before they complete.
  2. Daraz’s Inventory Management (Spark + MongoDB)

    • Tool: Spark processes sales data to forecast demand, while MongoDB stores product attributes (e.g., size, color).
    • How: Uses Spark’s MLlib to predict stockouts and MongoDB’s flexible schema to handle new product types (e.g., adding "material" field for blankets).
    • Impact: Reduces overstocking by 30% during Dashain sales.
  3. NEPSE’s Stock Market Analytics (AWS EMR)

    • Tool: AWS EMR runs Spark jobs on 20+ years of stock price data.
    • How: Analyzes trends like "stocks rise 10% before Dashain" to help traders.
    • Impact: Used by 80% of brokerage firms in Nepal.
  4. Pathao’s Driver Matching (Cassandra + Spark)

    • Tool: Cassandra stores driver locations (updated every 5 seconds), while Spark matches riders to nearest drivers.
    • How: Spark queries Cassandra to find the closest available driver within 2km.
    • Impact: Reduces wait times from 5 minutes to <90 seconds.

Exam Tip

This unit is conceptual + applied. Expect:

  1. Short Definitions (2 marks):

    • "Explain HDFS replication."
    • "What is an RDD in Spark?"
    • Tip: Memorize key terms like NameNode, DataNode, MapReduce phases, Kafka topics/partitions.
  2. Scenario-Based Questions (5–7 marks):

    • "NTC wants to analyze 50TB of call logs. Compare Hadoop vs. Spark for this task."
    • "Daraz stores user reviews in MongoDB. Why not MySQL?"
    • Tip: Use the decision table above to structure answers. Always tie tools to real-world Nepalese examples.
  3. Diagram Questions (3–5 marks):

    • Draw HDFS architecture or Kafka producer-consumer flow.
    • Tip: Sketch the NameNode/DataNode or Kafka broker diagrams from this note.
  4. Code Snippets (3 marks):

    • Write a MapReduce word count pseudocode or a Spark SQL query.
    • Tip: For Spark SQL, use the Nabil Bank fraud detection example.
  5. Comparison Tables (5 marks):

    • "Compare Hadoop and Spark for a machine learning task."
    • Tip: Use the Hadoop vs. Spark table in this note as a template.

Common Pitfalls to Avoid:

  • Confusing HDFS (storage) with MapReduce (processing).
  • Forgetting that NoSQL is for unstructured data (e.g., don’t use MongoDB for financial ledgers).
  • Overlooking real-time vs. batch: Kafka/Flink for streams, Hadoop/Spark for batch.
  • Ignoring cloud options: AWS EMR/GCP Dataproc are valid answers for "scalable big data" questions.

Final Advice:

  • Relate everything to Nepal: Examiners love answers like "NTC uses Hadoop for CDRs" or "Daraz uses Spark for recommendations."
  • Practice drawing:
    • HDFS cluster (NameNode + DataNodes).
    • Kafka producer-consumer flow.
    • Spark RDD/DataFrame lifecycle.
  • Memorize 3 tools per category:
    • Batch: Hadoop, Spark.
    • NoSQL: MongoDB, Cassandra.
    • Streaming: Kafka, Flink.

Based on the PU BE Computer (PU) syllabus for Data Science and Analytics (CMP422), unit 9.

Discussion

Loading…