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:
- Train more accurate models (e.g., deep learning on millions of images).
- Detect patterns invisible in small datasets (e.g., fraud in bank transactions).
- 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 APIKey Steps:
- Data Ingestion: Stream from Kafka, batch from S3 (e.g., Khalti’s transaction logs).
- Preprocessing: Clean, normalize, and sample (e.g., NEPSE’s missing stock data).
- Feature Store: Reuse features across models (e.g., user behavior features for Daraz recommendations).
- Training: Use Spark MLlib or TensorFlow on Kubernetes.
- 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
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.
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.
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
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?"
Link to Nepalese context:
- NEPSE: Time-series forecasting for stock prices.
- NTC: Predictive maintenance for power grids.
- eSewa: NLP for customer support chatbots.
Diagrams are worth marks:
- Draw data parallelism vs. model parallelism.
- Sketch a big data ML pipeline (ingest → preprocess → train → serve).
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
Gradient Descent Update (Distributed):
- : Number of data shards.
- : Learning rate.
AUC-ROC (for classification):
- TPR: True Positive Rate.
Mean Squared Error (Regression):
Based on the TU BIM syllabus for Big Data and Analytics (IT278), unit 9.
Discussion
Loading…