ETL Designer
Workflow en étapes
1. Analyse des sources et destinations
- Inventorier chaque source : type (SGBD, API REST, fichier plat, Kafka, SaaS), volume moyen/max, fréquence de rafraîchissement, contraintes d'accès réseau, latence acceptable.
- Identifier la destination : data warehouse (Snowflake, BigQuery, Redshift, SQL Server DWH), data lake (S3, ADLS), data mart ou base opérationnelle.
- Consigner les SLA de fraîcheur (batch quotidien, quasi-temps-réel <5 min, streaming) et les fenêtres de maintenance source.
2. Choix ETL vs ELT
| Critère |
ETL |
ELT |
| Puissance de la destination |
faible (on-prem DWH) |
forte (Snowflake, BigQuery) |
| Données sensibles |
masquage/anonymisation tôt |
plus difficile à isoler |
| Flexibilité exploration |
faible |
forte (raw zone conservée) |
| Coût compute |
moteur ETL dédié |
warehouse paie la transformation |
Décision à documenter : justifier le choix dans un ADR ou un commentaire de pipeline.
3. Extraction
Full load — snapshot complet ; simple, utiliser quand la table source est petite (<1 M lignes) ou sans colonne de delta fiable.
Incrémental par timestamp/ID :
-- Extraction SQL incrémentale
SELECT *
FROM source_table
WHERE updated_at > :last_run_ts
AND updated_at <= :current_run_ts
ORDER BY updated_at;
CDC (Change Data Capture) — Debezium sur PostgreSQL/MySQL/SQL Server, Oracle GoldenGate, AWS DMS. Capturer INSERT/UPDATE/DELETE sans polling. Nécessite que le WAL ou le binlog soit activé.
Extraction API paginée (Python) :
def extract_all_pages(url, headers, page_size=200):
results, cursor = [], None
while True:
params = {"limit": page_size, **({"cursor": cursor} if cursor else {})}
r = requests.get(url, headers=headers, params=params, timeout=30)
r.raise_for_status()
data = r.json()
results.extend(data["items"])
cursor = data.get("next_cursor")
if not cursor:
break
return results
4. Transformation
Ordre recommandé :
- Nettoyage — trim, normalisation casse, formats dates (ISO 8601), suppression des caractères invalides.
- Déduplication — window function sur clé naturelle +
ROW_NUMBER().
- Enrichissement — lookup tables, API tiers (géocodage, scoring).
- Mapping de codes — table de correspondance versionnée en base ou fichier YAML.
- Calcul de KPIs — uniquement en fin de chaîne, sur données propres.
-- Déduplication avec window function
WITH ranked AS (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY updated_at DESC) AS rn
FROM staging.orders
)
INSERT INTO dw.orders
SELECT * EXCLUDE (rn) FROM ranked WHERE rn = 1;
5. Stratégie de chargement (Loading)
| Mode |
Quand l'utiliser |
Snippet |
| Full refresh |
petits référentiels, pas de SCD |
TRUNCATE + INSERT |
| Upsert/Merge |
tables transactionnelles, clé naturelle stable |
MERGE SQL ou INSERT … ON CONFLICT |
| SCD Type 1 |
on écrase, pas d'historique requis |
UPDATE direct |
| SCD Type 2 |
historisation obligatoire (audit, BI temps) |
colonnes valid_from, valid_to, is_current |
| Append-only |
événements immuables (logs, factures) |
INSERT pur |
-- MERGE Snowflake / SQL Server
MERGE INTO dw.customers AS tgt
USING staging.customers AS src
ON tgt.customer_id = src.customer_id
WHEN MATCHED THEN
UPDATE SET tgt.email = src.email, tgt.updated_at = src.updated_at
WHEN NOT MATCHED THEN
INSERT (customer_id, email, created_at, updated_at)
VALUES (src.customer_id, src.email, src.created_at, src.updated_at);
6. Error handling & idempotence
- Retry avec backoff exponentiel sur erreurs transitoires (timeout, rate-limit 429).
- Dead letter table pour les enregistrements invalides : stocker la ligne brute + message d'erreur + timestamp + job_id.
- Idempotence : chaque run sur le même intervalle temporel doit produire exactement le même résultat — utiliser une fenêtre
[start_ts, end_ts[ explicite passée en paramètre.
- Table de réconciliation : compter source vs destination après chaque chargement.
# Réconciliation source vs destination
def reconcile(src_count: int, dst_count: int, tolerance_pct: float = 0.01):
diff_pct = abs(src_count - dst_count) / max(src_count, 1)
if diff_pct > tolerance_pct:
raise ValueError(f"Reconciliation failed: src={src_count}, dst={dst_count}, diff={diff_pct:.1%}")
7. Performance
- Bulk loading : préférer
COPY INTO (Snowflake), LOAD DATA INFILE (MySQL), bcp ou BULK INSERT (SQL Server) aux INSERT ligne par ligne.
- Partitionner les fichiers intermédiaires (Parquet, ORC) par date ou entité avant chargement.
- Paralléliser l'extraction : diviser par partition, shard ou plage de dates ; utiliser
ThreadPoolExecutor ou Spark.
- Comprimer les fichiers en transit (gzip, snappy) pour réduire le temps de transfert réseau.
- Indexer uniquement les colonnes de jointure et de filtre sur la table destination — pas d'index sur toutes les colonnes.
# Chargement bulk Snowflake via connector
conn.cursor().execute(f"""
COPY INTO {schema}.{table}
FROM @my_stage/{filename}.parquet.gz
FILE_FORMAT = (TYPE=PARQUET)
PURGE = TRUE
""")
8. Orchestration & monitoring
Outils : Apache Airflow (DAG Python), Dagster (assets), dbt Cloud (scheduling dbt), Prefect, SSIS (on-prem legacy).
DAG Airflow minimal :
from airflow.decorators import dag, task
from pendulum import datetime
@dag(schedule="0 3 * * *", start_date=datetime(2026, 1, 1), catchup=False)
def orders_etl():
@task
def extract(): ...
@task
def transform(raw): ...
@task
def load(clean): ...
load(transform(extract()))
orders_etl()
Monitoring obligatoire :
- Durée d'exécution par étape (comparer à la baseline).
- Volumes extraits/chargés (alerte si ±20 % de la moyenne mobile).
- Data lineage tracé (outil : OpenLineage / Marquez, dbt docs).
- SLA de fraîcheur : alerter si la table destination n'est pas mise à jour dans la fenêtre attendue.
Garde-fous & anti-patterns
| Anti-pattern |
Problème |
Solution |
SELECT * en extraction sans schéma fixé |
rupture silencieuse à l'ajout de colonnes |
sélectionner les colonnes explicitement, versionner le schéma |
| Transformation dans la requête d'extraction |
couplage fort, difficile à tester |
séparer extraction brute et transformation |
| Chargement sans idempotence |
doublons à chaque retry |
fenêtre temporelle explicite + upsert ou truncate |
| Pas de dead letter |
perte silencieuse d'enregistrements invalides |
table d'erreurs systématique |
| ETL monolithique sans étapes atomiques |
rollback impossible en cas d'échec partiel |
découper en tâches redémarrables indépendamment |
| Secrets hard-codés dans le code |
fuite de credentials |
variables d'environnement ou vault (AWS Secrets Manager, Azure Key Vault) |
| Ignorer les fuseaux horaires |
décalages silencieux autour du changement d'heure |
stocker et transmettre en UTC, convertir en aval |
Bonnes pratiques 2026
- Contract-first : définir un contrat de schéma (JSON Schema, Avro, Protobuf) entre source et pipeline ; valider dès l'extraction.
- Data observability : intégrer un outil comme Great Expectations ou Soda pour des tests de qualité automatisés à chaque run.
- Metadata-driven ETL : piloter les pipelines par configuration (YAML/DB) plutôt que par duplication de code.
- Least privilege : chaque job utilise un compte dédié avec droits minimaux (lecture seule sur la source, écriture uniquement sur le schéma de staging).
- Gestion de la dette technique : documenter les workarounds source (champs mal typés, encodages exotiques) dans un fichier
QUIRKS.md versionné à côté du pipeline.
1---2name: etl-designer3description: Conception de processus ETL/ELT pour l'intégration de données. Se déclenche avec "ETL", "ELT", "extraction", "transformation", "load", "data integration", "data warehouse", "SSIS", "Talend", "Informatica". Also triggers on "ETL process", "ELT design", "data integration job".4---56# ETL Designer78## Workflow en étapes910### 1. Analyse des sources et destinations11- Inventorier chaque source : type (SGBD, API REST, fichier plat, Kafka, SaaS), volume moyen/max, fréquence de rafraîchissement, contraintes d'accès réseau, latence acceptable.12- Identifier la destination : data warehouse (Snowflake, BigQuery, Redshift, SQL Server DWH), data lake (S3, ADLS), data mart ou base opérationnelle.13- Consigner les SLA de fraîcheur (batch quotidien, quasi-temps-réel <5 min, streaming) et les fenêtres de maintenance source.1415### 2. Choix ETL vs ELT1617| Critère | ETL | ELT |18|---|---|---|19| Puissance de la destination | faible (on-prem DWH) | forte (Snowflake, BigQuery) |20| Données sensibles | masquage/anonymisation tôt | plus difficile à isoler |21| Flexibilité exploration | faible | forte (raw zone conservée) |22| Coût compute | moteur ETL dédié | warehouse paie la transformation |2324**Décision à documenter** : justifier le choix dans un ADR ou un commentaire de pipeline.2526### 3. Extraction2728**Full load** — snapshot complet ; simple, utiliser quand la table source est petite (<1 M lignes) ou sans colonne de delta fiable.2930**Incrémental par timestamp/ID** :31```sql32-- Extraction SQL incrémentale33SELECT *34FROM source_table35WHERE updated_at > :last_run_ts36 AND updated_at <= :current_run_ts37ORDER BY updated_at;38```3940**CDC (Change Data Capture)** — Debezium sur PostgreSQL/MySQL/SQL Server, Oracle GoldenGate, AWS DMS. Capturer INSERT/UPDATE/DELETE sans polling. Nécessite que le WAL ou le binlog soit activé.4142**Extraction API paginée (Python)** :43```python44def extract_all_pages(url, headers, page_size=200):45 results, cursor = [], None46 while True:47 params = {"limit": page_size, **({"cursor": cursor} if cursor else {})}48 r = requests.get(url, headers=headers, params=params, timeout=30)49 r.raise_for_status()50 data = r.json()51 results.extend(data["items"])52 cursor = data.get("next_cursor")53 if not cursor:54 break55 return results56```5758### 4. Transformation5960Ordre recommandé :611. **Nettoyage** — trim, normalisation casse, formats dates (ISO 8601), suppression des caractères invalides.622. **Déduplication** — window function sur clé naturelle + `ROW_NUMBER()`.633. **Enrichissement** — lookup tables, API tiers (géocodage, scoring).644. **Mapping de codes** — table de correspondance versionnée en base ou fichier YAML.655. **Calcul de KPIs** — uniquement en fin de chaîne, sur données propres.6667```sql68-- Déduplication avec window function69WITH ranked AS (70 SELECT *,71 ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY updated_at DESC) AS rn72 FROM staging.orders73)74INSERT INTO dw.orders75SELECT * EXCLUDE (rn) FROM ranked WHERE rn = 1;76```7778### 5. Stratégie de chargement (Loading)7980| Mode | Quand l'utiliser | Snippet |81|---|---|---|82| Full refresh | petits référentiels, pas de SCD | `TRUNCATE + INSERT` |83| Upsert/Merge | tables transactionnelles, clé naturelle stable | `MERGE` SQL ou `INSERT … ON CONFLICT` |84| SCD Type 1 | on écrase, pas d'historique requis | UPDATE direct |85| SCD Type 2 | historisation obligatoire (audit, BI temps) | colonnes `valid_from`, `valid_to`, `is_current` |86| Append-only | événements immuables (logs, factures) | `INSERT` pur |8788```sql89-- MERGE Snowflake / SQL Server90MERGE INTO dw.customers AS tgt91USING staging.customers AS src92 ON tgt.customer_id = src.customer_id93WHEN MATCHED THEN94 UPDATE SET tgt.email = src.email, tgt.updated_at = src.updated_at95WHEN NOT MATCHED THEN96 INSERT (customer_id, email, created_at, updated_at)97 VALUES (src.customer_id, src.email, src.created_at, src.updated_at);98```99100### 6. Error handling & idempotence101102- **Retry avec backoff exponentiel** sur erreurs transitoires (timeout, rate-limit 429).103- **Dead letter table** pour les enregistrements invalides : stocker la ligne brute + message d'erreur + timestamp + job_id.104- **Idempotence** : chaque run sur le même intervalle temporel doit produire exactement le même résultat — utiliser une fenêtre `[start_ts, end_ts[` explicite passée en paramètre.105- **Table de réconciliation** : compter source vs destination après chaque chargement.106107```python108# Réconciliation source vs destination109def reconcile(src_count: int, dst_count: int, tolerance_pct: float = 0.01):110 diff_pct = abs(src_count - dst_count) / max(src_count, 1)111 if diff_pct > tolerance_pct:112 raise ValueError(f"Reconciliation failed: src={src_count}, dst={dst_count}, diff={diff_pct:.1%}")113```114115### 7. Performance116117- **Bulk loading** : préférer `COPY INTO` (Snowflake), `LOAD DATA INFILE` (MySQL), `bcp` ou `BULK INSERT` (SQL Server) aux `INSERT` ligne par ligne.118- **Partitionner** les fichiers intermédiaires (Parquet, ORC) par date ou entité avant chargement.119- **Paralléliser** l'extraction : diviser par partition, shard ou plage de dates ; utiliser `ThreadPoolExecutor` ou Spark.120- **Comprimer** les fichiers en transit (gzip, snappy) pour réduire le temps de transfert réseau.121- **Indexer** uniquement les colonnes de jointure et de filtre sur la table destination — pas d'index sur toutes les colonnes.122123```python124# Chargement bulk Snowflake via connector125conn.cursor().execute(f"""126 COPY INTO {schema}.{table}127 FROM @my_stage/{filename}.parquet.gz128 FILE_FORMAT = (TYPE=PARQUET)129 PURGE = TRUE130""")131```132133### 8. Orchestration & monitoring134135**Outils** : Apache Airflow (DAG Python), Dagster (assets), dbt Cloud (scheduling dbt), Prefect, SSIS (on-prem legacy).136137**DAG Airflow minimal** :138```python139from airflow.decorators import dag, task140from pendulum import datetime141142@dag(schedule="0 3 * * *", start_date=datetime(2026, 1, 1), catchup=False)143def orders_etl():144 @task145 def extract(): ...146147 @task148 def transform(raw): ...149150 @task151 def load(clean): ...152153 load(transform(extract()))154155orders_etl()156```157158**Monitoring obligatoire** :159- Durée d'exécution par étape (comparer à la baseline).160- Volumes extraits/chargés (alerte si ±20 % de la moyenne mobile).161- Data lineage tracé (outil : OpenLineage / Marquez, dbt docs).162- SLA de fraîcheur : alerter si la table destination n'est pas mise à jour dans la fenêtre attendue.163164---165166## Garde-fous & anti-patterns167168| Anti-pattern | Problème | Solution |169|---|---|---|170| `SELECT *` en extraction sans schéma fixé | rupture silencieuse à l'ajout de colonnes | sélectionner les colonnes explicitement, versionner le schéma |171| Transformation dans la requête d'extraction | couplage fort, difficile à tester | séparer extraction brute et transformation |172| Chargement sans idempotence | doublons à chaque retry | fenêtre temporelle explicite + upsert ou truncate |173| Pas de dead letter | perte silencieuse d'enregistrements invalides | table d'erreurs systématique |174| ETL monolithique sans étapes atomiques | rollback impossible en cas d'échec partiel | découper en tâches redémarrables indépendamment |175| Secrets hard-codés dans le code | fuite de credentials | variables d'environnement ou vault (AWS Secrets Manager, Azure Key Vault) |176| Ignorer les fuseaux horaires | décalages silencieux autour du changement d'heure | stocker et transmettre en UTC, convertir en aval |177178## Bonnes pratiques 2026179180- **Contract-first** : définir un contrat de schéma (JSON Schema, Avro, Protobuf) entre source et pipeline ; valider dès l'extraction.181- **Data observability** : intégrer un outil comme Great Expectations ou Soda pour des tests de qualité automatisés à chaque run.182- **Metadata-driven ETL** : piloter les pipelines par configuration (YAML/DB) plutôt que par duplication de code.183- **Least privilege** : chaque job utilise un compte dédié avec droits minimaux (lecture seule sur la source, écriture uniquement sur le schéma de staging).184- **Gestion de la dette technique** : documenter les workarounds source (champs mal typés, encodages exotiques) dans un fichier `QUIRKS.md` versionné à côté du pipeline.