Plateforme de traitement de données embarquées pour une flotte de véhicules connectés. Projet pédagogique de data engineering — EPITA, 4 contributeurs.
Des millions de capteurs émettent leur position et leur comportement de conduite plusieurs fois par seconde (~200 Go/jour). VigiRoute alimente en parallèle un service Alert temps réel (notifier les conducteurs proches d'un danger en moins de quelques secondes) et un service Stats batch (cartographier les zones à risque sur le long terme pour des analystes).
Les choix d'architecture et leurs justifications sont détaillés dans docs/architecture.md. Le diagramme mermaid associé est dans docs/architecture.mmd.
| Brique | Outil |
|---|---|
| Langage | Scala 2.13 fonctionnel |
| Build | sbt multi-modules (6 build.sbt, un par sous-projet + racine) |
| Ingestion stream | Apache Kafka (topic vehicle-telemetry, format Avro, partitionné par zone géo) |
| Traitement temps réel | fs2-kafka + cats-effect (jobs JVM purement fonctionnels) |
| Index géospatial | Redis (GeoSpatial) via redis4cats |
| Bus d'alertes | Kafka (topic driver-alerts, JSON) → consumer notif |
| Data lake | S3 — en local MinIO — couches Bronze (Avro container files) / Silver (Parquet) / Gold (agrégations) |
| Batch ETL | Apache Spark 3.5 (DataFrame API) |
| Module sbt | Main |
|---|---|
producer |
com.vigiroute.producer.Main |
streaming/danger_detector |
com.vigiroute.detector.Main |
streaming/notifier |
com.vigiroute.notifier.Main |
bronze_writer |
com.vigiroute.bronze.Main |
batch |
com.vigiroute.batch.{BronzeToSilver,SilverToGold,Analytics,AnalyticsExport} |
streaming/position_updater |
com.vigiroute.positionupdater.Main |
demo-bridge |
com.vigiroute.demo.Main (pont SSE/REST de démo) |
- Docker ≥ 24 et Docker Compose v2
- JDK 17+ (Temurin recommandé)
- sbt ≥ 1.10 (
brew install sbt,sdk install sbt, ou via Coursier) - ~6 Go de RAM disponibles pour la stack Docker complète
# 1. Lancer Kafka, Redis, MinIO, Spark
docker compose up -d
# 2. Vérifier les topics Kafka
docker compose exec kafka kafka-topics --bootstrap-server kafka:29092 --list
# Attendu : driver-alerts, vehicle-telemetry
# 3. Reset total (volumes inclus)
docker compose down -vDepuis la racine du repo :
# Compiler tout
sbt compile
# Lancer chaque service (dans un terminal séparé)
sbt "producer/run"
sbt "positionUpdater/run"
sbt "dangerDetector/run"
sbt "notifier/run"
sbt "bronzeWriter/run"
# Jobs batch Spark (via le cluster du compose)
# NB : depuis le conteneur Spark, l'endpoint MinIO est http://minio:9000
# (pas localhost), d'où le -e S3_ENDPOINT=…
sbt "batch/assembly"
docker compose exec -e S3_ENDPOINT=http://minio:9000 spark-master /opt/spark/bin/spark-submit \
--class com.vigiroute.batch.BronzeToSilver \
/workspace/batch/target/scala-2.13/batch.jar
# Idem pour SilverToGold puis Analytics.
# (Pas de --packages : spark-avro et hadoop-aws sont déjà dans le fat jar.)La procédure de test bout-en-bout (infra → flux → Spark → analyse) est détaillée
dans docs/TESTING.md.
Interface web de démo (carte temps réel + dashboard analytics + vue archi),
branchée en direct sur la stack (SSE) ou en replay d'une trace enregistrée.
Procédure complète : docs/DEMO.md.
scripts/demo.sh # tout-en-un, mode REPLAY (aucune stack requise)
scripts/demo.sh live # stack Docker + pipeline complet + pont + frontLancement manuel équivalent :
sbt "demoBridge/run" # pont SSE/REST (port 8090)
cd web && npm install && npm run dev # front (port 5173)| Service | URL | Identifiants |
|---|---|---|
| MinIO (console S3) | http://localhost:9001 | minioadmin / minioadmin |
| Spark master UI | http://localhost:8080 | — |
| Spark worker UI | http://localhost:8081 | — |
.
├── README.md Ce fichier
├── build.sbt Racine sbt — agrège les 6 sous-projets
├── project/
│ ├── build.properties Version sbt
│ ├── plugins.sbt sbt-assembly
│ └── Dependencies.scala Versions centralisées
├── docker-compose.yml Stack locale (Kafka, Redis, MinIO, Spark)
│
├── schemas/
│ └── telemetry.avsc Schéma Avro de référence (figé)
│
├── producer/ Simulateur IoT
├── streaming/
│ ├── position_updater/ MAJ Redis GEOADD
│ ├── danger_detector/ Détection d'alertes
│ └── notifier/ Push notif
├── bronze_writer/ Kafka -> S3 (Avro container files)
├── batch/ Spark Bronze/Silver/Gold + analytics
│
└── docs/
├── architecture.md Justifications d'architecture
├── architecture.mmd Diagramme mermaid
└── TESTING.md Guide de test & démo de bout en bout
Chaque sous-dossier a son README.md détaillant son rôle.
Les 5 composants du PoC tournent de bout en bout :
- le
producersimule la flotte et pousse la télémétrie Avro dans Kafka (clé =zone_id) ; position_updater+danger_detectorsélectionnent les alertes (4 règles, voisins via RedisGEORADIUS) ;- le
notifiertraite les alertes (push mocké, retry/backoff, rate-limit, DLQ) ; bronze_writerpuis les jobsBronzeToSilver/SilverToGoldalimentent le data lake bronze/silver/gold ;batch/Analyticsrépond aux 4 questions en Spark (le SQL Athena dansanalytics/est un chemin alternatif).
Tests : sbt test couvre la logique pure, les codecs Avro, les règles de détection et
les transformations Spark. Le reste (déploiement cloud, dashboard) relève de la partie perso.
- Branche
mainprotégée, pas de push direct. - Une tâche = une branche
feature/…= une PR avec au moins un relecteur. - Convention de commits :
feat: …,fix: …,docs: …,refactor: …. - Chaque dev commite sous sa propre identité (
git config user.email).