728x90
반응형
SMALL
https://github.com/hosose/DE/tree/airflow
GitHub - hosose/DE at airflow
Contribute to hosose/DE development by creating an account on GitHub.
github.com
1. 실전 시나리오 개요
스마트 팩토리의 오븐 센서에서 10개의 측정 데이터(온도 및 생성 시간)가 수집되어 파일 형태로 적재된다고 가정합니다.
우리는 이 데이터를 정제하고, 단위 변환 파생변수(섭씨 ➔ 화씨)를 추가한 뒤 MySQL 데이터 웨어하우스에 자동 적재하는 ETL 파이프라인을 구축합니다.
flowchart LR
create_table[SQLExecuteQueryOperator<br/>1. MySQL 테이블 생성] --> extract[PythonOperator<br/>2. Extract: JSON 생성]
extract --> transform[PythonOperator<br/>3. Transform: Pandas 정제]
transform --> load[PythonOperator<br/>4. Load: MySqlHook 벌크 적재]
2. 파이프라인 전체 구현 코드 (05._mysql_etl.py)
소스 코드 위치: dags/05._mysql_etl.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
from airflow.providers.mysql.hooks.mysql import MySqlHook
from datetime import datetime, timedelta
from contextlib import closing
import logging, pendulum, json, random, os
import pandas as pd
KST = pendulum.timezone("Asia/Seoul")
DATA_PATH = "/opt/airflow/dags/data"
os.makedirs(DATA_PATH, exist_ok=True)
# 1. EXTRACT: 더미 센서 데이터 수집 (JSON 파일 저장 후 파일 경로 XCom 반환)
def _extract(**kwargs):
data = [
{
"sensor_id": f"SENSOR_{i+1}",
"timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
"temperature": round(random.uniform(20.0, 150.0), 2),
"status": "on"
}
for i in range(10)
]
file_path = f"{DATA_PATH}/sensor_data_{kwargs['ds_nodash']}.json"
with open(file_path, 'w') as f:
json.dump(data, f)
return file_path
# 2. TRANSFORM: Pandas 데이터 클리닝 및 파생변수 생성
def _transform(**kwargs):
ti = kwargs["ti"]
json_path = ti.xcom_pull(task_ids="extract")
df = pd.read_json(json_path)
# 이상치 제거: 섭씨 100도 이하 센서 데이터만 필터링
clean_df = df[df['temperature'] <= 100].copy()
# 파생변수 생성: 섭씨(°C) ➔ 화씨(°F) 변환
clean_df["temperature_f"] = (clean_df['temperature'] * 9/5) + 32
csv_path = f"{DATA_PATH}/preprocessing_data_{kwargs['ds_nodash']}.csv"
clean_df.to_csv(csv_path, index=False)
return csv_path
# 3. LOAD: MySqlHook 및 executemany를 활용한 DB 적재
def _load(**kwargs):
ti = kwargs["ti"]
csv_path = ti.xcom_pull(task_ids="transform")
df = pd.read_csv(csv_path)
hooks = MySqlHook(mysql_conn_id="mysql_default")
with closing(hooks.get_conn()) as conn:
with closing(conn.cursor()) as cursor:
sql = """
INSERT INTO sensor_readings
(sensor_id, timestamp, temperature_c, temperature_f)
VALUES (%s, %s, %s, %s)
"""
params = [
(row['sensor_id'], row['timestamp'], row['temperature'], row['temperature_f'])
for _, row in df.iterrows()
]
cursor.executemany(sql, params)
conn.commit()
# DAG 설정
with DAG(
dag_id="05_mysql_etl",
schedule_interval="@daily",
start_date=pendulum.datetime(2026, 6, 29, tz=KST),
catchup=False,
tags=['etl', 'mysql']
) as dag:
# 0. DDL: MySQL 테이블 사전 생성
task_create_table = SQLExecuteQueryOperator(
task_id="create_table",
conn_id="mysql_default",
sql="""
CREATE TABLE IF NOT EXISTS sensor_readings (
id INT AUTO_INCREMENT PRIMARY KEY,
sensor_id VARCHAR(50),
timestamp DATETIME,
temperature_c FLOAT,
temperature_f FLOAT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
"""
)
task_extract = PythonOperator(task_id="extract", python_callable=_extract)
task_transform = PythonOperator(task_id="transform", python_callable=_transform)
task_load = PythonOperator(task_id="load", python_callable=_load)
# 순서 연결
task_create_table >> task_extract >> task_transform >> task_load
3. 핵심 기술 요소 분석
- 데이터 전달 최소화 기법: 대용량 Raw Data 자체를 XCom으로 주고받으면 메타 DB에 과부하가 걸립니다. 대신 파일 경로(
file_path)만 XCom으로 전달하고 파일 기반 통신을 수행합니다. - 안전한 DB 커넥션 관리:
contextlib.closing()을 적용하여 예외 상황 발생 시에도 Connection과 Cursor 자원이 안전하게 반납되도록 작성했습니다. - 벌크 삽입 (
executemany): 반복적인 단건execute()대신executemany()를 사용하여 데이터베이스 I/O 성능을 최적화했습니다.
4. 4편 요약 및 다음 편 예고
단일 DAG 내에서 Extract ➔ Transform ➔ Load 전체 과정을 성공적으로 구현하였습니다.
마지막 5편에서는 단일 DAG의 한계를 극복하는 Multi-DAG 모듈화 연동(TriggerDagRunOperator)과 외부 REST API(FastAPI) 서비스 통합을 다루어 연재를 마무리하겠습니다.
반응형
LIST