Big Data and AnalyticsUnit 99 min read

Machine Learning on Big Data: Algorithms, Scalability & Applications

Unit 9 of Big Data and Analytics explores how machine learning (ML) techniques are adapted for big data—scalable algorithms, distributed training, feature engineering for massive datasets, and real-world deployments in analytics pipelines. Covers supervised/unsupervised deep learning, model optimization, and ethical co

Core Concepts: Why ML Needs Big Data (and Vice Versa)

Machine learning thrives on data volume, velocity, and variety—three pillars of big data. Traditional ML models (e.g., logistic regression, decision trees) fail at scale due to:

  • Computational limits: A single server cannot process terabytes of data in memory.
  • Feature explosion: High-dimensional data (e.g., text, images) requires distributed feature extraction.
  • Real-time demands: Streaming data (e.g., social media, IoT) needs online learning.

Big data enables ML to:

  1. Train more accurate models (e.g., deep learning on millions of images).
  2. Detect patterns invisible in small datasets (e.g., fraud in bank transactions).
  3. Adapt to dynamic environments (e.g., recommendation systems updating hourly).

1. Scalable ML Algorithms for Big Data

Not all ML algorithms scale. Here’s how they adapt:

Algorithm Type Traditional Approach Big Data Adaptation Tools/Frameworks
Supervised Learning Batch training (e.g., scikit-learn) Distributed stochastic gradient descent (SGD) Spark MLlib, TensorFlow
Unsupervised Learning K-means (single-node) Mini-batch K-means, approximate clustering Apache Flink, Faiss (Facebook)
Deep Learning GPU-accelerated (e.g., PyTorch) Distributed training (parameter servers) Horovod, Ray, BigDL
Reinforcement Learning Q-learning (tabular) Deep Q-Networks (DQN) with experience replay RLlib (AWS)

2. Distributed Training: How Models Learn Across Clusters

Problem: A single GPU cannot fit a model trained on 100GB of data. Solution: Data parallelism and model parallelism.

A. Data Parallelism (Most Common)

  • How it works: Split data into chunks (shards), train identical models on each shard, then aggregate gradients (e.g., via AllReduce).
  • Example: Training a language model on 1PB of text (e.g., Google’s BERT).
    flowchart TD
      A["Data Shard 1"] --> B["Model Copy 1"]
      A --> C["Model Copy 2"]
      A --> D["Model Copy N"]
      B --> E["Gradient Aggregator"]
      C --> E
      D --> E
      E --> F["Updated Global Model"]

B. Model Parallelism

  • How it works: Split the model itself (e.g., different layers on different GPUs).
  • Example: Training a Transformer with 100+ layers (used in Google’s T5).

3. Feature Engineering for Big Data

Big data often comes in unstructured formats (text, images, logs). Key techniques:

A. Text Data

  • TF-IDF → Word2Vec/Glove: Convert words to dense vectors (e.g., for sentiment analysis).
  • Example: eSewa’s chatbot uses Word2Vec to classify user queries (e.g., "bill payment" vs. "complaint").

B. Image Data

  • CNNs (Convolutional Neural Networks): Extract features from pixels (e.g., Daraz’s product tagging).
  • Example: Ncell’s OCR uses CNNs to read handwritten numbers on SIM cards.

C. Time-Series Data

  • LSTMs/Transformers: Capture temporal patterns (e.g., NTC’s electricity demand forecasting).
  • Example: Pathao’s surge pricing uses LSTMs to predict driver availability during festivals.

4. Big Data ML Workflow: From Raw Data to Model

sequenceDiagram
    participant DataSource as "Data Source (e.g., NEPSE Stock Data)"
    participant Preprocess as "Preprocessing (Spark/PySpark)"
    participant FeatureStore as "Feature Store (Feast)"
    participant Train as "Distributed Training (Spark MLlib)"
    participant Evaluate as "Evaluation (MLflow)"
    participant Deploy as "Serving (Kubernetes)"

    DataSource->>Preprocess: Ingest (Kafka/S3)
    Preprocess->>FeatureStore: Extract features
    FeatureStore->>Train: Train (SGD/Adam)
    Train->>Evaluate: Metrics (AUC, RMSE)
    Evaluate->>Deploy: Deploy as API

Key Steps:

  1. Data Ingestion: Stream from Kafka, batch from S3 (e.g., Khalti’s transaction logs).
  2. Preprocessing: Clean, normalize, and sample (e.g., NEPSE’s missing stock data).
  3. Feature Store: Reuse features across models (e.g., user behavior features for Daraz recommendations).
  4. Training: Use Spark MLlib or TensorFlow on Kubernetes.
  5. Serving: Deploy as a REST API (e.g., Ncell’s fraud detection model).

