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 프로바이더가 있는 모든 데이터베이스와 함께 사용 가능합니다
  • PostgreSQL, Snowflake, BigQuery, DuckDB 등을 지원합니다
Airflow로 데이터 파이프라인 구축하기

연결(Connections)

  • 자격 증명을 코드 외부에 저장하세요
  • 각 연결에는 ID, 유형, 호스트, 포트, 로그인 정보가 있습니다
  • 다양한 방법으로 생성 가능: UI, CLI, API, 환경 변수
  • 오퍼레이터에서 conn_id로 참조합니다

Airflow 연결 UI

Airflow로 데이터 파이프라인 구축하기

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
        """,
    )
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",
    )
  • @dagtemplate_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 vs parameters

params (Jinja 렌더링)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • 값이 SQL 문자열에 직접 삽입됩니다
  • SQL 인젝션 취약점이 있습니다

parameters (DB 수준 바인딩)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • 값이 데이터베이스 드라이버에 전달됩니다
  • 인젝션 방어: 드라이버가 이스케이프 처리를 담당합니다
Airflow로 데이터 파이프라인 구축하기

Astro로 로컬에서 프로덕션까지

$$

  • Astronomer의 Astro: 관리형 Airflow 플랫폼

$$

  • Astro CLI: 명령어 하나로 로컬 Airflow 실행

$$

  • 로컬에서 개발하고 프로덕션에 원활하게 배포

Astro CLI 및 플랫폼

Airflow로 데이터 파이프라인 구축하기

연습해 봅시다!

Airflow로 데이터 파이프라인 구축하기

Preparing Video For Download...