Big Data and AnalyticsUnit 715 min read

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

Unit 7 of Big Data and Analytics explores how raw data is collected, processed in real-time, and stored for analytics. Learn about batch vs. stream processing, ingestion tools (Kafka, Flume), stream processing frameworks (Spark Streaming, Flink), and real-world applications like fraud detection in eSewa or traffic rout

Key Concepts & Workflow

1. Data Ingestion: The Gateway to Big Data

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

Types of Data Ingestion

Type Description Example Use Case Tools Used
Batch Processing Data collected and processed in fixed intervals (hours/days). Daily sales reports for Daraz. Apache NiFi, Sqoop
Stream Processing Real-time data processed as it arrives (milliseconds/seconds). Fraud detection in Khalti transactions. Apache Kafka, Flume
Micro-batch Hybrid approach: small batches processed frequently (e.g., every 5 seconds). Stock price updates on NEPSE. Spark Streaming

Why does ingestion matter?

  • Volume: Handles terabytes/petabytes from IoT sensors (e.g., NTC’s smart meters), logs (e.g., Daraz server logs), or social media (e.g., Twitter feeds).
  • Velocity: Stream processing enables real-time decisions (e.g., Pathao’s dynamic pricing based on live traffic).
  • Variety: Ingests structured (SQL databases), semi-structured (JSON/XML), and unstructured (images, videos) data.

2. Data Ingestion Tools: How Data Moves

