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
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:
- Producers (e.g., eSewa app) publish data to topics (logical channels like
transactionsoruser_activity). - Brokers (Kafka servers) store data in partitions (scalable storage units).
- 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
transactionstopic. - 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.
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:
- Discretized Streams (DStreams): Breaks data into micro-batches (e.g., 1-second windows).
- Transformations: Applies operations like
map,filter, orjointo each batch. - 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 |
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
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.
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.
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. |
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).
- Web app (orders placed via
- 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 ordersStep 4: Handling Edge Cases
- Abandoned Cart Detection:
- Flink maintains a stateful window (10-minute tumbling window) for each order.
- If no
payment_confirmationevent arrives, it triggers an email/SMS to the user.
- Delayed Shipments:
- Flink joins the
ordersstream with alogisticsstream (from Nabil Courier). - If
logistics.status = "in_transit"for >24 hours, it alerts the logistics team.
- Flink joins the
Exam Tip
What Examiners Look For
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).
Tool Selection:
- Kafka for high-throughput ingestion.
- Flink/Spark Streaming for real-time processing.
- Flume for log aggregation (e.g., server logs).
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").
Diagrams:
- Draw pipeline flows (producers → Kafka → Spark → output).
- Label components (topics, partitions, consumers).
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:
- Subscribe to the
transactionstopic.- Apply rules (e.g., ">5 transactions in 1 minute from one IP").
- 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
Throughput (Events/sec):
- Example: If Spark processes 10,000 transactions in 5 seconds, throughput = 2,000 events/sec.
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…