StreamML — Distributed ML Inference Platform
Real-time fraud detection pipeline built on distributed systems principles.
Target: 50k predictions/sec | <10ms p99 latency | <15s node recovery
[Producer] → [Kafka: 3 brokers] → [Consumer Workers]
↓
[ML Inference]
↓
[FastAPI + Monitoring]
Messaging: Apache Kafka (KRaft mode, 3 brokers)
Workers: Python confluent-kafka (Phase 2) → Ray (Phase 4)
Pipeline: Apache Spark Streaming (Phase 3)
API: FastAPI + Nginx (Phase 5)
Monitoring: Prometheus + Grafana (Phase 6)
Deployment: Docker Compose → Kubernetes (Phase 7-8)
3-broker Kafka cluster in KRaft mode (no ZooKeeper)
Topic: transactions — 3 partitions, replication factor 3
Producer: generates synthetic fraud detection events, keys by user_id
Consumer: ml-inference-group, manual offset commit, fault isolated
Fault tolerance demo: killed broker mid-stream, zero message loss,
automatic leader re-election, recovery under 15 seconds
Scenario
Result
Kill 1 of 3 brokers mid-stream
Zero message loss
Broker recovery time
< 15 seconds
Min brokers needed for writes
2 of 3 (MIN_INSYNC_REPLICAS)
Kafka Cluster — 3 brokers, 159/159 in-sync replicas
Topic: transactions — 3 partitions, replication factor 3
Fault tolerance — broker killed mid-stream, zero message loss
Spark Structured Streaming pipeline consuming Kafka transactions topic
Windowed feature computation: 10-minute sliding window, 2-minute slide
Features computed per user: txn_count_10min, avg_amount_10min, max_amount_10min, min_amount_10min
Fraud signal identified: users with high txn_count + high avg_amount in short windows
Runs locally with local[*] master using manually downloaded Kafka connector JARs
Spark computing user features from live Kafka stream
Producer + Spark pipeline running side by side
3 Ray actor workers, each loading XGBoost fraud model once into RAM
Real dataset: 590,540 IEEE-CIS transactions, 3.5% fraud rate
Model upgraded: Random Forest → XGBoost with feature engineering
35 features including time encoding, log transforms, card metadata
AUC-ROC: 0.9343 | Recall: 0.82 | Precision: 0.23
Avg inference latency: 29-41ms on laptop
Fault tolerance: worker crash detected, Kafka offset uncommitted, auto-reprocessed
FRAUD ALERT triggered for transactions with fraud_probability > 0.5
Phase 4 — Model Improvement Journey
Version
Model
Features
AUC-ROC
Recall
v1
Random Forest
15 raw
0.8805
0.30
v2
Random Forest + balanced
15 raw
0.8923
0.69
v3
XGBoost + engineering
35 features
0.9343
0.82
Metric
Value
Training dataset
590,540 real transactions
Fraud rate
3.5%
AUC-ROC
0.9343
Recall (fraud)
0.82
Ray workers
3 actors
Avg latency (laptop)
~30ms
Model load strategy
Once per actor at startup
XGBoost model performance — AUC 0.9343
Ray workers + producer — fraud alerts in real time
FastAPI serving layer with 3 endpoints: /predict, /predict/batch, /metrics
Ray workers serve XGBoost model via HTTP — no script needed, just a POST request
Redis caching — repeated transactions return cached result instantly
Background Kafka logging — every prediction logged asynchronously without blocking response
Risk scoring: LOW / MEDIUM / HIGH / CRITICAL based on fraud probability
Auto-generated interactive API docs at /docs
Phase 5 — Live Prediction Example
POST /predict
{
"user_id" : " user_42" ,
"amount" : 4981.64 ,
"card1" : 9500 ,
"C1" : 3 ,
"D1" : 14
Response:
{
"fraud_probability" : 0.5419 ,
"is_fraud" : true ,
"risk_level" : " HIGH" ,
"latency_ms" : 732.96 ,
"worker_id" : 1 ,
"cached" : false
}
Interactive API docs — auto-generated by FastAPI
Live fraud prediction — HIGH risk detected
Real-time metrics endpoint— live worker stats
Phase 5 Performance Numbers
Metric
Value
Note
Avg latency (laptop)
516ms
Ray IPC overhead on Windows
Worker 1 best latency
155ms
After model warmup
Concurrent workers
3
Round-robin load balanced
Uptime
1337s
Stable API
Fraud detection rate
100% on test data
High-value test transactions
Phase 5 — Load Test Results
Metric
Value
Throughput
102.1 req/sec
Success rate
100% (200/200)
Avg latency
184ms
p50 latency
140ms
p95 latency
547ms
p99 latency
603ms
Platform
Single laptop, 8 Docker containers
Prometheus scraping FastAPI metrics every 15 seconds
Grafana dashboard with 4 live panels
Custom metrics: streamml_predictions_total, streamml_active_workers, inference latency histogram
Load test spike visible in dashboard — 200 predictions at 23:20
Real-time worker scaling: 0 → 3 Ray workers on demand
Metric
Value
Scrape interval
15 seconds
Dashboard panels
4
Metrics tracked
Request rate, p99 latency, predictions, active workers
Data retention
7 days
Live Grafana dashboard — load test spike visible at 23:20
# Start Kafka cluster
docker compose up -d
# Terminal 1 — consumer
python consumer/consumer.py
# Terminal 2 — producer
python producer/producer.py
# Kafka UI
open http://localhost:8080
streamml/
├── docker-compose.yml # 3-broker Kafka + UI
├── producer/
│ └── producer.py # Transaction event generator
├── consumer/
│ └── consumer.py # Fault-isolated consumer
├── spark/
│ └── feature_pipeline.py
├── ml/
│ ├── fraud_detector.py
│ ├── train_model.py
│ ├── train_modelSYN.py
│ └── models/
├── jars/
└── docs/
├── TROUBLESHOOTING.md
└── screenshots/
└── README.md