Explorer
KNOW-PAT-233

ETL & Analytics Pipeline Design — Ingestion, transformation, warehousing, observability

Domaine
data
Type
pattern
Priorité
P1

ETL & Analytics Pipeline Design

Problème

Les pipelines data sont construits au fil de l'eau : pas de schéma, pas de monitoring, transformations en SQL inline, pas de tests. Résultat : données silencieusement fausses, pipelines qui cassent sans alerte.

Solution

Architecture ETL/ELT en couches + tests de qualité + observability + idempotence.

1. ETL vs ELT — choisir

Approche Quand Outils
ETL (Extract → Transform → Load) Source limitée, transformations complexes, data warehouse léger Airflow + Pandas + SQL
ELT (Extract → Load → Transform) Data warehouse puissant (BigQuery, Snowflake), transformations en SQL dbt + Fivetran/Airbyte

Tendance 2026 : ELT avec dbt (transformations versionnées, testables, documentées).

2. Architecture en couches (medallion)

Raw (Bronze)  →  Staging (Silver)  →  Mart (Gold)
─────────────    ─────────────────    ───────────
Données brutes   Nettoyées + typées   Modélisées métier
Pas de filtre    Dédoublonnées        Aggrégées, prêtes BI
Toutes sources   Conformées           Star schema

dbt — implémentation ELT

-- models/staging/stg_orders.sql
SELECT
  order_id::bigint AS order_id,
  customer_id::bigint AS customer_id,
  order_date::date AS order_date,
  amount::numeric(19,4) AS amount,
  status::varchar AS status
FROM {{ source('raw', 'orders') }}
WHERE order_date >= '2026-01-01'

-- models/mart/fct_revenue.sql
SELECT
  d.date,
  d.month,
  SUM(f.amount) AS revenue,
  COUNT(DISTINCT f.customer_id) AS unique_customers
FROM {{ ref('fct_orders') }} f
JOIN {{ ref('dim_date') }} d ON f.order_date = d.date
GROUP BY 1, 2

3. Ingestion — patterns

# Airflow DAG — ingestion quotidienne
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta

default_args = {
    'owner': 'data',
    'retries': 3,
    'retry_delay': timedelta(minutes=5),
    'email_on_failure': True,
}

dag = DAG('daily_ingestion', default_args=default_args,
          schedule='0 2 * * *', start_date=datetime(2026, 1, 1),
          catchup=False)

def extract_api(**context):
    ds = context['ds']  # date d'exécution
    data = requests.get(f'https://api.example.com/orders?date={ds}').json()
    # Écrire en zone raw (append, pas overwrite)
    write_to_raw(data, table='orders', partition=ds)

def transform_staging(**context):
    ds = context['ds']
    run_dbt_models(select=f'staging', vars={'execution_date': ds})

extract = PythonOperator(task_id='extract_api', python_callable=extract_api, dag=dag)
transform = PythonOperator(task_id='transform_staging', python_callable=transform_staging, dag=dag)
extract >> transform

4. Tests de qualité data

# dbt — tests sur les modèles
# models/staging/schema.yml
version: 2
models:
  - name: stg_orders
    columns:
      - name: order_id
        tests:
          - unique
          - not_null
      - name: amount
        tests:
          - not_null
          - dbt_utils.accepted_range:
              min_value: 0
              max_value: 1000000
      - name: status
        tests:
          - accepted_values:
              values: ['pending', 'paid', 'shipped', 'cancelled']

5. Idempotence — règle critique

-- ❌ Mal : append sans déduplication
INSERT INTO orders SELECT * FROM raw_orders WHERE date = '2026-07-21';

-- ✅ Bien : merge/upsert (idempotent)
MERGE INTO orders AS target
USING (SELECT * FROM raw_orders WHERE date = '2026-07-21') AS source
ON target.order_id = source.order_id
WHEN MATCHED THEN UPDATE SET amount = source.amount, status = source.status
WHEN NOT MATCHED THEN INSERT (order_id, customer_id, order_date, amount, status)
  VALUES (source.order_id, source.customer_id, source.order_date, source.amount, source.status);
  • Un pipeline doit pouvoir être re-joué sans dupliquer ni corrompre les données
  • Partionnement par date → re-processer une partition sans toucher les autres

6. Observability

Métrique Outil Alert si
Freshness (dernier run) Airflow, dbt > 24h sans run
Volume (lignes ingérées) Custom test > 50% de variation vs J-1
Null rate dbt test > 5% sur colonne critique
Latence pipeline Airflow > 2x durée normale
Échec tâche Airflow Any failure
# Airflow — alerting Slack
from airflow.providers.slack.operators.slack import SlackAPIPostOperator

alert = SlackAPIPostOperator(
    task_id='slack_alert',
    slack_conn_id='slack',
    channel='#data-alerts',
    text=f'Pipeline {{dag.dag_id}} failed on {{ds}}',
    trigger_rule='one_failed',  # se déclenche si une tâche échoue
)
extract >> transform >> alert

7. Modélisation — star schema

           dim_date
              │
dim_customer ─┼─ fct_orders ─ dim_product
              │
           dim_store
  • Tables de faits : événements mesurables (orders, clicks, payments)
  • Tables de dimensions : entités descriptives (customers, products, dates)
  • Jamais de jointures dans le mart → pré-calculer les aggrégats

Anti-patterns

  • Pas de tests de qualité (données silencieusement fausses)
  • Pipeline non-idempotent (duplications au re-run)
  • Transformations en SQL inline non versionné
  • Pas de monitoring (pipeline cassé pendant des jours)
  • Append sans déduplication
  • Une seule grosse table au lieu d'un star schema
  • Pas de partitionnement (full scan à chaque requête)

Références

  • [[KNOW-PAT-227]] — Database Schema Design (types, indexing)
  • [[KNOW-PAT-229]] — Caching Strategies (cache des aggrégats)
  • [[KNOW-PAT-230]] — CI/CD Pipeline (même logique pour data pipelines)
  • [[KNOW-REF-014]] — System Design Primer