Pipeline Big Data temps réel de détection de fraude bancaire, construit autour de la Modern Data Stack : Kafka → Spark → Snowflake → Metabase, avec un modèle XGBoost intégré au flux.
Projet Big Data — ENSIAS, 2026
Détecter en temps réel les transactions frauduleuses sur un flux bancaire simulé, en couvrant toute la Big Data Value Chain :
- Ingestion — collecte du flux de transactions
- Preprocessing — nettoyage et structuration à la volée
- Storage — persistance dans un Data Warehouse cloud
- Processing — traitement distribué
- Analysis + ML — prédiction de fraude par modèle entraîné
- Visualization — dashboard live des fraudes détectées
CSV (IEEE-CIS Fraud)
│
▼
Producer Python ──► KAFKA ──► SPARK (cluster master + worker) ──► SNOWFLAKE ──► METABASE
(simule le flux (broker) (Structured Streaming + ML) (cloud DW) (dashboard live)
temps réel) │
▼
MLflow + XGBoost
(tracking + modèle)
Tout l'environnement (sauf Snowflake, qui est cloud) tourne dans Docker Compose : un seul docker-compose up lance Kafka, Zookeeper, Spark master et Spark worker.
| Couche | Outils |
|---|---|
| Ingestion | Apache Kafka + Zookeeper |
| Streaming | Apache Spark Structured Streaming (cluster master + worker) |
| Machine Learning | XGBoost + MLflow (tracking et versioning du modèle) |
| Storage | Snowflake (Data Warehouse cloud) |
| Visualization | Metabase (dashboard temps réel) |
| Orchestration | Docker Compose |
| Données | IEEE-CIS Fraud Detection (Kaggle, ~590k transactions, >500 MB) |
Dataset : IEEE-CIS Fraud Detection (Kaggle)
- ~590 000 transactions réelles (Vesta Corporation)
- Plus de 400 features (numériques, catégorielles, temporelles)
- Label binaire
isFraud(très déséquilibré, ~3,5 % de fraudes)
Le producer lit le CSV ligne par ligne et envoie chaque transaction en JSON dans Kafka, simulant un flux temps réel à débit configurable.
- Cluster Spark distribué (master + worker) lancé via Docker Compose, et non en mode local — pour montrer un vrai pipeline scalable.
- Snowflake comme Data Warehouse cloud, branché en sortie du flux Spark via le connecteur officiel
spark-snowflake. - Modèle XGBoost entraîné en batch et chargé dans le flux Spark pour prédire la fraude à la volée sur chaque micro-batch.
- MLflow pour tracker les expériences et versionner le modèle, ce qui permet au binôme de travailler en parallèle sur le ML pendant que le pipeline streaming est développé en parallèle.
- Identifiants Snowflake gérés via
.env(jamais versionnés sur GitHub) pour respecter les bonnes pratiques de sécurité.
fraud-detection-bigdata/
├── docker-compose.yml # Kafka + Zookeeper + Spark master + worker
├── Dockerfile.spark # Image Spark personnalisée
├── producer/
│ └── producer.py # Lit le CSV et envoie les transactions dans Kafka
├── streaming/
│ └── spark_stream.py # Lit Kafka, applique le modèle, écrit dans Snowflake
├── ml/ # Entraînement XGBoost + tracking MLflow
├── jars/ # Connecteurs Kafka et Snowflake (non versionnés)
├── data/ # Dataset IEEE-CIS (non versionné)
├── .env # Identifiants Snowflake (non versionné)
└── README.md
- Docker Desktop installé et lancé
- Un compte Snowflake (essai gratuit)
- Le dataset IEEE-CIS Fraud Detection téléchargé dans
data/ - Python 3.x pour le producer :
pip install kafka-python pandas
Crée un fichier .env à la racine avec tes identifiants Snowflake :
SNOWFLAKE_ACCOUNT=TON-ACCOUNT-ID
SNOWFLAKE_USER=ton_user
SNOWFLAKE_PASSWORD=ton_password
SNOWFLAKE_WAREHOUSE=COMPUTE_WH
SNOWFLAKE_DATABASE=FRAUD_DB
SNOWFLAKE_SCHEMA=PUBLIC
Puis crée la table dans un SQL Worksheet Snowflake :
CREATE DATABASE IF NOT EXISTS FRAUD_DB;
CREATE TABLE IF NOT EXISTS FRAUD_DB.PUBLIC.TRANSACTIONS (
TransactionID FLOAT, isFraud FLOAT, TransactionDT FLOAT,
TransactionAmt FLOAT, ProductCD STRING, card4 STRING, card6 STRING,
prediction FLOAT
);docker-compose up -dL'interface Spark est disponible sur http://localhost:8090.
docker exec kafka kafka-topics --create --topic transactions \
--bootstrap-server localhost:9092 --partitions 1 --replication-factor 1docker exec -it spark-master /opt/spark/bin/spark-submit \
--master spark://spark-master:7077 \
--conf spark.jars.ivy=/tmp/.ivy2 \
--jars /app/jars/spark-sql-kafka-0-10_2.12-3.5.0.jar,/app/jars/spark-token-provider-kafka-0-10_2.12-3.5.0.jar,/app/jars/kafka-clients-3.4.1.jar,/app/jars/commons-pool2-2.11.1.jar \
--packages net.snowflake:snowflake-jdbc:3.16.1,net.snowflake:spark-snowflake_2.12:3.1.1 \
/app/streaming/spark_stream.pypython producer/producer.pySELECT COUNT(*) FROM FRAUD_DB.PUBLIC.TRANSACTIONS;
SELECT * FROM FRAUD_DB.PUBLIC.TRANSACTIONS WHERE prediction = 1 LIMIT 10;- Anas BENAMARA — Pipeline streaming (Kafka, Spark, Snowflake), Docker, intégration
- Aymane BOUGHALEB — Machine Learning (XGBoost, MLflow), dashboard Metabase
Projet académique réalisé dans le cadre du cours Big Data — ENSIAS, 2026. Distribué sous licence MIT.