간단한 데이터 파이프라인 만들기
간단한 데이터 파이프라인 만들기 (Building a Simple Data Pipeline)
이 튜토리얼은 외부 소스에서 데이터를 가져와 데이터베이스에 로드하고, 그 과정에서 정리하는 작지만 의미 있는 데이터 파이프라인을 만드는 방법을 보여드려요. SQLExecuteQueryOperator와 Postgres provider를 활용하는 실무 패턴을 배웁니다.
출처: 문서
본문
우리 시리즈의 세 번째 튜토리얼에 오신 것을 환영해요! 이 시점에서 여러분은 이미 첫 DAG을 작성하고 몇 가지 기본 operators를 사용해 봤을 거예요. 이제 외부 소스에서 데이터를 가져와 데이터베이스에 로드하고, 그 과정에서 정리하는 작지만 의미 있는 데이터 파이프라인을 만들 차례예요.
이 튜토리얼은 Airflow에서 SQL을 실행하는 유연하고 현대적인 방법인 SQLExecuteQueryOperator를 소개해요. 이를 사용해 Airflow UI에서 구성할 로컬 Postgres 데이터베이스와 상호작용할 거예요.
이 튜토리얼을 마치면 다음과 같은 동작하는 파이프라인을 갖게 돼요:
- CSV 파일 다운로드
- 스테이징 테이블로 데이터 로드
- 데이터 정리 후 대상 테이블로 upsert
그 과정에서 Airflow의 UI, connection 시스템, SQL 실행, DAG 작성 패턴에 대한 실무 경험을 쌓게 될 거예요.
진행하면서 더 깊이 알고 싶다면 아래 두 참조가 유용해요:
시작해 볼까요!
초기 설정 (Initial setup)
Caution
이 튜토리얼을 실행하려면 Docker가 설치되어 있어야 해요. Docker Compose를 사용해 Airflow를 로컬에서 실행할 거예요. 설정에 도움이 필요하다면 Docker Compose 퀵스타트 가이드를 확인하세요.
파이프라인을 실행하려면 동작하는 Airflow 환경이 필요해요. Docker Compose는 시스템 전체 설치 없이 쉽고 안전하게 만들어줘요. 터미널을 열고 다음을 실행하기만 하면 돼요:
# Download the docker-compose.yaml file
curl -LfO 'https://airflow.apache.org/docs/apache-airflow/stable/docker-compose.yaml'
# Make expected directories and set an expected environment variable
mkdir -p ./dags ./logs ./plugins
echo -e "AIRFLOW_UID=$(id -u)" > .env
# Initialize the database
docker compose up airflow-init
# Start up all services
docker compose up
Airflow가 실행되면 UI를 http://localhost:8080에서 방문해요.
다음으로 로그인해요:
- Username:
airflow - Password:
airflow
Airflow 대시보드에 도착하면 DAG을 트리거하고, 로그를 탐색하고, 환경을 관리할 수 있어요.
Postgres Connection 만들기
파이프라인이 Postgres에 쓸 수 있으려면, Airflow에 어떻게 연결해야 하는지 알려줘야 해요. UI에서 Admin > Connections 페이지를 열고 + 버튼을 클릭해 새 connection을 추가해요.
다음 세부 정보를 채워요:
- Connection ID:
tutorial_pg_conn - Connection Type:
postgres - Host:
postgres - Database:
airflow(이것은 우리 컨테이너의 기본 데이터베이스) - Login:
airflow - Password:
airflow - Port:
5432

connection을 저장해요. 이렇게 하면 Airflow에 Docker 환경에서 실행 중인 Postgres 데이터베이스에 도달하는 방법을 알려줘요.
다음으로, 이 connection을 사용하는 파이프라인을 만들기 시작할게요.
스테이징 및 최종 데이터용 테이블 만들기
테이블 생성부터 시작해요. 두 개의 테이블을 만들 거예요:
employees_temp: 원본 데이터용 스테이징 테이블employees: 정리되고 중복 제거된 목적지
이 테이블들을 만드는 데 필요한 SQL 문을 실행하기 위해 SQLExecuteQueryOperator를 사용할 거예요.
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
);""",
)
선택적으로 dags/ 폴더 안의 .sql 파일에 이 SQL 문장들을 넣고 sql= 인자에 파일 경로를 전달할 수도 있어요. 이는 DAG 코드를 깔끔하게 유지하는 좋은 방법이에요.
스테이징 테이블로 데이터 로드하기
다음으로 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():
# NOTE: configure this as appropriate for your Airflow environment
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")
postgres_hook.copy_expert(
"COPY employees_temp FROM STDIN WITH CSV HEADER DELIMITER AS ',' QUOTE '\"'",
data_path,
)
이 Task는 Airflow를 네이티브 Python 및 SQL hooks와 결합하는 맛을 보여줘요 — 실제 파이프라인에서 흔한 패턴이에요.
데이터 병합 및 정리
이제 데이터를 중복 제거하고 최종 테이블로 병합해요. SQL INSERT … ON CONFLICT DO UPDATE를 실행하는 Task를 작성할 거예요.
from airflow.sdk import task
from airflow.providers.postgres.hooks.postgres import PostgresHook
@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
DAG 정의하기
이제 모든 Task를 정의했으니, 그것들을 DAG으로 묶을 차례예요.
import datetime
import pendulum
import os
import requests
from airflow.sdk import dag, task
from airflow.providers.postgres.hooks.postgres import PostgresHook
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
@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(
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
);""",
)
@task
def get_data():
# NOTE: configure this as appropriate for your Airflow environment
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")
postgres_hook.copy_expert(
"COPY employees_temp FROM STDIN WITH CSV HEADER DELIMITER AS ',' QUOTE '\"'",
data_path,
)
@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
[create_employees_table, create_employees_temp_table] >> get_data() >> merge_data()
dag = ProcessEmployees()
이 DAG을 dags/process_employees.py로 저장해요. 잠시 후 UI에 나타날 거예요.
DAG 트리거하고 살펴보기
Airflow UI를 열고 목록에서 process_employees DAG을 찾아요. 슬라이더로 "on"으로 토글한 뒤, 실행 버튼으로 run을 트리거해요.
Grid 뷰에서 각 Task가 실행되는 것을 보고, 각 단계의 로그를 탐색할 수 있어요.



성공하면, 외부 세계의 데이터를 통합해 Postgres에 로드하고 깨끗하게 유지하는 완전히 동작하는 파이프라인을 갖게 돼요.
다음은 무엇일까요? (What's Next?)
수고했어요! 이제 Airflow의 핵심 패턴과 도구로 실제 파이프라인을 만들었어요. 다음 단계를 위한 몇 가지 아이디어:
- MySQL이나 SQLite처럼 다른 SQL provider로 바꿔 보세요.
- DAG을 TaskGroups로 나누거나 더 유용한 패턴으로 리팩터링해 보세요.
- 데이터 처리 시 알림 단계를 추가하거나 알림을 보내 보세요.
참고 (See also)
- Airflow 문서에서 더 많은 how-to 가이드 둘러보기
- SQL provider 참조 탐색
- 자신만의 커스텀 operator 작성 배우기