Exécuter des charges SQL

Créer des pipelines de données avec Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Pourquoi SQL dans Airflow ?

 

  • Les charges SQL sont l'usage Airflow le plus courant
  • Laissez la base de données faire le gros du travail
  • Airflow orchestre le quand et le
  • La base de données gère le comment

Orchestration SQL avec Airflow

Créer des pipelines de données avec Airflow

SQLExecuteQueryOperator

from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator

SQLExecuteQueryOperator( task_id="load_sales", conn_id="duckdb_analytics", sql="INSERT INTO sales SELECT * FROM staging", )

 

  • Fonctionne avec toute base de données ayant un fournisseur Airflow compatible
  • PostgreSQL, Snowflake, BigQuery, DuckDB, et plus encore
Créer des pipelines de données avec Airflow

Connexions

  • Stockez les informations d'identification hors du code
  • Chacune a un ID, type, hôte, port, identifiant
  • Créées de multiples façons : interface, CLI, API ou variables d'environnement
  • Référencez-les par conn_id dans les opérateurs

Interface des connexions Airflow

Créer des pipelines de données avec Airflow

Créer un pipeline SQL

from airflow.sdk import dag, task
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator

@dag(schedule="@daily")
def sales_pipeline():

    aggregate = SQLExecuteQueryOperator(
        task_id="aggregate_daily_sales",
        conn_id="duckdb_analytics",
        sql="""
            INSERT INTO daily_summary (order_date, total_orders, total_revenue)
            SELECT order_date, COUNT(*), SUM(revenue)
            FROM raw_orders GROUP BY order_date
        """,
    )
Créer des pipelines de données avec Airflow

Fichiers SQL externes

@dag(
    schedule="@daily",
    template_searchpath="/path/to/include/sql",
)
def sales_pipeline():
    aggregate = SQLExecuteQueryOperator(
        task_id="aggregate_daily_sales",
        conn_id="duckdb_analytics",
        sql="aggregate_sales.sql",
    )
  • Définissez template_searchpath sur @dag
  • Référencez le nom de fichier dans sql
  • Gardez les scripts et la logique d'affaires hors du dossier dags/ 💡

Structure des fichiers du projet

Créer des pipelines de données avec Airflow

Gabarits Jinja dans les fichiers SQL

DELETE FROM daily_summary WHERE order_date = '{{ ds }}';

INSERT INTO daily_summary (order_date, total_orders, total_revenue)
SELECT order_date, COUNT(*), SUM(revenue)
FROM raw_orders
WHERE order_date = '{{ ds }}'
GROUP BY order_date;

 

  • {{ ds }} rend la date logique (YYYY-MM-DD)
  • DELETE puis INSERT : le motif idempotent du chapitre 2
  • Relancer la même date produit le même résultat
Créer des pipelines de données avec Airflow

params vs parameters

params (rendu Jinja)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • Valeur interpolée dans la chaîne SQL
  • Vulnérable à l'injection SQL

parameters (liaison au niveau BD)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • Valeur transmise au pilote de base de données
  • À l'épreuve des injections : le pilote gère l'échappement
Créer des pipelines de données avec Airflow

Du local à la production avec Astro

$$

  • Astronomer's Astro : plateforme Airflow gérée

$$

  • Astro CLI : Airflow local en une commande

$$

  • Développez localement, déployez en production sans heurts

Astro CLI et plateforme

Créer des pipelines de données avec Airflow

Passons à la pratique !

Créer des pipelines de données avec Airflow

Preparing Video For Download...