Big Data and AnalyticsUnit 75 min read

Data Ingestion & Stream Processing: Pipelines, Tools & Real-Time Analytics

Unit 7 of Big Data and Analytics explores how to ingest massive data streams (batch vs. real-time), key tools like Kafka and Flume, stream processing frameworks (Spark Streaming, Flink), and their applications in fraud detection, IoT, and social media analytics—with Nepalese examples like Ncell’s network monitoring and

What is Data Ingestion?

Data ingestion is the process of collecting, transferring, and loading raw data from diverse sources into a storage system (e.g., HDFS, data lakes) for processing or analysis. It is the first critical step in big data pipelines, ensuring data is available for analytics or machine learning.

Types of Data Ingestion

  1. Batch Processing: Data is ingested in fixed intervals (e.g., hourly/daily). Examples: Log files, nightly database backups.
  2. Real-Time/Stream Processing: Data is ingested continuously as it is generated. Examples: Stock market ticks, social media feeds, IoT sensor data.
flowchart TD
    A["Data Sources"] --> B["Batch Ingestion"]
    A --> C["Stream Ingestion"]
    B --> D["HDFS/S3\n(Stored for later analysis)"]
    C --> E["Stream Processor\n(Spark/Flink)"]
    E --> F["Real-Time Analytics\nor Storage"]

Why does this matter?

  • Batch: Suitable for historical analysis (e.g., monthly sales reports).
  • Stream: Critical for time-sensitive decisions (e.g., fraud detection in bank transactions).

Key Data Ingestion Tools

1. Apache Kafka

  • A distributed event streaming platform designed for high-throughput, fault-tolerant data pipelines.
  • Uses a publish-subscribe model where producers send data to topics, and consumers read from them.
  • Example: Ncell uses Kafka to ingest real-time call detail records (CDRs) to monitor network traffic and detect anomalies.
classDiagram
    class Producer {
        +send(data)
    }
    class Topic {
        +partitioned logs
    }
    class Consumer {
        +consume(data)
    }
    Producer --> Topic : publishes
    Topic --> Consumer : subscribes

2. Apache Flume

  • A reliable, scalable tool for batch data ingestion (e.g., logs, clickstreams).
  • Uses agents to collect data from sources (e.g., web servers) and channel it to HDFS or HBase.
  • Example: Daraz uses Flume to ingest user clickstream data from its website to analyze shopping trends.

3. Apache NiFi

  • A data flow automation tool for ETL (Extract, Transform, Load) pipelines.
  • Provides a visual interface to design data flows (e.g., moving data from APIs to databases).
  • Example: Pathao uses NiFi to aggregate ride data from drivers and passengers for dynamic pricing.

Stream Processing Frameworks

Stream processing analyzes data in motion, enabling real-time insights. Key frameworks:

Framework Language Key Feature Use Case
Apache Spark Streaming Scala/Python Micro-batch processing (DStreams) Fraud detection in bank transactions
Apache Flink Java/Scala Low-latency, stateful processing Real-time IoT sensor analytics
Apache Storm Java Fault-tolerant, distributed Social media trend analysis

Worked Example: Fraud Detection in NMB Bank

Scenario: NMB Bank wants to detect real-time fraud in credit card transactions.

  1. Data Source: Transactions streamed via Kafka topics.
  2. Processing: Spark Streaming analyzes transactions for anomalies (e.g., sudden large purchases).
  3. Action: Flag suspicious transactions for review.
# Pseudocode for fraud detection in Spark Streaming
from pyspark.streaming import StreamingContext

ssc = StreamingContext(sparkContext, batchDuration=10)  # 10-second batches
transactions = ssc.socketTextStream("kafka-broker", 9999)
fraudulent = transactions.filter(lambda x: is_fraud(x))  # Custom fraud logic
fraudulent.pprint()

In the Real World

  1. Ncell’s Network Monitoring

    • Tool: Kafka + Spark Streaming
    • How: Ingests CDR data (call logs) in real-time to detect SIM box fraud (illegal call routing). Alerts are sent to network ops within seconds.
  2. Daraz’s Order Processing

    • Tool: Flume + Hadoop
    • How: Batch-ingests order data nightly to update inventory and predict stockouts. Uses stream processing for real-time order status updates.
  3. NTC’s Traffic Management

    • Tool: Apache Flink
    • How: Processes GPS data from buses to optimize routes and reduce congestion in Kathmandu. Detects accidents or delays instantly.

Challenges in Data Ingestion

Challenge Solution Example
Data Volume Distributed tools (Kafka, Flume) Ncell handles millions of CDRs/hour
Latency Stream processing (Flink, Spark) Pathao updates prices in <100ms
Data Variety Schema-less ingestion (Avro, JSON) Daraz processes unstructured logs
Fault Tolerance Replication (Kafka partitions) NTC’s traffic system never crashes

Exam Tip

  1. Differentiate batch vs. stream: Batch is for historical analysis; stream is for real-time decisions.
  2. Tools matter: Know when to use Kafka (streaming), Flume (batch), or Spark/Flink (processing).
  3. Real-world mapping: Relate Kafka to Ncell/CDRs, Flume to Daraz logs, and Spark to fraud detection.
  4. Diagrams: Draw Kafka’s producer-consumer model or Spark Streaming’s DStreams in exams.
  5. Code snippets: Be ready to write pseudocode for stream processing (e.g., fraud detection logic).

Based on the TU BITM syllabus for Big Data and Analytics (IT278), unit 7.

Discussion

Loading…