|
| 1 | +# Architecture — Pipeline Homepedia |
| 2 | + |
| 3 | +## Contexte |
| 4 | + |
| 5 | +Stack : **DuckDB** comme moteur unique (ingestion + spatial + scoring), **GCS** comme stockage objet, **Cloud Composer (Airflow)** comme orchestrateur. Cycle annuel, ~1,3 Go de sources brutes, 182 Mo de résultat. |
| 6 | + |
| 7 | +Le code `exploration/src/` est la référence d'implémentation — les jobs sont des adaptations directes des fonctions existantes, avec GCS en entrée/sortie au lieu du système de fichiers local. |
| 8 | + |
| 9 | +--- |
| 10 | + |
| 11 | +## DAG Airflow — `homepedia_pipeline` |
| 12 | + |
| 13 | +``` |
| 14 | +@yearly ───────────────────────────────────────────────────────────────────────────────────────────── |
| 15 | +
|
| 16 | + ┌─────────────────────────────── GROUPE : ingest (parallèle) ──────────────────────────────────┐ |
| 17 | + │ │ |
| 18 | + │ ingest_dvf ingest_geometries ingest_transport ingest_climat │ |
| 19 | + │ (DVF CSV.gz (GeoJSON Etalab (GTFS national (GeoJSON stations │ |
| 20 | + │ → bronze/dvf/) → bronze/geom/) → bronze/ MF + parsing .data │ |
| 21 | + │ transport/) → bronze/climat/) │ |
| 22 | + │ │ |
| 23 | + │ ┌──────────────────── SOUS-GROUPE : sources communales (parallèle) ──────────────────────┐ │ |
| 24 | + │ │ │ │ |
| 25 | + │ │ ingest_emploi ingest_securite ingest_tourisme ingest_bpe ingest_revenus │ │ |
| 26 | + │ │ ingest_risques ingest_dpe ingest_proximite │ │ |
| 27 | + │ │ (CSV/ZIP → bronze/emploi/, bronze/securite/, …) │ │ |
| 28 | + │ └─────────────────────────────────────────────────────────────────────────────────────────┘ │ |
| 29 | + └───────────────────────────────────────────────────────────────────────────────────────────────┘ |
| 30 | + │ |
| 31 | + ▼ |
| 32 | + ┌─────────────────────────── GROUPE : preprocess (séquentiel partiel) ──────────────────────────┐ |
| 33 | + │ │ |
| 34 | + │ preprocess_geometries │ |
| 35 | + │ (clean_geometries → silver/communes_geom/) │ |
| 36 | + │ │ │ |
| 37 | + │ ▼ │ |
| 38 | + │ preprocess_dvf ──► preprocess_commune_agg │ |
| 39 | + │ (clean_dvf (agrégat médiane/p25/p75 │ |
| 40 | + │ → silver/ par commune │ |
| 41 | + │ dvf_clean/) → silver/commune_agg/) │ |
| 42 | + │ │ |
| 43 | + │ preprocess_transport ← dépend de preprocess_geometries (ST_Contains) │ |
| 44 | + │ (clean_transport_stops + jointure spatiale → silver/transport_commune/) │ |
| 45 | + │ │ |
| 46 | + │ preprocess_climat ← dépend de preprocess_geometries (CROSS JOIN Haversine) │ |
| 47 | + │ (clean_climat → silver/climat_commune/) │ |
| 48 | + │ │ |
| 49 | + │ preprocess_proximite ← dépend de preprocess_geometries (CROSS JOIN Haversine) │ |
| 50 | + │ (clean_proximite_metropole → silver/proximite_commune/) │ |
| 51 | + │ │ |
| 52 | + │ ┌──── sources communales (parallèle, dépendent de preprocess_geometries) ────────────────┐ │ |
| 53 | + │ │ preprocess_emploi preprocess_securite preprocess_tourisme │ │ |
| 54 | + │ │ preprocess_bpe preprocess_revenus preprocess_risques preprocess_dpe │ │ |
| 55 | + │ │ (clean_* → silver/<source>_commune/) │ │ |
| 56 | + │ └────────────────────────────────────────────────────────────────────────────────────────┘ │ |
| 57 | + └────────────────────────────────────────────────────────────────────────────────────────────────┘ |
| 58 | + │ |
| 59 | + ▼ |
| 60 | + ┌──────────────────────────────── validate_silver ───────────────────────────────────────────────┐ |
| 61 | + │ quality.validate() + quality.coverage() sur toutes les tables silver │ |
| 62 | + │ → gold/dq_reports/silver_<run_date>.json │ |
| 63 | + │ Échoue le DAG (task FAILED) si une règle critique est KO │ |
| 64 | + └────────────────────────────────────────────────────────────────────────────────────────────────┘ |
| 65 | + │ |
| 66 | + ▼ |
| 67 | + ┌──────────────────────────────────── score ─────────────────────────────────────────────────────┐ |
| 68 | + │ Fusionne les 9 tables silver (commune_agg + 8 dimensions) │ |
| 69 | + │ Normalisation _norm() p01–p99, composer_score(), gap, gap_pondere │ |
| 70 | + │ → gold/score_territoire/run_date=YYYY/score.parquet │ |
| 71 | + └────────────────────────────────────────────────────────────────────────────────────────────────┘ |
| 72 | + │ |
| 73 | + ▼ |
| 74 | + ┌──────────────────────────────── validate_gold ─────────────────────────────────────────────────┐ |
| 75 | + │ Contrôles sur le gold : nb communes scorées, gap dans [-1, 1], │ |
| 76 | + │ pas de commune avec score NULL, Top 25 stable vs run précédent (τ Kendall > 0.8) │ |
| 77 | + │ → gold/dq_reports/gold_<run_date>.json │ |
| 78 | + └────────────────────────────────────────────────────────────────────────────────────────────────┘ |
| 79 | + │ |
| 80 | + ▼ |
| 81 | + ┌──────────────────────────────── publish ───────────────────────────────────────────────────────┐ |
| 82 | + │ Copie gold/score_territoire/run_date=YYYY/ → gold/score_territoire/latest/ │ |
| 83 | + │ (GCSToGCSOperator — l'API FastAPI lit toujours le chemin /latest/) │ |
| 84 | + └────────────────────────────────────────────────────────────────────────────────────────────────┘ |
| 85 | +``` |
| 86 | + |
| 87 | +--- |
| 88 | + |
| 89 | +## Détail des tâches Airflow |
| 90 | + |
| 91 | +| Tâche | Opérateur Airflow | Code source réutilisé | Sortie GCS | |
| 92 | +|---|---|---|---| |
| 93 | +| `ingest_dvf` | `PythonOperator` | `ingest.ensure_dvf` adapté | `bronze/dvf/year={{year}}/` | |
| 94 | +| `ingest_geometries` | `PythonOperator` | `ingest.ensure_geometries` | `bronze/geom/` | |
| 95 | +| `ingest_transport` | `PythonOperator` | `ingest.ensure_transport` (partie download) | `bronze/transport/` | |
| 96 | +| `ingest_climat` | `PythonOperator` | `ingest_extra._build_stations_csv` | `bronze/climat/` | |
| 97 | +| `ingest_<source>` (×7) | `PythonOperator` | `ingest_extra.ensure_*` (partie download) | `bronze/<source>/` | |
| 98 | +| `preprocess_geometries` | `PythonOperator` | `preprocess.clean_geometries` | `silver/communes_geom/` | |
| 99 | +| `preprocess_dvf` | `PythonOperator` | `preprocess.clean_dvf` | `silver/dvf_clean/` | |
| 100 | +| `preprocess_commune_agg` | `PythonOperator` | `ingest.ensure_commune_agg` | `silver/commune_agg/` | |
| 101 | +| `preprocess_transport` | `PythonOperator` | `preprocess.clean_transport_stops` + `ST_Contains` | `silver/transport_commune/` | |
| 102 | +| `preprocess_climat` | `PythonOperator` | `preprocess.clean_climat` | `silver/climat_commune/` | |
| 103 | +| `preprocess_proximite` | `PythonOperator` | `preprocess.clean_proximite_metropole` | `silver/proximite_commune/` | |
| 104 | +| `preprocess_<source>` (×7) | `PythonOperator` | `preprocess.clean_*` | `silver/<source>_commune/` | |
| 105 | +| `validate_silver` | `PythonOperator` | `quality.validate`, `quality.coverage` | `gold/dq_reports/` | |
| 106 | +| `score` | `PythonOperator` | `section_score()` de `exploration/notebooks/exploration.py` | `gold/score_territoire/` | |
| 107 | +| `validate_gold` | `PythonOperator` | `quality.validate` + règles gold custom | `gold/dq_reports/` | |
| 108 | +| `publish` | `GCSToGCSOperator` | — | `gold/score_territoire/latest/` | |
| 109 | + |
| 110 | +--- |
| 111 | + |
| 112 | +## Structure GCS |
| 113 | + |
| 114 | +``` |
| 115 | +gs://homepedia-data/ |
| 116 | +├── bronze/ |
| 117 | +│ ├── dvf/year=2024/dept=75/ … (partitionné par année et département) |
| 118 | +│ ├── geom/communes.parquet |
| 119 | +│ ├── transport/arrets.parquet |
| 120 | +│ ├── climat/stations.parquet |
| 121 | +│ ├── emploi/year=2021/ |
| 122 | +│ ├── securite/year=2022/ |
| 123 | +│ ├── tourisme/year=2021/ |
| 124 | +│ ├── bpe/year=2024/ |
| 125 | +│ ├── revenus/year=2021/ |
| 126 | +│ ├── risques/ |
| 127 | +│ └── dpe/ |
| 128 | +├── silver/ |
| 129 | +│ ├── communes_geom/ (GeoParquet) |
| 130 | +│ ├── dvf_clean/year=2024/ |
| 131 | +│ ├── commune_agg/year=2024/ |
| 132 | +│ ├── transport_commune/ |
| 133 | +│ ├── climat_commune/ |
| 134 | +│ ├── proximite_commune/ |
| 135 | +│ ├── emploi_commune/ |
| 136 | +│ ├── securite_commune/ |
| 137 | +│ ├── tourisme_commune/ |
| 138 | +│ ├── equipements_commune/ |
| 139 | +│ ├── revenus_commune/ |
| 140 | +│ ├── risques_commune/ |
| 141 | +│ └── dpe_commune/ |
| 142 | +└── gold/ |
| 143 | + ├── score_territoire/ |
| 144 | + │ ├── run_date=2025-01-15/score.parquet |
| 145 | + │ └── latest/score.parquet ← lu par l'API FastAPI |
| 146 | + └── dq_reports/ |
| 147 | + ├── silver_2025-01-15.json |
| 148 | + └── gold_2025-01-15.json |
| 149 | +``` |
| 150 | + |
| 151 | +--- |
| 152 | + |
| 153 | +## Décisions arrêtées |
| 154 | + |
| 155 | +> Les décisions d'architecture sont désormais tracées dans [`adr/`](adr/) |
| 156 | +> (norme MADR). En cas de divergence avec ce document, les ADR font foi. |
| 157 | +
|
| 158 | +| Décision | Choix | |
| 159 | +|---|---| |
| 160 | +| Moteur | **DuckDB** seul (pas de Spark) — 1,3 Go tient en RAM, DuckDB natif sur GCS ([ADR-0001](adr/0001-duckdb-moteur-unique.md)) | |
| 161 | +| Stockage | **GCS** `gs://homepedia-data/` ([ADR-0002](adr/0002-stockage-gcs-medaillon-parquet.md)) | |
| 162 | +| Orchestration | ~~Cloud Composer (Airflow)~~ **Remplacée** : Cloud Run Jobs + Workflows + Scheduler ([ADR-0009](adr/0009-orchestration-cloud-run-jobs.md)) | |
| 163 | +| Format | **Parquet** (GeoParquet pour `communes_geom`) ([ADR-0002](adr/0002-stockage-gcs-medaillon-parquet.md)) | |
| 164 | +| API | Lit `gold/score_territoire/latest/` via `duckdb.read_parquet('gs://...')` | |
| 165 | +| Ré-exécution partielle | Via paramètre `year` passé au job DVF (équivalent du `dag_run.conf` initialement prévu) | |
0 commit comments