SQL 기반 데이터 파이프라인 구축하기
Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation docs/Airflow/docs/tutorial/pipeline.rst
이 모듈을 다 읽으면
- Airflow의 Connection과 SQLExecuteQueryOperator를 이용해 외부 DB에 SQL을 실행하는 방법을 설명할 수 있다
- 스테이징 테이블에 적재한 뒤 정제·업서트하는 2단계 패턴이 왜 필요한지 설명할 수 있다
- Hook을 이용한 커스텀 TaskFlow 태스크와 오퍼레이터를 같은 Dag 의존성 체인 안에서 함께 쓰는 방법을 설명할 수 있다
SQLExecuteQueryOperator와 PostgresHook을 이용해 CSV 파일을 다운로드하고, 스테이징 테이블에 적재한 뒤, 중복을 제거하며 최종 테이블에 업서트하는 작은 데이터 파이프라인을 구축한다. Admin UI에서 Connection을 등록하는 절차, 오퍼레이터와 TaskFlow 태스크를 한 Dag 안에서 함께 쓰는 패턴을 다룬다.
Admin > Connections로 외부 DB 연결 등록하기
파이프라인이 Postgres에 쓰기 작업을 하려면, 먼저 Airflow에게 어떻게 그 DB에 연결할지 알려줘야 한다. UI의 Admin → Connections 페이지에서 + 버튼을 눌러 새 커넥션을 추가한다. Connection ID(예: ``tutorial_pg_conn``), Connection Type(``postgres``), Host, Database, Login, Password, Port 같은 세부 정보를 입력하고 저장하면, 이 커넥션 정보가 Docker 환경에서 실행 중인 Postgres에 도달하는 방법을 Airflow에게 알려준다. 이렇게 커넥션을 UI(또는 환경변수/시크릿 백엔드)로 등록해두면, Dag 코드에는 자격증명이 하드코딩되지 않고 ``conn_id``만 참조하면 된다는 것이 핵심이다.
핵심 포인트
- Admin → Connections에서 Connection ID/Type/Host/Database/Login/Password/Port를 등록해 Airflow가 외부 DB에 연결하는 방법을 알게 한다
- Dag 코드는 자격증명을 직접 갖지 않고 conn_id만 참조하므로, 코드와 자격증명 관리가 분리된다
SQLExecuteQueryOperator로 테이블 준비하기
``SQLExecuteQueryOperator``는 Airflow에서 SQL을 실행하는 유연하고 현대적인 방법이다. 이 튜토리얼에서는 원본 데이터를 담을 스테이징 테이블 ``employees_temp``와, 정제되고 중복 제거된 최종 목적지인 ``employees`` 테이블 두 개를 만든다.
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
create_employees_table = SQLExecuteQueryOperator(
task_id="create_employees_table",
conn_id="tutorial_pg_conn",
sql="""
CREATE TABLE IF NOT EXISTS employees (
"Serial Number" NUMERIC PRIMARY KEY,
"Company Name" TEXT,
"Employee Markme" TEXT,
"Description" TEXT,
"Leave" INTEGER
);""",
)
create_employees_temp_table = SQLExecuteQueryOperator(
task_id="create_employees_temp_table",
conn_id="tutorial_pg_conn",
sql="""
DROP TABLE IF EXISTS employees_temp;
CREATE TABLE employees_temp (
"Serial Number" NUMERIC PRIMARY KEY,
"Company Name" TEXT,
"Employee Markme" TEXT,
"Description" TEXT,
"Leave" INTEGER
);""",
)
스테이징 테이블(``employees_temp``)은 매 실행마다 ``DROP TABLE IF EXISTS`` 후 재생성되어 항상 이번 적재분만 담고, 최종 테이블(``employees``)은 ``CREATE TABLE IF NOT EXISTS``로 한 번만 만들어져 누적된 정제 데이터를 보관한다는 점이 두 SQL의 차이다. 이 SQL 문들을 ``.sql`` 파일로 ``dags/`` 폴더 안에 두고 ``sql=`` 인자에 파일 경로를 넘겨 Dag 코드를 더 깔끔하게 유지할 수도 있다.
핵심 포인트
- SQLExecuteQueryOperator는 conn_id로 등록된 커넥션을 참조해 SQL을 실행하는 범용 오퍼레이터다
- 스테이징 테이블은 매번 DROP 후 재생성(이번 적재분만 보관), 최종 테이블은 CREATE IF NOT EXISTS(누적 보관)로 역할이 다르다
PostgresHook과 @task로 데이터 적재하기
다음으로 CSV 파일을 다운로드해 로컬에 저장하고, ``PostgresHook``을 이용해 ``employees_temp``에 적재하는 태스크를 작성한다.
import os
import requests
from airflow.sdk import task
from airflow.providers.postgres.hooks.postgres import PostgresHook
@task
def get_data():
data_path = "/opt/airflow/dags/files/employees.csv"
os.makedirs(os.path.dirname(data_path), exist_ok=True)
url = "https://raw.githubusercontent.com/apache/airflow/main/airflow-core/docs/tutorial/pipeline_example.csv"
response = requests.request("GET", url)
with open(data_path, "w") as file:
file.write(response.text)
postgres_hook = PostgresHook(postgres_conn_id="tutorial_pg_conn")
conn = postgres_hook.get_conn()
cur = conn.cursor()
with open(data_path, "r") as file:
cur.copy_expert(
"COPY employees_temp FROM STDIN WITH CSV HEADER DELIMITER AS ',' QUOTE '\"'",
file,
)
conn.commit()
이 태스크는 requests로 파일을 내려받고, PostgresHook.get_conn()으로 얻은 커넥션의 커서에서 ``COPY ... FROM STDIN``으로 CSV를 대량 적재(bulk load)한다. Airflow를 순수 Python 코드 및 SQL 훅과 결합해 쓰는 이 패턴은 실무 파이프라인에서 매우 흔하게 등장한다.
핵심 포인트
- PostgresHook은 conn_id로 등록된 커넥션을 얻어와 psycopg2 커서 등 저수준 DB 접근을 제공하는 훅이다
- cursor.copy_expert("COPY ... FROM STDIN ...")는 CSV 파일을 행 단위 INSERT보다 훨씬 빠르게 대량 적재하는 Postgres의 벌크 로드 방식이다
업서트 패턴으로 정제하고 병합하기
이제 데이터를 중복 제거하고 최종 테이블에 병합한다. ``INSERT ... ON CONFLICT DO UPDATE`` 형태의 SQL을 실행하는 태스크를 작성한다.
@task
def merge_data():
query = """
INSERT INTO employees
SELECT *
FROM (
SELECT DISTINCT *
FROM employees_temp
) t
ON CONFLICT ("Serial Number") DO UPDATE
SET
"Employee Markme" = excluded."Employee Markme",
"Description" = excluded."Description",
"Leave" = excluded."Leave";
"""
try:
postgres_hook = PostgresHook(postgres_conn_id="tutorial_pg_conn")
conn = postgres_hook.get_conn()
cur = conn.cursor()
cur.execute(query)
conn.commit()
return 0
except Exception as e:
return 1
``SELECT DISTINCT``로 스테이징 테이블 안의 중복을 먼저 제거한 다음, 기본 키(``"Serial Number"``)가 이미 존재하면 지정한 컬럼들만 최신 값으로 갱신(UPDATE)하고 없으면 새로 삽입(INSERT)하는 것이 이 업서트 패턴의 핵심이다. 이 태스크는 성공하면 0, 예외가 발생하면 1을 반환해 실행 결과를 후속 로직이나 모니터링에서 참조할 수 있게 한다.
핵심 포인트
- INSERT ... ON CONFLICT (pk) DO UPDATE는 기본 키 충돌 시 갱신, 없으면 삽입하는 업서트를 한 문장으로 표현한다
- SELECT DISTINCT로 스테이징 데이터를 먼저 중복 제거한 뒤 최종 테이블에 병합하는 것이 이 정제 단계의 핵심이다
오퍼레이터와 TaskFlow 태스크를 함께 연결하기
마지막으로 이 모든 태스크를 하나의 Dag으로 묶는다.
@dag(
dag_id="process_employees",
schedule="0 0 * * *",
start_date=pendulum.datetime(2021, 1, 1, tz="UTC"),
catchup=False,
dagrun_timeout=datetime.timedelta(minutes=60),
)
def ProcessEmployees():
create_employees_table = SQLExecuteQueryOperator(...)
create_employees_temp_table = SQLExecuteQueryOperator(...)
@task
def get_data():
...
@task
def merge_data():
...
[create_employees_table, create_employees_temp_table] >> get_data() >> merge_data()
dag = ProcessEmployees()
의존성 체인 ``[create_employees_table, create_employees_temp_table] >> get_data() >> merge_data()``는 오퍼레이터(SQLExecuteQueryOperator)와 TaskFlow 태스크(@task)가 같은 Dag 안에서 자연스럽게 함께 연결될 수 있음을 보여준다. 두 테이블 생성 태스크는 서로 독립적이므로 병렬로 실행되고, 둘 다 끝난 뒤에야 데이터 적재(``get_data``)와 병합(``merge_data``)이 순서대로 이어진다. ``dagrun_timeout=datetime.timedelta(minutes=60)``은 이 Dag 실행 전체가 60분을 넘기면 강제로 실패 처리되도록 상한을 두는 안전장치다.
핵심 포인트
- 리스트를 이용한 `[taskA, taskB] >> taskC`는 taskA/B를 병렬로 먼저 실행한 뒤 그 둘이 모두 끝나야 taskC를 실행하는 팬인(fan-in) 구조를 만든다
- SQLExecuteQueryOperator(오퍼레이터)와 @task 함수(TaskFlow)는 >> 로 자연스럽게 하나의 의존성 체인에 섞일 수 있다
- dagrun_timeout은 Dag 실행 전체의 소요 시간에 상한을 두어, 무한정 멈춰 있는 실행을 방지하는 안전장치다