執行 SQL 工作負載

使用 Airflow 建置資料管線

Volker Janz

Senior Developer Advocate at Astronomer

為什麼在 Airflow 用 SQL?

 

  • SQL 工作負載是最常見的 Airflow 用例
  • 資料庫負責繁重運算
  • Airflow 協調何時在哪裡
  • 資料庫負責如何執行

Airflow 中的 SQL 編排

使用 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", )

 

  • 適用於具相容 Airflow provider 的任何資料庫
  • PostgreSQL、Snowflake、BigQuery、DuckDB 等
使用 Airflow 建置資料管線

Connections

  • 憑證存放在程式碼之外
  • 每個連線都有 ID、type、host、port、login
  • 可用多種方式建立:UI、CLI、API,或環境變數
  • 在 operators 以 conn_id 參照

Airflow Connections 介面

使用 Airflow 建置資料管線

建立 SQL pipeline

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
        """,
    )
使用 Airflow 建置資料管線

外部 SQL 檔案

@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",
    )
  • @dag 設定 template_searchpath
  • sql 參照檔名
  • 將腳本與商業邏輯放在 dags/ 資料夾之外 💡

專案檔案結構

使用 Airflow 建置資料管線

SQL 檔中的 Jinja 範本

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 }} 會渲染為邏輯日期YYYY-MM-DD
  • 先 DELETE 後 INSERT:第 2 章的冪等模式
  • 針對同一日期重跑會得到相同結果
使用 Airflow 建置資料管線

params 與 parameters 的差異

params(Jinja 渲染)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • 值被插入 SQL 字串
  • 易受 SQL injection 攻擊

parameters(資料庫層綁定)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • 值傳給資料庫驅動程式
  • 防注入:驅動程式負責跳脫
使用 Airflow 建置資料管線

用 Astro 從本機到正式環境

$$

  • Astronomer's Astro:受管的 Airflow 平台

$$

  • Astro CLI:一行指令啟動本機 Airflow

$$

  • 本機開發,無縫部署到正式環境

Astro CLI 與平台

使用 Airflow 建置資料管線

一起來練習吧!

使用 Airflow 建置資料管線

Preparing Video For Download...