การรัน SQL workloads

การสร้าง Data Pipeline ด้วย Airflow

Volker Janz

Senior Developer Advocate at Astronomer

ทำไมต้องใช้ SQL ใน Airflow?

 

  • workload SQL คือ use case ที่พบบ่อยที่สุดใน Airflow
  • ให้ ฐานข้อมูล รับภาระงานหนัก
  • Airflow จัดการ เมื่อไหร่ และ ที่ไหน
  • ฐานข้อมูลจัดการ อย่างไร

การ orchestrate SQL ใน Airflow

การสร้าง Data Pipeline ด้วย 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 และอื่น ๆ
การสร้าง Data Pipeline ด้วย Airflow

Connections

  • เก็บ ข้อมูลรับรอง ไว้นอกโค้ด
  • แต่ละรายการมี ID, ประเภท, host, port, login
  • สร้างได้ หลายวิธี: UI, CLI, API หรือ environment variables
  • อ้างอิงด้วย conn_id ใน operator

UI Connections ของ Airflow

การสร้าง Data Pipeline ด้วย 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
        """,
    )
การสร้าง Data Pipeline ด้วย 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",
    )
  • กำหนด template_searchpath ใน @dag
  • ระบุ ชื่อไฟล์ ใน sql
  • เก็บ script และ business logic ไว้นอกโฟลเดอร์ dags/ 💡

โครงสร้างไฟล์ของโปรเจกต์

การสร้าง Data Pipeline ด้วย Airflow

Jinja templates ในไฟล์ 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 }} แทนค่า logical date (YYYY-MM-DD)
  • DELETE แล้ว INSERT คือ รูปแบบ idempotent จาก Chapter 2
  • รันซ้ำวันเดิมได้ ผลลัพธ์เหมือนเดิมเสมอ
การสร้าง Data Pipeline ด้วย Airflow

params vs parameters

params (Jinja rendering)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • ค่าถูก แทรกเข้าใน SQL string โดยตรง
  • เสี่ยงต่อ SQL injection

parameters (DB-level binding)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • ค่าถูกส่งไปยัง database driver
  • ปลอดภัยจาก injection: driver จัดการการ escape เอง
การสร้าง Data Pipeline ด้วย Airflow

จาก local สู่ production ด้วย Astro

$$

  • Astro ของ Astronomer: แพลตฟอร์ม Airflow แบบ managed

$$

  • Astro CLI: รัน Airflow ในเครื่องได้ด้วยคำสั่งเดียว

$$

  • พัฒนาใน local แล้ว deploy ขึ้น production ได้อย่างราบรื่น

Astro CLI และแพลตฟอร์ม

การสร้าง Data Pipeline ด้วย Airflow

มาฝึกกันเถอะ!

การสร้าง Data Pipeline ด้วย Airflow

Preparing Video For Download...