5. Real-World Applications in Nepal

A. Fraud Detection in Banks (e.g., NMB, Global IME)

  • Problem: 10,000+ transactions/sec; fraudsters mimic patterns.
  • Solution: Isolation Forest (unsupervised) trained on Spark to detect anomalies.
  • Example: Global IME’s credit card fraud model flags transactions in <100ms using approximate nearest neighbors (ANN).

B. Traffic Prediction (Kathmandu Metropolitan City)

  • Problem: 500+ traffic cameras generate 1TB/day of video data.
  • Solution: YOLO (You Only Look Once) + LSTM to predict congestion.
  • Example: KMC’s smart traffic lights adjust timings based on real-time predictions.

C. Agricultural Yield Prediction (AGRIBANK, Nepal)

  • Problem: 2M+ small farmers; weather + soil data varies.
  • Solution: XGBoost on Spark to predict rice/wheat yield.
  • Example: AGRIBANK’s loan approval uses ML to assess farmer risk.

6. Challenges and Ethical Considerations

Challenge Big Data ML Solution Ethical Risk
Bias in Training Data Stratified sampling, fairness-aware algorithms Excludes minority groups (e.g., rural users in Khalti)
Privacy (GDPR/PDPA) Federated learning, differential privacy Ncell’s call logs could leak user locations
Model Interpretability SHAP/LIME for explainability NEPSE’s stock predictions must justify trades
Concept Drift Online learning (e.g., Pathao’s dynamic pricing) Model becomes obsolete (e.g., COVID-19 changed travel patterns)

In the Real World

  1. Khalti’s Fraud Detection

    • Idea: Real-time anomaly detection using Isolation Forest on Spark.
    • How: Every transaction is scored for fraud probability within 50ms using pre-trained embeddings of user behavior.
  2. Daraz’s Recommendation Engine

    • Idea: Collaborative filtering + deep learning (two-tower model).
    • How: 100M+ user-item interactions are processed via Apache Flink to update recommendations hourly.
  3. Ncell’s Churn Prediction

    • Idea: XGBoost on historical call/SMS data to predict customer attrition.
    • How: 300M+ records are processed nightly to identify at-risk users for retention offers.

Exam Tip

  1. Compare algorithms: Always contrast batch vs. online learning, centralized vs. distributed training.

    • Example: "Why does Spark MLlib use mini-batch gradient descent instead of full-batch for big data?"
  2. Link to Nepalese context:

    • NEPSE: Time-series forecasting for stock prices.
    • NTC: Predictive maintenance for power grids.
    • eSewa: NLP for customer support chatbots.
  3. Diagrams are worth marks:

    • Draw data parallelism vs. model parallelism.
    • Sketch a big data ML pipeline (ingest → preprocess → train → serve).
  4. Common pitfalls:

    • ❌ Saying "big data = more data always improves ML" (overfitting risk).
    • ❌ Ignoring feature scaling in distributed SGD (diverges if unscaled).

Worked Example: Predicting Loan Defaults for a Nepalese Bank

Scenario: A bank has 500K loan records with features like income, employment status, and past defaults. Goal: Predict default probability.

Step 1: Data Splitting (Spark)

from pyspark.ml.feature import VectorAssembler
from pyspark.ml.classification import LogisticRegression

# Assemble features
assembler = VectorAssembler(inputCols=["income", "employment_status", "past_defaults"], outputCol="features")

# Split data (80-20 train-test)
train_data, test_data = df.randomSplit([0.8, 0.2])

Step 2: Distributed Training (Spark MLlib)

lr = LogisticRegression(featuresCol="features", labelCol="default", maxIter=10, regParam=0.01)
model = lr.fit(train_data)

Step 3: Evaluation

predictions = model.transform(test_data)
print(f"AUC: {predictions.select('prediction', 'default').stat.areaUnderROC()}")
# Output: AUC = 0.89 (good!)

Real-World Twist:

  • NMB Bank uses a similar model but with XGBoost on Spark for better non-linearity.
  • Challenge: Class imbalance (only 5% defaults). Solution: SMOTE oversampling in Spark.

Key Formulas to Remember

  1. Gradient Descent Update (Distributed):

    • : Number of data shards.
    • : Learning rate.
  2. AUC-ROC (for classification):

    • TPR: True Positive Rate.
  3. Mean Squared Error (Regression):

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

Discussion

Loading…