Airflow로 데이터 파이프라인 구축하기
Volker Janz
Senior Developer Advocate at Astronomer

from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperatorSQLExecuteQueryOperator( task_id="load_sales", conn_id="duckdb_analytics", sql="INSERT INTO sales SELECT * FROM staging", )
conn_id로 참조합니다
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
""",
)
@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/ 폴더 외부에 보관하세요 💡
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)로 렌더링됩니다params (Jinja 렌더링)
SQLExecuteQueryOperator(
sql="""SELECT * FROM orders
WHERE product = '{{ params.product }}'""",
params={"product": user_input},
)
parameters (DB 수준 바인딩)
SQLExecuteQueryOperator(
sql="SELECT * FROM orders
WHERE product = $product",
parameters={"product": user_input},
)
$$
$$
$$

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