← 학습 카테고리

Learn

Airflow

151개 모듈 · 현재 7번째

Airflow 모듈 7/151 airflow-learn-07

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 실행 전체의 소요 시간에 상한을 두어, 무한정 멈춰 있는 실행을 방지하는 안전장치다