२०७१ (Kafka 0.8)Kafka introducespartitions for scalabi२०७२ (Flume 1.6)Flume adds HDFSsink (Ncell CDRs store२०७४ (Spark 2.0)Spark Streamingintegrates with Kafka
Key milestones in Nepali case studies (timeline of tool adoption in BS years)

A. Apache Kafka

Kafka is a distributed event streaming platform designed for high-throughput, fault-tolerant ingestion. It acts as a buffer between data producers (e.g., mobile apps, sensors) and consumers (e.g., Spark, Flink).

How Kafka Works:

  1. Producers (e.g., eSewa app) publish data to topics (logical channels like transactions or user_activity).
  2. Brokers (Kafka servers) store data in partitions (scalable storage units).
  3. Consumers (e.g., fraud detection system) subscribe to topics and process data in real-time.

Example: eSewa Fraud Detection

  • Producer: eSewa app logs every transaction (amount, time, user ID) to Kafka’s transactions topic.
  • Consumer: A Spark Streaming job reads the topic, applies ML models to flag suspicious transactions (e.g., sudden large transfers), and alerts admins in <1 second.

B. Apache Flume

Flume is a lightweight, reliable tool for collecting and aggregating log data (e.g., web server logs, application logs). It uses a pipeline of agents to move data from sources to HDFS or other sinks.

Components:

  • Source: Receives data (e.g., syslog from Ncell’s servers).
  • Channel: Temporary storage (memory/disk) for data in transit.
  • Sink: Writes data to HDFS, HBase, or another system.

Example: Ncell Call Detail Records (CDR) Ingestion

  • Source: Ncell’s base stations generate CDRs (call duration, tower location, timestamp).
  • Flume Agent: Collects CDRs from 100+ towers and batches them every 5 minutes.
  • Sink: Writes to HDFS for later analytics (e.g., predicting network congestion).

3. Stream Processing: Turning Data into Action

Stream processing analyzes data as it arrives, enabling real-time insights. Unlike batch processing (which waits for data to accumulate), stream processing triggers actions instantly.

Event-time processingStateful operations (windows/sessions)Apache FlinkMicro-batch (DStreams)Structured Streaming (SQL API)Apache Spark StreamingTopology-based processingGuaranteed message processingApache StormStream Processing Frameworks
Comparison of major stream processing frameworks
classDiagram
  class KafkaTopic {
    +publish(data: JSON) void
    +partitions: int
  }
  class SparkStreaming {
    +DStream: micro-batch
    +transformations: map/filter/join
  }
  class FlinkJob {
    +eventTimeProcessing: boolean
    +statefulOps: window/session
  }
  KafkaTopic --> SparkStreaming : "feeds into"
  SparkStreaming --> FlinkJob : "triggers"
  FlinkJob --> "Pathao Dashboard" : "updates"
  note for KafkaTopic "Example: ride_requests topic\n  {lat, lng, ride_type}"
Class diagram of Pathao’s real-time pipeline components (Kafka → Spark → Flink)

Key Frameworks

Framework Language Use Case Example in Nepal
Apache Spark Streaming Scala/Java/Python Real-time analytics (e.g., traffic routing). Pathao’s dynamic pricing based on live rides.
Apache Flink Java/Scala Event-time processing (e.g., fraud detection). Khalti’s instant transaction validation.
Apache Storm Any High-throughput processing (e.g., IoT). NTC’s smart grid monitoring.

How Spark Streaming Works:

  1. Discretized Streams (DStreams): Breaks data into micro-batches (e.g., 1-second windows).
  2. Transformations: Applies operations like map, filter, or join to each batch.
  3. Output: Sends results to databases, dashboards, or triggers alerts.

Example: Pathao’s Traffic Routing

flowchart TD
    A["Pathao App (User Requests)"] -->|"JSON: {lat, lng, ride_type}"| B["Kafka Topic: ride_requests"]
    B --> C["Spark Streaming Job"]
    C -->|"Filter: high-demand zones"| D["Flink Job: Rebalance Drivers"]
    D --> E["Pathao Dashboard: Live Driver Locations"]
    E --> F["User: Faster Matching"]
  • Step 1: User requests a ride → data sent to Kafka.
  • Step 2: Spark Streaming filters requests in high-demand areas (e.g., Thapathali during rush hour).
  • Step 3: Flink rebalances drivers in real-time, reducing wait times.

4. Batch vs. Stream Processing: When to Use Which

Criteria Batch Processing Stream Processing
Latency High (hours/days) Low (milliseconds)
Use Case Historical analysis (e.g., yearly sales trends) Real-time decisions (e.g., fraud alerts)
Tools Hadoop MapReduce, Spark Batch Kafka, Flink, Spark Streaming
Data Volume Large, static datasets Continuous, unbounded streams
Example in Nepal Daraz’s end-of-month inventory reports Ncell’s real-time network congestion alerts
06121824Batch (Daraz)24Stream (Pathao)1Micro-batch (NEPSE)5
Average processing latency in hours (batch) vs. seconds (stream/micro-batch) for Nepali examples

When to Choose Stream Processing?

  • Fraud detection (e.g., Khalti blocking unauthorized transactions).
  • IoT monitoring (e.g., NTC detecting power outages in real-time).
  • Personalization (e.g., YouTube recommending videos as you watch).

In the Real World

  1. eSewa’s Fraud Detection

    • Idea Used: Stream processing with Kafka + Spark.
    • How: Every transaction is published to a Kafka topic. A Spark Streaming job checks for anomalies (e.g., 10 transactions in 1 minute from one phone) and blocks the account instantly. Saved $2M/year in fraud losses.
  2. Pathao’s Dynamic Pricing

    • Idea Used: Real-time data ingestion + micro-batching.
    • How: Ride requests flow into Kafka. Spark Streaming aggregates demand by location (e.g., 500 requests in Kathmandu’s Ring Road). Flink adjusts surge pricing dynamically, ensuring drivers earn more during peak times.
  3. NTC’s Smart Grid Monitoring

    • Idea Used: Stream processing for IoT data.
    • How: Smart meters send power usage data every second to Flume. A Flink job detects sudden drops (e.g., a transformer failure) and reroutes power, preventing blackouts. Reduced outages by 40% in 2023.

5. Challenges in Data Ingestion & Stream Processing

Challenge Cause Solution
Data Skew Uneven distribution of keys (e.g., one user generates 90% of transactions). Use salting (adding random prefixes to keys) or partitioning strategies.
Late Data Network delays or slow producers. Use watermarks (Flink) or late-event handling (Spark).
State Management Tracking state (e.g., user session) in streams. Use checkpointing (periodic snapshots of state).
Scalability Handling millions of events per second. Distribute across Kafka partitions or Spark executors.
02.254.56.759Data Volume9Latency7Fault Tolerance8Scalability6Cost5
Common challenges in Nepali data pipelines (1-10 scale, 10=most critical)

Example: Handling Late Data in Khalti

  • Problem: A transaction takes 3 seconds to arrive due to network lag.
  • Solution: Flink uses a watermark (e.g., "no data after 5 seconds is late") and buffers the transaction. If it arrives within 5 seconds, it’s processed; otherwise, it’s discarded or logged for later.

Worked Example: Designing a Data Pipeline for Daraz

Scenario: Daraz wants to analyze real-time order statuses (placed, shipped, delivered) to improve logistics and detect abandoned carts.

Step 1: Define Requirements

  • Data Sources:
    • Web app (orders placed via daraz.com.np).
    • Mobile app (orders via Daraz app).
    • Third-party logistics (e.g., Nabil Courier API).
  • Processing Needs:
    • Detect abandoned carts (no payment in 10 minutes).
    • Alert logistics team if an order is stuck in transit for >24 hours.
  • Output:
    • Real-time dashboard for customer service.
    • Historical data in HDFS for monthly reports.

Step 2: Choose Tools

Component Tool Reason
Ingestion Apache Kafka Handles high-throughput order events.
Stream Processing Apache Flink Low-latency event-time processing.
Storage HDFS + HBase Scalable storage for historical data.
Dashboard Apache Superset Real-time visualization.

Step 3: Pipeline Design

sequenceDiagram
    participant User
    participant DarazApp
    participant Kafka
    participant Flink
    participant HBase
    participant Dashboard

    User->>DarazApp: Places order (JSON: {order_id, user_id, items, status})
    DarazApp->>Kafka: Publish to topic "orders"
    Kafka->>Flink: Stream of orders
    Flink->>Flink: Rule 1: If status="paid" and no activity in 10 mins → "abandoned"
    Flink->>Flink: Rule 2: If status="shipped" and no update in 24 hrs → "delayed"
    Flink->>HBase: Store order history
    Flink->>Dashboard: Update real-time metrics
    Dashboard->>Support: Alert for abandoned/delayed orders

Step 4: Handling Edge Cases

  1. Abandoned Cart Detection:
    • Flink maintains a stateful window (10-minute tumbling window) for each order.
    • If no payment_confirmation event arrives, it triggers an email/SMS to the user.
  2. Delayed Shipments:
    • Flink joins the orders stream with a logistics stream (from Nabil Courier).
    • If logistics.status = "in_transit" for >24 hours, it alerts the logistics team.

Exam Tip

What Examiners Look For

  1. Distinguish Batch vs. Stream:

    • Batch: "Processed in fixed intervals (e.g., hourly sales reports)."
    • Stream: "Processed as it arrives (e.g., fraud detection)."
    • Common Mistake: Confusing Kafka (stream ingestion) with HDFS (batch storage).
  2. Tool Selection:

    • Kafka for high-throughput ingestion.
    • Flink/Spark Streaming for real-time processing.
    • Flume for log aggregation (e.g., server logs).
  3. Real-World Applications:

    • Link concepts to Nepali examples (e.g., "Pathao uses Kafka for ride requests").
    • Explain why a tool is chosen (e.g., "Flink’s event-time processing is better than Spark’s micro-batching for fraud detection").
  4. Diagrams:

    • Draw pipeline flows (producers → Kafka → Spark → output).
    • Label components (topics, partitions, consumers).
  5. Challenges:

    • Mention data skew, late data, or state management and how they’re solved (e.g., "watermarks in Flink").

Sample Exam Questions & Answers

Q1: "Explain how eSewa can use stream processing to detect fraudulent transactions." Answer:

eSewa can use Apache Kafka to ingest transaction data in real-time. A Spark Streaming job would:

  1. Subscribe to the transactions topic.
  2. Apply rules (e.g., ">5 transactions in 1 minute from one IP").
  3. Flag suspicious transactions and block the account via a Kafka consumer that triggers an alert. Tools: Kafka (ingestion), Spark (processing), Redis (fast lookup for user profiles).

Q2: "Compare batch and stream processing with examples from Nepali companies." Answer:

Aspect Batch Processing Stream Processing
Example Daraz’s monthly sales reports (Hadoop MapReduce). Pathao’s dynamic pricing (Spark Streaming).
Latency High (daily/weekly). Low (<1 second).
Use Case Year-end financial audits. Real-time fraud detection.

Key Formulas & Metrics

  1. Throughput (Events/sec):

    • Example: If Spark processes 10,000 transactions in 5 seconds, throughput = 2,000 events/sec.
  2. End-to-End Latency:

    • Example: Kafka (100ms) + Flink (200ms) + Dashboard (50ms) = 350ms total.

Common Pitfalls

  • Ignoring Late Data: Always design for out-of-order events (e.g., Flink’s allowedLateness).
  • Over-Partitioning Kafka: Too many partitions increase overhead; too few cause bottlenecks.
  • Stateful Processing Without Checkpoints: Flink/Spark Streaming jobs must checkpoint state to avoid losing progress on failure.

Further Reading

  • Books:
    • Designing Data-Intensive Applications (Martin Kleppmann) – Chapter 5 (Stream Processing).
    • Apache Kafka (Neha Narkhede) – For deep dives into Kafka internals.
  • Tools to Try:
    • Set up a local Kafka cluster using Docker.
    • Run a Spark Streaming job on a sample dataset (e.g., Twitter feeds).

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

Discussion

Loading…