Skip to content

Latest commit

 

History

21 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Pipeline ETL Météo Paris

kafka-python requests google-cloud-bigquery pyspark Kafka Zookeeper Schema Registry Control Center Airflow PostgreSQL Spark Bitnami BigQuery NumPy Pandas PyArrow Maven Central

Flux Open‑Meteo vers Kafka, traitement avec Spark Structured Streaming, chargement dans BigQuery et visualisation via Looker Studio. Orchestration par Airflow.

Vous pouvez changer la ville selon la latitude et la longitude dans l’API. Le ville actuel est Paris.

Architecture

System Architecture

Topics Kafka

  • Horaire: weather_paris
  • Quotidien: daily_paris

Produits par dags/kafka_stream.py (appel API Open‑Meteo, normalisation et envoi vers Kafka).

Consommateur Spark Streaming

  • Entrée: spark_stream.py
  • Bootstrap Kafka:
    • Dans les conteneurs: broker:29092
    • Depuis l’hôte: localhost:9092

BigQuery

  • Project ID: VOTRE_PROJECT_ID (modifiable dans spark_stream.py)
  • Dataset: VOTRE_NOM_DATASET (modifiable dans spark_stream.py)
  • Tables:
    • Horaire: meteo_hourly
    • Quotidien: meteo_daily

Schémas selon spark_stream.py:

  • meteo_hourly: id(STRING, REQUIRED), time_text(STRING), time_ts(TIMESTAMP), latitude(FLOAT64), longitude(FLOAT64), timezone(STRING), timezone_abbreviation(STRING), temperature_2m(FLOAT64), relative_humidity_2m(FLOAT64), apparent_temperature(FLOAT64), precipitation(FLOAT64), surface_pressure(FLOAT64), cloud_cover(FLOAT64), wind_speed_10m(FLOAT64)
  • meteo_daily: id(STRING, REQUIRED), date_text(STRING), date_ts(DATE), latitude(FLOAT64), longitude(FLOAT64), timezone(STRING), timezone_abbreviation(STRING), sunrise_time(STRING), sunset_time(STRING), sunshine_duration(FLOAT64), sunshine_duration_time(STRING)

Démarrage rapide

  1. Préparer les credentials BigQuery (placer config/config.json).
  2. Lancer les services:
docker compose up -d
  1. Ouvrir Airflow UI: http://localhost:8080
    • Déclencher le DAG producteur: weather_daily (planification @daily)

Outils & URLs

Organisation du code

  • DAGs: dags/
  • Job Spark Streaming: spark_stream.py
  • Entrypoint Airflow: script/entrypoint.sh
  • Dépendances Python: requirements.txt

Visualisation par Looker Studio

Le rapport se met automatiquement à jour chaque jour.

alt text

J’ai filtré les données à la date d’aujourd’hui.

alt text

Commentaire

C’est un pipeline ETL de type streaming, mais j’ai choisi de le mettre en batch en raison des caractéristiques des données. En effet, l’API Open-Meteo envoie des données une fois par jour. Ainsi, il n’est pas nécessaire d’utiliser Kafka, mais je souhaite l’employer dans un but d’apprentissage et bien comprendre le fonctionnement.

About

C’est un pipeline de bout en bout (end-to-end) dans le domaine de l’ingénierie des données (data engineering).

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages