Chạy workload SQL

Xây dựng Data Pipeline với Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Vì sao dùng SQL trong Airflow?

 

  • Workload SQL là use case phổ biến nhất của Airflow
  • Để cơ sở dữ liệu gánh phần nặng
  • Airflow điều phối khi nàoở đâu
  • Cơ sở dữ liệu xử lý như thế nào

Điều phối SQL trong Airflow

Xây dựng Data Pipeline với 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", )

 

  • Làm việc với mọi cơ sở dữ liệu có Airflow provider tương thích
  • PostgreSQL, Snowflake, BigQuery, DuckDB, và hơn thế nữa
Xây dựng Data Pipeline với Airflow

Connections

  • Lưu thông tin xác thực ngoài mã nguồn
  • Mỗi connection có ID, loại, host, port, login
  • Tạo theo nhiều cách: UI, CLI, API hoặc biến môi trường
  • Tham chiếu bằng conn_id trong operator

Giao diện Connections của Airflow

Xây dựng Data Pipeline với Airflow

Xây dựng 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
        """,
    )
Xây dựng Data Pipeline với Airflow

Tệp SQL bên ngoài

@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",
    )
  • Đặt template_searchpath trên @dag
  • Tham chiếu tên tệp trong sql
  • Để script và logic nghiệp vụ ngoài thư mục dags/ 💡

Cấu trúc tệp dự án

Xây dựng Data Pipeline với Airflow

Mẫu Jinja trong tệp 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 }} hiển thị ngày logic (YYYY-MM-DD)
  • Xóa rồi chèn: mẫu idempotent từ Chương 2
  • Chạy lại cùng ngày cho kết quả như nhau
Xây dựng Data Pipeline với Airflow

params và parameters

params (rendering Jinja)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • Giá trị được chèn vào chuỗi SQL
  • Dễ bị SQL injection

parameters (ràng buộc ở cấp DB)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • Giá trị được truyền cho trình điều khiển DB
  • An toàn trước injection: driver xử lý escaping
Xây dựng Data Pipeline với Airflow

Từ local đến production với Astro

$$

  • Astronomer's Astro: nền tảng Airflow được quản lý

$$

  • Astro CLI: chạy Airflow cục bộ với một lệnh

$$

  • Phát triển cục bộ, triển khai production mượt mà

Astro CLI và nền tảng

Xây dựng Data Pipeline với Airflow

Cùng luyện tập nào!

Xây dựng Data Pipeline với Airflow

Preparing Video For Download...