🌦️ Mini ETL Météo - Paris¶
Open-Meteo → PostgreSQL/PostGIS → Agrégats & Alertes → Carte¶
Ce notebook présente un pipeline de données de bout en bout, construit pour un exercice technique de Data Engineer :
- Extraction - récupération de données météo horaires via l'API Open-Meteo (température, vent, précipitations)
- Chargement - persistance des mesures brutes dans PostgreSQL, avec PostGIS pour la géométrie des points de mesure
- Transformation - agrégation des mesures à l'échelle de la ville, par fenêtres de 3 heures
- Alerting - génération d'alertes simples à partir de seuils météorologiques
- Visualisation - cartographie interactive des points de mesure
Stack technique : Python · SQLAlchemy (Core) · PostgreSQL/PostGIS · pandas · folium
Périmètre : ville de Paris, données du mois de février 2026.
1. Modèle de données¶
Le schéma repose sur 5 tables (voir datamodel.dbml) :
city ───< measure_point ───< device_measurement (mesures brutes, 1 par point / heure)
│
└────< aggregated_measurement (agrégats 3h, à l'échelle de la ville)
│
└────< alert (alertes générées sur les agrégats)
city/measure_pointstockent leur position via PostGIS (geometry).device_measurementcontient une ligne par point de mesure et par heure.aggregated_measurementagrège tous les points d'une ville sur une fenêtre de 3h (moyenne, min, max).alertréférence unaggregated_measurement(ou undevice_measurement) qui dépasse un seuil défini.
from sqlalchemy import create_engine, MetaData, Table, insert, text
import geoalchemy2 # requis pour que SQLAlchemy reconnaisse le type "geometry" de PostGIS
import pandas as pd
import requests
import folium
from datetime import datetime, timezone
2. Connexion à la base de données¶
Connexion via SQLAlchemy à la base PostgreSQL/PostGIS fournie par l'environnement Docker de l'exercice.
db_user = "weather_user"
db_password = "weather_password"
db_host = "localhost"
db_port = "5435"
db_name = "weather"
# Ceci est un exercice, ne pas le faire dans un vrai projet.
def connect_db(db_user: str, db_password: str, db_host: str, db_port: str, db_name: str):
"""Crée un engine SQLAlchemy vers la base PostgreSQL du projet."""
return create_engine(
f"postgresql+psycopg2://{db_user}:{db_password}@{db_host}:{db_port}/{db_name}"
)
conn = connect_db(db_user, db_password, db_host, db_port, db_name)
metadata_obj = MetaData()
print("Connexion établie :", conn)
Connexion établie : Engine(postgresql+psycopg2://weather_user:***@localhost:5435/weather)
3. Localisation des points de mesure (PostGIS)¶
Les positions sont stockées en geometry (PostGIS). On les convertit en latitude/longitude standard (SRID 4326, WGS84) avec ST_Transform, puis on extrait les coordonnées avec ST_X / ST_Y.
def locate_devices(conn) -> pd.DataFrame:
"""Récupère les points de mesure de Paris avec leurs coordonnées WGS84."""
query = """
SELECT
ST_X(ST_Transform(location, 4326)) AS long,
ST_Y(ST_Transform(location, 4326)) AS lat,
id AS measure_point_id,
city_id
FROM weather.measure_point
WHERE name LIKE 'Paris%%';
"""
return pd.read_sql_query(query, con=conn)
devices = locate_devices(conn)
print(f"{len(devices)} points de mesure trouvés pour Paris")
devices.head()
10 points de mesure trouvés pour Paris
| long | lat | measure_point_id | city_id | |
|---|---|---|---|---|
| 0 | 2.3488 | 48.8534 | b1eebc99-9c0b-4ef8-bb6d-6bb9bd380a01 | a0eebc99-9c0b-4ef8-bb6d-6bb9bd380a11 |
| 1 | 2.3590 | 48.8820 | b1eebc99-9c0b-4ef8-bb6d-6bb9bd380a02 | a0eebc99-9c0b-4ef8-bb6d-6bb9bd380a11 |
| 2 | 2.3412 | 48.8290 | b1eebc99-9c0b-4ef8-bb6d-6bb9bd380a03 | a0eebc99-9c0b-4ef8-bb6d-6bb9bd380a11 |
| 3 | 2.3960 | 48.8530 | b1eebc99-9c0b-4ef8-bb6d-6bb9bd380a04 | a0eebc99-9c0b-4ef8-bb6d-6bb9bd380a11 |
| 4 | 2.2930 | 48.8560 | b1eebc99-9c0b-4ef8-bb6d-6bb9bd380a05 | a0eebc99-9c0b-4ef8-bb6d-6bb9bd380a11 |
4. Récupération des données météo (API Open-Meteo)¶
Pour chaque point de mesure, on interroge l'API Open-Meteo (endpoint archive) sur la période du mois de février 2026, au pas horaire, pour la température, la vitesse du vent et les précipitations.
def api_meteo(latitude: float, longitude: float) -> dict:
"""Récupère les données météo horaires pour une position donnée (février 2026)."""
url = (
"https://archive-api.open-meteo.com/v1/archive"
f"?latitude={latitude}&longitude={longitude}"
"&start_date=2026-02-01&end_date=2026-02-28"
"&hourly=temperature_2m,wind_speed_10m,precipitation"
"&timezone=Europe%2FLondon"
)
return requests.get(url).json()
def fetch_all_meteo(conn) -> dict:
"""Interroge l'API pour tous les points de mesure et structure le résultat par point puis par heure."""
location_devices = locate_devices(conn)
result = {}
for _, point in location_devices.iterrows():
meteo = api_meteo(point["lat"], point["long"])
hourly = meteo["hourly"]
times = hourly["time"]
metrics = [k for k in hourly if k != "time"]
result[point["measure_point_id"]] = {
datetime.fromisoformat(t).replace(tzinfo=timezone.utc).isoformat(): {
k: hourly[k][i] for k in metrics
}
for i, t in enumerate(times)
}
return result
5. Stockage des mesures brutes¶
Chaque mesure horaire est insérée dans device_measurement. L'insertion se fait en masse (une liste de dicts passée à execute), bien plus efficace qu'une boucle de insert().values().
L'étape est rendue idempotente : si les mesures sont déjà présentes en base, on ne relance pas l'appel API ni l'insertion.
device_measurement_table = Table(
"device_measurement", metadata_obj, schema="weather", autoload_with=conn
)
existing_count = pd.read_sql_query(
"SELECT count(*) AS n FROM weather.device_measurement;", con=conn
)["n"][0]
if existing_count == 0:
raw_meteo = fetch_all_meteo(conn)
rows = [
{
"measure_point_id": mpid,
"date_observed": ts,
"temperature": vals.get("temperature_2m"),
"wind_speed": vals.get("wind_speed_10m"),
"precipitation": vals.get("precipitation"),
}
for mpid, series in raw_meteo.items()
for ts, vals in series.items()
]
with conn.begin() as connection:
connection.execute(insert(device_measurement_table), rows)
print(f"{len(rows)} mesures insérées dans device_measurement.")
else:
print(f"{existing_count} mesures déjà présentes dans device_measurement — insertion ignorée.")
6720 mesures déjà présentes dans device_measurement — insertion ignorée.
6. Agrégation à l'échelle de la ville (fenêtres de 3h)¶
Choix de conception : plutôt que de rappeler l'API pour ré-agréger, on agrège directement en SQL depuis device_measurement, déjà en base. C'est plus rapide, plus robuste (pas de dépendance réseau), et garantit la cohérence entre données brutes et agrégats.
La fonction Postgres date_bin('3 hours', colonne, origine) regroupe chaque timestamp dans sa fenêtre de 3h — l'équivalent SQL du (heure // 3) * 3 qu'on ferait en Python.
def aggregate_city_by_3h(conn) -> dict:
"""Agrège les mesures de device_measurement par ville et par fenêtre de 3h (UTC)."""
query = """
SELECT
mp.city_id,
date_bin('3 hours', dm.date_observed, TIMESTAMPTZ '2000-01-01') AS aggregated_at,
count(DISTINCT dm.measure_point_id) AS point_count,
avg(dm.temperature) AS temperature_avg,
min(dm.temperature) AS temperature_min,
max(dm.temperature) AS temperature_max,
avg(dm.wind_speed) AS wind_speed_avg,
min(dm.wind_speed) AS wind_speed_min,
max(dm.wind_speed) AS wind_speed_max,
sum(dm.precipitation) AS precipitation_total,
min(dm.precipitation) AS precipitation_min,
max(dm.precipitation) AS precipitation_max
FROM weather.device_measurement dm
JOIN weather.measure_point mp ON mp.id = dm.measure_point_id
GROUP BY mp.city_id, date_bin('3 hours', dm.date_observed, TIMESTAMPTZ '2000-01-01')
"""
df = pd.read_sql_query(query, con=conn)
result = {}
for row in df.itertuples():
bucket_key = row.aggregated_at.isoformat()
result.setdefault(row.city_id, {})[bucket_key] = {
"point_count": int(row.point_count),
"temperature_avg": round(row.temperature_avg, 3) if row.temperature_avg is not None else None,
"temperature_min": round(row.temperature_min, 3) if row.temperature_min is not None else None,
"temperature_max": round(row.temperature_max, 3) if row.temperature_max is not None else None,
"wind_speed_avg": round(row.wind_speed_avg, 3) if row.wind_speed_avg is not None else None,
"wind_speed_min": round(row.wind_speed_min, 3) if row.wind_speed_min is not None else None,
"wind_speed_max": round(row.wind_speed_max, 3) if row.wind_speed_max is not None else None,
"precipitation_total": round(row.precipitation_total, 3) if row.precipitation_total is not None else None,
"precipitation_min": round(row.precipitation_min, 3) if row.precipitation_min is not None else None,
"precipitation_max": round(row.precipitation_max, 3) if row.precipitation_max is not None else None,
}
return result
data_aggr = aggregate_city_by_3h(conn)
n_buckets = sum(len(buckets) for buckets in data_aggr.values())
print(f"{len(data_aggr)} ville(s), {n_buckets} fenêtres de 3h agrégées")
1 ville(s), 224 fenêtres de 3h agrégées
7. Persistance des agrégats¶
On réinsère systématiquement (DELETE puis INSERT) pour que cette cellule soit rejouable sans erreur ni doublon — la contrainte unique (city_id, aggregation_window, aggregated_at) garantit qu'on ne peut pas avoir deux agrégats pour la même fenêtre.
aggregated_measurement_table = Table(
"aggregated_measurement", metadata_obj, schema="weather", autoload_with=conn
)
rows = [
{
"city_id": city_id,
"aggregation_window": "three_hours",
"aggregated_at": ts,
**vals,
}
for city_id, series in data_aggr.items()
for ts, vals in series.items()
]
with conn.begin() as connection:
# `alert` référence `aggregated_measurement` par clé étrangère : on vide d'abord les alertes
# (elles seront de toute façon régénérées à la section 8) avant de purger les agrégats.
connection.execute(text("DELETE FROM weather.alert;"))
connection.execute(text("DELETE FROM weather.aggregated_measurement WHERE aggregation_window = 'three_hours';"))
connection.execute(insert(aggregated_measurement_table), rows)
print(f"{len(rows)} agrégats (ré)insérés dans aggregated_measurement.")
224 agrégats (ré)insérés dans aggregated_measurement.
8. Génération d'alertes¶
Méthodologie de seuillage : plutôt que des seuils arbitraires, ils ont été choisis en observant la distribution réelle des données (top 20 des valeurs les plus hautes par métrique), pour rester cohérents avec le climat parisien de février 2026.
| Métrique | Seuil | Type d'alerte |
|---|---|---|
| Température moyenne (3h) | ≥ 12.0 °C | high_temperature |
| Vitesse du vent moyenne (3h) | ≥ 19.0 km/h | high_wind_speed |
| Précipitations cumulées (3h) | ≥ 20.0 mm | high_precipitation |
La logique de détection est isolée dans une fonction pure (generate_alerts), sans effet de bord ni dépendance à la base — ce qui la rend directement testable (voir section 10).
THRESHOLDS = {
"temperature_avg": {"threshold": 12.0, "alert_type": "high_temperature", "unit": "°C"},
"wind_speed_avg": {"threshold": 19.0, "alert_type": "high_wind_speed", "unit": "km/h"},
"precipitation_total": {"threshold": 20.0, "alert_type": "high_precipitation", "unit": "mm"},
}
def generate_alerts(data_aggr: dict, id_lookup: dict, thresholds: dict) -> list[dict]:
"""
Compare chaque agrégat 3h aux seuils définis et retourne la liste des
lignes à insérer dans `alert`. Fonction pure : aucun accès base de données.
"""
alert_rows = []
for city_id, timestamps in data_aggr.items():
for ts, valeurs in timestamps.items():
agg_id = id_lookup.get((city_id, ts))
if agg_id is None:
continue
for metric, config in thresholds.items():
observed = valeurs.get(metric)
if observed is not None and observed >= config["threshold"]:
alert_rows.append({
"source_level": "aggregated",
"alert_type": config["alert_type"],
"aggregated_measurement_id": agg_id,
"severity": "high",
"observed_value": observed,
"threshold_value": config["threshold"],
"unit": config["unit"],
"date_issued": datetime.now(timezone.utc),
})
return alert_rows
alert_table = Table("alert", metadata_obj, schema="weather", autoload_with=conn)
ids_df = pd.read_sql_query(
"SELECT id, city_id, aggregated_at FROM weather.aggregated_measurement "
"WHERE aggregation_window = 'three_hours';",
con=conn,
)
id_lookup = {(row.city_id, row.aggregated_at.isoformat()): row.id for row in ids_df.itertuples()}
alert_rows = generate_alerts(data_aggr, id_lookup, THRESHOLDS)
with conn.begin() as connection:
connection.execute(text("DELETE FROM weather.alert;"))
if alert_rows:
connection.execute(insert(alert_table), alert_rows)
print(f"{len(alert_rows)} alertes générées.")
pd.read_sql_query("SELECT alert_type, count(*) FROM weather.alert GROUP BY alert_type ORDER BY 2 DESC;", con=conn)
73 alertes générées.
| alert_type | count | |
|---|---|---|
| 0 | high_temperature | 33 |
| 1 | high_wind_speed | 20 |
| 2 | high_precipitation | 20 |
9. Visualisation cartographique¶
Carte interactive (Leaflet, via folium) affichant chaque point de mesure avec sa dernière mesure connue. La dernière mesure par point est récupérée avec DISTINCT ON, une syntaxe Postgres qui garde une seule ligne par groupe selon l'ORDER BY donné.
query_latest = """
SELECT DISTINCT ON (dm.measure_point_id)
dm.measure_point_id,
dm.date_observed,
dm.temperature,
dm.wind_speed,
dm.precipitation,
ST_X(ST_Transform(mp.location, 4326)) AS long,
ST_Y(ST_Transform(mp.location, 4326)) AS lat
FROM weather.device_measurement dm
JOIN weather.measure_point mp ON mp.id = dm.measure_point_id
ORDER BY dm.measure_point_id, dm.date_observed DESC;
"""
latest = pd.read_sql_query(query_latest, con=conn)
carte = folium.Map(location=[latest["lat"].mean(), latest["long"].mean()], zoom_start=12, tiles="cartodbpositron")
for _, row in latest.iterrows():
popup_html = f"""
<b>Point :</b> {row['measure_point_id']}<br>
<b>Dernière mesure :</b> {row['date_observed']}<br>
<b>Température :</b> {row['temperature']} °C<br>
<b>Vent :</b> {row['wind_speed']} km/h<br>
<b>Précipitation :</b> {row['precipitation']} mm
"""
folium.CircleMarker(
location=[row["lat"], row["long"]],
radius=7,
color="#2a6f97",
fill=True,
fill_opacity=0.85,
popup=folium.Popup(popup_html, max_width=250),
).add_to(carte)
carte
10. Tests unitaires¶
La logique métier d'alerting (generate_alerts) est une fonction pure : pas d'appel réseau, pas d'accès base de données. Elle peut donc être testée avec des données factices, sans dépendre d'une infrastructure — c'est le principal intérêt d'avoir isolé cette logique du reste du pipeline.
(Dans un vrai projet, ces tests vivraient dans tests/test_alerts.py et seraient lancés via pytest. Ils sont inclus ici pour la démonstration.)
def test_generate_alerts_detects_threshold_breach():
data = {"city-1": {"2026-02-01T00:00:00+00:00": {"temperature_avg": 20.0}}}
lookup = {("city-1", "2026-02-01T00:00:00+00:00"): "agg-1"}
thresholds = {"temperature_avg": {"threshold": 12.0, "alert_type": "high_temperature", "unit": "°C"}}
alerts = generate_alerts(data, lookup, thresholds)
assert len(alerts) == 1
assert alerts[0]["alert_type"] == "high_temperature"
assert alerts[0]["observed_value"] == 20.0
assert alerts[0]["aggregated_measurement_id"] == "agg-1"
def test_generate_alerts_respects_threshold():
data = {"city-1": {"2026-02-01T00:00:00+00:00": {"temperature_avg": 5.0}}}
lookup = {("city-1", "2026-02-01T00:00:00+00:00"): "agg-1"}
thresholds = {"temperature_avg": {"threshold": 12.0, "alert_type": "high_temperature", "unit": "°C"}}
alerts = generate_alerts(data, lookup, thresholds)
assert alerts == []
def test_generate_alerts_skips_missing_lookup():
data = {"city-1": {"2026-02-01T00:00:00+00:00": {"temperature_avg": 20.0}}}
thresholds = {"temperature_avg": {"threshold": 12.0, "alert_type": "high_temperature", "unit": "°C"}}
alerts = generate_alerts(data, {}, thresholds)
assert alerts == []
def test_generate_alerts_handles_multiple_metrics():
data = {"city-1": {"2026-02-01T00:00:00+00:00": {"temperature_avg": 20.0, "wind_speed_avg": 25.0}}}
lookup = {("city-1", "2026-02-01T00:00:00+00:00"): "agg-1"}
thresholds = {
"temperature_avg": {"threshold": 12.0, "alert_type": "high_temperature", "unit": "°C"},
"wind_speed_avg": {"threshold": 19.0, "alert_type": "high_wind_speed", "unit": "km/h"},
}
alerts = generate_alerts(data, lookup, thresholds)
assert len(alerts) == 2
assert {a["alert_type"] for a in alerts} == {"high_temperature", "high_wind_speed"}
test_generate_alerts_detects_threshold_breach()
test_generate_alerts_respects_threshold()
test_generate_alerts_skips_missing_lookup()
test_generate_alerts_handles_multiple_metrics()
print("Tous les tests sont passés ✅")
Tous les tests sont passés ✅
Conclusion¶
Ce pipeline couvre l'ensemble du cycle de vie de la donnée météo : extraction (API), chargement (PostgreSQL/PostGIS), transformation (agrégation SQL), détection (alertes) et restitution (carte interactive).
Points de conception à retenir :
- Agrégation effectuée côté base (SQL) plutôt que par rappel d'API, pour la robustesse et la performance.
- Étapes d'insertion idempotentes (delete + insert) pour permettre de rejouer le notebook sans erreur.
- Logique d'alerting extraite en fonction pure, testable indépendamment de l'infrastructure.
- Seuils d'alerte choisis à partir de l'observation des données réelles, pas arbitrairement.