Hands-on: Mini Data Pipeline
4ํธ ยท ์จ๋ณด๋ฉ ์ฌ๋ผ์ด๋์ ๊ฐ๋ ์ง๋๋ฅผ ์ง์ ๋๋ ค๋ณด๊ธฐ
- Haram Lee
- 2026-07-05
- work / Daangn / data-team-notes
๐ ๊ฐ๋ ยท์ฝ๋ ์ ๋ฆฌ ์๋ฃ(์ด์ ๋ฐํ). ใ4. ์คํ ๊ฒฐ๊ณผใ์ ์คํฌ๋ฆฐ์ท/๋ก๊ทธ๋ ๋ก์ปฌ์์ DAG๋ฅผ ์ง์ ๋๋ฆฐ ๋ค ์ด์ด์ ์ฑ์ด๋ค.
์์ ์ธ ํธ(DT Platform ํธ, MongoDB CDC ํธ, DBT + Airflow ํธ)์ด ๋น๊ทผ ๋ฐ์ดํฐ ๊ฐ์นํํ์ ๊ธฐ์ ๋ธ๋ก๊ทธ๋ฅผ ์ฝ๊ณ ์ง๋๋ฅผ ๊ทธ๋ฆฐ ๋ ธํธ์๋ค๋ฉด, ์ด๋ฒ ํธ์ ๊ทธ ์ง๋๋ฅผ ์ง์ ์๊ฒ ํ ๋ฒ ๋๋ ค๋ณด๋ ์ค์ต์ด๋ค.
์จ๋ณด๋ฉ์์ ๋ฐ์ ์ฌ๋ผ์ด๋ ํ ์ฅ์ด ์ถ๋ฐ์ ์ด์๋ค. ๋ฐ์ดํฐ ์ธํ๋ผ ์ ์ฒด๋ฅผ ํ ํ๋ฉด์ ๊ทธ๋ฆฐ ์ง๋์๊ณ , ๊ตฌ์์ “100 of 300” ์ด๋ผ๊ณ ์ ํ ์์๋ค. ์ฒ์์ ๋ง๋งํ๋๋ฐ, ๊ฒฐ๊ตญ ํ๊ณ ์ถ์ ๋ง์ ์ด๊ฑฐ์๋ค.
“์ด๊ฑธ ์ ๋ถ ๋ง์คํฐํด๋ผ"๊ฐ ์๋๋ผ, ์ ์ฒด ํ(300) ์ค ํต์ฌ ํ๋ฆ(100)๋งํผ์ ์ง์ ๋ณด๊ณ ๊ฐ์ ์ก์๋ผ.
๊ทธ๋์ ํด๋ผ์ฐ๋ ๊ณ์ ์์ด๋ ๊ตฌ์กฐ๋ง ํ์ฌ์์ผ๋ก ๋๊ฐ์ด ๊ฐ์ ธ๊ฐ๋ ๋ฏธ๋ ํ์ดํ๋ผ์ธ์ ๋ก์ปฌ์ ์ธ์๋ดค๋ค.
1. ์ฌ๋ผ์ด๋๊ฐ ๋งํ๋ ๊ฒ โ ๋ฐ์ดํฐ ์ธํ๋ผ ์ ์ฒด ์ง๋
์ฌ๋ผ์ด๋๋ ํฌ๊ฒ Pipelines(๋ฐ์ดํฐ๊ฐ ํ๋ฅด๋ ๊ธธ) ๊ณผ Infrastructure(๊ทธ ๊ธธ์ด ์ค์ ๋ก ๋์๊ฐ๋ ๋ฐ๋ฅ) ์ผ๋ก ๋๋๊ณ , ๊ทธ ์์ ๋ค์ฏ ์ถ์ด ์์๋ค.
Pipelines โโฌโ Orchestration : Airflow
โโ Batch : Spark, RDBMS, NoSQL
โโ Streaming : Flink, Kafka, Pub/Sub, Dataflow
Infra โโโโโโฌโ Kubernetes, Karpenter
โโ DPaaS : DT Platform + @๊ฐ ์ถ์ ์์ ๋ ธํธ๋ค๊ณผ ์ด๋ ๊ฒ ์ฐ๊ฒฐ๋๋ค.
| ์ฌ๋ผ์ด๋ ์ถ | ํต์ฌ ๋๊ตฌ | ํ ์ค ์ ์ | ์ด์ด์ง๋ ๋ ธํธ |
|---|---|---|---|
| Orchestration | Airflow | ์ธ์ ยท์ด๋ค ์์๋ก ์์ ์ ์คํํ ์ง ๊ด๋ฆฌ | DBT + Airflow ํธ |
| Batch | Spark / RDBMS / NoSQL | ์์ธ ๋ฐ์ดํฐ๋ฅผ ์ฃผ๊ธฐ์ ์ผ๋ก ํ ๋ฒ์ ์ฒ๋ฆฌ | DT Platform ํธ |
| Streaming | Flink / Kafka / Pub/Sub / Dataflow | ๋ฐ์ดํฐ๊ฐ ์๊ธฐ๋ ์ฆ์ ํ๋ ค๋ณด๋ด๋ฉฐ ์ฒ๋ฆฌ | MongoDB CDC ํธ |
| Infra | Kubernetes / Karpenter | ์ ์์ ๋ค์ด ์ค์ ๋ก ์คํ๋๋ ํด๋ฌ์คํฐ | โ |
| DPaaS | DT Platform | ์ ์ ๋ถ๋ฅผ ๊ตฌ์ฑ์์ด ์ฝ๊ฒ ์ฐ๊ฒ ๋ง๋ ๋ด๋ถ ํ๋ซํผ | DT Platform ํธ |
ํ ์ค๋ก ๋ค์ ๋ฌถ์ผ๋ฉด:
Airflow = ์์
์งํ์
Spark = ํฐ ๋ฐ์ดํฐ ์ฒ๋ฆฌ ๋
ธ๋์
Kafka = ์ค์๊ฐ ์ด๋ฒคํธ ์ฐ์ฒด๊ตญ
Flink = ์ค์๊ฐ ์ฒ๋ฆฌ ๊ณต์ฅ
Kubernetes = ์ด๋ค์ด ์ฌ๋ ์์
์ฅ
Karpenter = ์์
์ฅ์ด ๋ถ์กฑํ๋ฉด ๊ฑด๋ฌผ์ ๋ ๋น๋ ค์ค๋ ์ 2. ๊ธฐ์ ์คํ ์ฌ์
์ฌ๋ผ์ด๋์ ์ ํ ๋๊ตฌ๋ค์ ๊ณต์ ๋ฌธ์ ์ ์๋ฅผ ๊ทผ๊ฑฐ๋ก, ํต์ฌ ํค์๋ ์ค์ฌ์ผ๋ก ์ ๋ฆฌํ๋ค. (์ธ์ฉ ๋งํฌ๋ ๋งจ ์๋ ์ฐธ๊ณ ์ ๋ชจ์)
Airflow โ Orchestration
๊ณต์ ์ ์: “์ํฌํ๋ก๋ฅผ ๊ฐ๋ฐยท์ค์ผ์คยท๋ชจ๋ํฐ๋งํ๊ธฐ ์ํ ์คํ์์ค ํ๋ซํผ.” ์ํฌํ๋ก๋ฅผ ํ์ด์ฌ ์ฝ๋๋ก ํํํ๋ค(“workflows as code”). โ Apache Airflow docs
- DAG (Directed Acyclic Graph): ์์ ๋ค์ ๋ฐฉํฅ์ฑ ๋น์ํ ๊ทธ๋ํ. ๋ฌด์์ยท์ด๋ค ์์๋ก ์คํํ ์ง ์ ์. ๊ณง ํ์ดํ๋ผ์ธ ๊ทธ ์์ฒด.
- Task: DAG๋ฅผ ์ด๋ฃจ๋ ์ต์ ์คํ ๋จ์. Task ๊ฐ ํ์ดํ๊ฐ ์์กด์ฑ.
- Operator: Task๋ฅผ ๋ง๋๋ ํ
ํ๋ฆฟ (
PythonOperator,BashOperator,@task๋ฑ). - Scheduler: DAG๋ฅผ ํ์ฑํ๊ณ ๋๊ฐ ๋๋ฉด Task๋ฅผ ์คํ ํ์ ๋ฃ๋ ์ฌ์ฅ. / Executor: Task๋ฅผ ์ค์ ๋ก ์คํํ๋ ๋ฐฉ์(LocalยทCeleryยทKubernetes Executor).
- XCom (cross-communication): Task ์ฌ์ด์ ์์ ๊ฐ์ ์ฃผ๊ณ ๋ฐ๋ ํต๋ก. TaskFlow์์ ํจ์ returnโ๋ค์ ํจ์ ์ธ์๊ฐ ๊ณง XCom.
- Sensor: ์ด๋ค ์กฐ๊ฑด(ํ์ผ ๋์ฐฉ, ์ DAG ์๋ฃ ๋ฑ)์ด ๋ ๋๊น์ง ๊ธฐ๋ค๋ฆฌ๋ ํน์ Task. โ DBT + ํธ์์
ExternalTaskSensor๋ก ์์ฒ ์ ์ฌ ์๋ฃ๋ฅผ ๊ธฐ๋ค๋ฆฐ ์ฌ๋ก. - TaskFlow API:
@dag/@task๋ฐ์ฝ๋ ์ดํฐ๋ก ํ์ด์ฌ ํจ์๋ฅผ ๊ทธ๋๋ก Task/DAG๋ก (3์ ์ค์ต์ด ์ด ๋ฐฉ์).
Helm โ ์ฟ ๋ฒ๋คํฐ์ค ํจํค์ง ๋งค๋์
๊ณต์ ์ ์: “์ฟ ๋ฒ๋คํฐ์ค๋ฅผ ์ํ ํจํค์ง ๋งค๋์ (the package manager for Kubernetes).” โ Helm docs
- Chart: ์ฟ ๋ฒ๋คํฐ์ค ๋ฆฌ์์ค ๋ฌถ์์ ๊ธฐ์ ํ ํจํค์ง. (apt์ .deb, brew์ formula ๊ฐ์ ๊ฒ)
- values.yaml: Chart์ ์ฃผ์ ํ๋ ์ค์ ๊ฐ ๋ชจ์. Airflow๋ผ๋ฉด executorยทimageยทworker ์ ๋ฑ.
- Template:
values๋ฅผ ๋ฐ์ ์ต์ข K8s manifest(YAML)๋ฅผ ๋ ๋๋งํ๋ ํ (templates/*.yaml). - Release: ํด๋ฌ์คํฐ์ ์ค์น๋ Chart์ ์ธ์คํด์ค(= ์ค์ ๋ฐฐํฌ๋ณธ). / Repository: Chart๋ฅผ ๋ฐฐํฌํ๋ ์ ์ฅ์.
- ๊ฐ๊ฐ:
values.yaml์์ โtemplates/์ ์ฃผ์ โ ์ต์ข manifest ์์ฑ โ ํด๋ฌ์คํฐ ์ ์ฉ. Airflow๋ ๊ณต์ Helm Chart ์ ๊ณต. (3์ ์์helm template๋ก ์ค์ต)
Terraform โ Infrastructure as Code
๊ณต์ ์ ์: “์ธํ๋ผ๋ฅผ ์ฌ๋์ด ์ฝ์ ์ ์๋ ์ค์ ํ์ผ๋ก ์ ์ํ๋ IaC ๋๊ตฌ.” โ HashiCorp docs
- HCL: Terraform์ ์ ์ธํ ์ค์ ์ธ์ด(“๋ฌด์์ ์ํ๋์ง"๋ฅผ ์ ์).
- Provider: AWSยทGCPยทKubernetes ๋ฑ ์ธ๋ถ API์ ์ํธ์์ฉํ๋ ํ๋ฌ๊ทธ์ธ.
- Resource: ๊ด๋ฆฌ ๋์ ์ธํ๋ผ ๊ฐ์ฒด(S3 bucket, BigQuery dataset, IAM role, EKS cluster ๋ฑ).
- State: ์ค์ ์ธํ๋ผ โ ์ฝ๋ ์ค์ ์ ๋งคํ์ ์ ์ฅ โ ๋ณ๊ฒฝ(diff) ํ๋จ ๊ทผ๊ฑฐ.
- plan / apply: ๋ฐ๋ ๋ด์ฉ ๋ฏธ๋ฆฌ๋ณด๊ธฐ(plan) โ ์ค์ ๋ฐ์(apply).
Batch โ Spark / RDBMS / NoSQL
- Spark: ๋๋ ๋ฐ์ดํฐ๋ฅผ ์ฌ๋ฌ ๋ ธ๋์ ๋ถ์ฐ ์ฒ๋ฆฌํ๋ ์์ง. ๋น๊ทผ์์ DB โ BigQuery ๋๋ ์ ์ฌ์ ์คํ ์์ง์ด๋ฉฐ EMR on EKS(์ฟ ๋ฒ๋คํฐ์ค ์)๋ก ๋๋ค. โ DT Platform ํธ
- RDBMS / NoSQL: ์๋น์ค ํ ์ด์ DB = ํ์ดํ๋ผ์ธ์ Source. MySQLยทPostgreSQL(๊ด๊ณํ), MongoDBยทRedisยทDynamoDB(๋น๊ด๊ณํ).
Streaming โ Kafka / Pub/Sub / Flink / Dataflow
- Kafka / Pub/Sub (๋ฉ์์ง ๋ธ๋ก์ปค): ์ด๋ฒคํธ๋ฅผ ๋ด์ ์๋น์์๊ฒ ์ ๋ฌํ๋ ์คํธ๋ฆผ ์ ์ฅ์/ํ. Kafka๋ ์คํ์์ค, Pub/Sub์ GCP ๊ด๋ฆฌํ.
- Flink / Dataflow (์ฒ๋ฆฌ ์์ง): ๋ธ๋ก์ปค์์ ์ฝ์ด ์ค์๊ฐ์ผ๋ก ๊ฒ์ฆยท์ค๋ณต์ ๊ฑฐยท๋ณํ. ๋น๊ทผ์ ํ๋ ๋ก๊ทธ๋ฅผ Pub/SubโDataflow๋ก, MongoDB ๋ณ๊ฒฝ๋ถ์ Flink CDC๋ก ์ฒ๋ฆฌ. โ MongoDB CDC ํธ
Storage & Warehouse โ S3 / BigQuery
- S3 (Object Storage): ํ์ผ์ ๊ฐ์ฒด(object) ๋ก ์ ์ฅ. ํด๋์ฒ๋ผ ๋ณด์ด๋ ๊ฑด ์ฌ์ค prefix(ํค ์ ๋์ด)์ผ ๋ฟ ์ง์ง ๋๋ ํฐ๋ฆฌ๊ฐ ์๋๋ค(AWS docs).
dt=2026-07-05/๊ฐ์ prefix๊ฐ ํํฐ์ ์ญํ . - BigQuery (Data Warehouse): Cloud Storage/๋ก์ปฌ ํ์ผ์์ load job์ผ๋ก CSVยทJSONยทParquet ์ ์ฌ(GCP docs). ๋ ์ง ๋ฑ์ผ๋ก ํํฐ์ ํ ์ด๋ธ ๊ตฌ์ฑ. ๋น๊ทผ ๋ฐ์ดํฐ ํ๋ฆ์ ์ข ์ฐฉ์ง.
Infra โ Kubernetes / Karpenter
- Kubernetes: ์ปจํ ์ด๋ํ๋ ์ํฌ๋ก๋๋ฅผ ๋ฐฐํฌยทํ์ฅยท๊ด๋ฆฌํ๋ ์คํ์์ค ์์คํ . Pod(๋ฐฐํฌ ์ต์ ๋จ์) โ Node(Pod๊ฐ ๋จ๋ ์๋ฒ) โ Deployment(์ํ๋ ์ํ ์ ์ง).
- Karpenter: ์ฟ ๋ฒ๋คํฐ์ค์ฉ just-in-time ๋ ธ๋ ํ๋ก๋น์ ๋. ์ค์ผ์ค ์ ๋ Pod๋ฅผ ๊ฐ์ง โ ํ์ํ ๋งํผ ๋ ธ๋๋ฅผ ์ฆ์ ๋์ฐ๊ณ , ๋น๋ฉด ์ ๋ฆฌ. Job์ด ๋ชฐ๋ฆฌ๋ฉด Podโ โ ์์ ๋ถ์กฑ โ ๋ ธ๋โ โ ๋๋๋ฉดโ.
DPaaS โ DT Platform + @
- ๊ตฌ์ฑ์์ด UI์์ Sourceยท์ ์ก ํ ์ด๋ธยทDestinationยท์ค์ผ์ค๋ง ๊ณ ๋ฅด๋ฉด ์ค์ ์ด Airflow์ ๋ฐ์๋๊ณ Spark Job์ด ๋๋ no-code ๋ฐ์ดํฐ ์ ์ก ํ๋ซํผ. “์ง์ ๋ค ์ฎ๊ฒจ์ค๋ค"๊ฐ ์๋๋ผ ๋ฐ๋ณต ์์ ์ ํ๋ซํผํํ๊ณ ๊ฐ๋๋ ์ผยทํ์งยท๊ด์ธก์ฑ์ ๋ง๋ ๋ค๋ ๋ฐฉํฅ. โ DT Platform ํธ
3. Hands-on โ ๋ก์ปฌ์ ๋ฏธ๋ ํ์ดํ๋ผ์ธ ์ธ์ฐ๊ธฐ
ํต์ฌ ์์ด๋์ด๋ ์ง์ง ํด๋ผ์ฐ๋๋ ์ ์ฐ๋ ๊ตฌ์กฐ๋ ํ์ฌ์ ๋๊ฐ์ด ๊ฐ์ ธ๊ฐ๋ ๊ฒ์ด๋ค. ์ด๋ ๊ฒ ์นํํ๋ค.
๋ก์ปฌ ํด๋ (ํ์ผ ์ฐ๊ธฐ) โ S3 (raw / curated)
SQLite ํ
์ด๋ธ โ BigQuery (fact_event)
Airflow DAG โ ์ค์ ์ค์ผ์คํธ๋ ์ด์
(๊ทธ๋๋ก)
ํ์ด์ฌ ๋ฆฌ์คํธ๋ก ๋ง๋ ์ด๋ฒคํธ โ Event Bus์์ ํ๋ฌ์จ ๋ก๊ทธํ์ดํ๋ผ์ธ ๋จ๊ณ๋ ์ค์ ๋ฐฐ์น ํ์ดํ๋ผ์ธ๊ณผ ๊ฐ์ ๊ณจ๊ฒฉ์ด๋ค.
extract_to_s3_like โ validate โ transform โ load_to_bq_like โ data_quality_check3.1 ํ๊ฒฝ: Docker Compose๋ก Airflow ๋์ฐ๊ธฐ
Airflow ๊ณต์ ๋ฌธ์์ Docker Compose quick start๋ฅผ ์ฌ์ฉํ๋ค. (ํ๋ก๋์ ์ฉ์ด ์๋๋ผ ๋ก์ปฌ ํ์ต์ฉ)
mkdir -p ~/airflow-study && cd ~/airflow-study
curl -LfO 'https://airflow.apache.org/docs/apache-airflow/stable/docker-compose.yaml'
mkdir -p ./dags ./logs ./plugins ./config ./data
echo "AIRFLOW_UID=$(id -u)" > .env
docker compose up airflow-init # ๋ฉํ๋ฐ์ดํฐ DB ์ด๊ธฐํ (์ต์ด 1ํ)
docker compose up -d # ์ ์ฒด ์คํ ๊ธฐ๋ (webserver/scheduler/worker...)๊ธฐ๋ณธ docker-compose.yaml์ dags/ logs/ config/ plugins/๋ง ํธ์คํธ์ ์ฐ๊ฒฐ(mount)ํ๋ค. ์ฐ๋ฆฌ DAG๋ /opt/airflow/data์ ๊ฒฐ๊ณผ๋ฌผ์ ์ฐ๋ฏ๋ก, ๊ทธ ๊ฒฐ๊ณผ๋ฅผ ๋งฅ์์ ๋์ผ๋ก ๋ณด๋ ค๋ฉด volume์ ํ ์ค ์ถ๊ฐํ๋ค.
# x-airflow-common ์ volumes: ์๋์ ์ถ๊ฐ
- ${AIRFLOW_PROJ_DIR:-.}/data:/opt/airflow/data๊ธฐ๋ ๋ค http://localhost:8080 ์ ์ (๊ธฐ๋ณธ ๊ณ์ airflow / airflow).
3.2 DAG ์ฝ๋
dags/mini_data_pipeline.py. TaskFlow API(@dag, @task)๋ก ์์ฑํ๋ค. ํจ์์ ๋ฐํ๊ฐ์ด ๋ค์ task์ ์ธ์๋ก ๋์ด๊ฐ๋ ๊ฒ์ด ๊ณง ์์กด์ฑ์ด ๋๋ค.
import csv
import json
import sqlite3
from pathlib import Path
import pendulum
# Airflow 3.x / 2.x ๋ ๋ค ๋์
try:
from airflow.sdk import dag, task, get_current_context
except ImportError:
from airflow.decorators import dag, task
from airflow.operators.python import get_current_context
BASE_DIR = Path("/opt/airflow/data")
@dag(
dag_id="mini_data_pipeline",
schedule=None,
start_date=pendulum.datetime(2026, 7, 1, tz="Asia/Seoul"),
catchup=False,
tags=["study", "s3", "bq", "dt-platform"],
)
def mini_data_pipeline():
@task
def extract_to_s3_like() -> str:
"""S3์ raw event๋ฅผ ์๋ ์ํฉ์ ๋ก์ปฌ ํด๋๋ก ํ๋ด. (์ค๋ฌด: s3://bucket/.../dt=.../)"""
ds = get_current_context()["ds"]
raw_dir = BASE_DIR / "s3_like" / "raw" / "event_bus" / f"dt={ds}"
raw_dir.mkdir(parents=True, exist_ok=True)
raw_path = raw_dir / "events.jsonl"
events = [
{"event_id": "e1", "user_id": 101, "event_name": "view_item", "amount": 0},
{"event_id": "e2", "user_id": 102, "event_name": "purchase", "amount": 13000},
{"event_id": "e3", "user_id": 101, "event_name": "purchase", "amount": 7000},
]
with raw_path.open("w") as f:
for event in events:
f.write(json.dumps(event) + "\n")
return str(raw_path)
@task
def validate(raw_path: str) -> str:
"""๋ฐ์ดํฐ ํ์ง ์ฒดํฌ: ํ์ ์ปฌ๋ผ ๋๋ฝ / row 0๊ฐ๋ฉด ์คํจ."""
required_keys = {"event_id", "user_id", "event_name", "amount"}
count = 0
with open(raw_path) as f:
for line in f:
row = json.loads(line)
missing = required_keys - row.keys()
if missing:
raise ValueError(f"Missing keys: {missing}")
count += 1
if count == 0:
raise ValueError("No data found")
return raw_path
@task
def transform(validated_raw_path: str) -> str:
"""raw JSONL โ curated CSV. (์ค๋ฌด: Spark job / dbt model / SQL transform)"""
ds = get_current_context()["ds"]
curated_dir = BASE_DIR / "s3_like" / "curated" / "event_bus" / f"dt={ds}"
curated_dir.mkdir(parents=True, exist_ok=True)
csv_path = curated_dir / "events.csv"
with open(validated_raw_path) as infile, csv_path.open("w", newline="") as outfile:
writer = csv.DictWriter(
outfile,
fieldnames=["dt", "event_id", "user_id", "event_name", "amount"],
)
writer.writeheader()
for line in infile:
row = json.loads(line)
writer.writerow({"dt": ds, **{k: row[k] for k in
["event_id", "user_id", "event_name", "amount"]}})
return str(csv_path)
@task
def load_to_bq_like(csv_path: str) -> str:
"""BigQuery ์ ์ฌ๋ฅผ SQLite insert๋ก ํ๋ด. (์ค๋ฌด: BigQuery load job / connector)"""
warehouse_dir = BASE_DIR / "warehouse"
warehouse_dir.mkdir(parents=True, exist_ok=True)
db_path = warehouse_dir / "bq_like.db"
conn = sqlite3.connect(db_path)
cur = conn.cursor()
cur.execute(
"""
CREATE TABLE IF NOT EXISTS fact_event (
dt TEXT, event_id TEXT PRIMARY KEY, user_id INTEGER,
event_name TEXT, amount INTEGER
)
"""
)
with open(csv_path) as f:
for row in csv.DictReader(f):
cur.execute(
"""INSERT OR REPLACE INTO fact_event
(dt, event_id, user_id, event_name, amount)
VALUES (?, ?, ?, ?, ?)""",
(row["dt"], row["event_id"], int(row["user_id"]),
row["event_name"], int(row["amount"])),
)
conn.commit()
conn.close()
return str(db_path)
@task
def data_quality_check(db_path: str) -> None:
"""์ ์ฌ ๊ฒฐ๊ณผ ๊ฒ์ฆ. (์ค๋ฌด: row count / null / ์ค๋ณต / partition check)"""
conn = sqlite3.connect(db_path)
cur = conn.cursor()
cur.execute("SELECT COUNT(*) FROM fact_event")
row_count = cur.fetchone()[0]
cur.execute("SELECT SUM(amount) FROM fact_event")
total_amount = cur.fetchone()[0]
conn.close()
if row_count == 0:
raise ValueError("Loaded table is empty")
print(f"row_count={row_count}")
print(f"total_amount={total_amount}")
# ๋ฐํ๊ฐ์ ๋ค์ task ์ธ์๋ก ๋๊ธฐ๋ ๊ฒ = ์์กด์ฑ ์ ์
raw_path = extract_to_s3_like()
validated_path = validate(raw_path)
csv_path = transform(validated_path)
db_path = load_to_bq_like(csv_path)
data_quality_check(db_path)
mini_data_pipeline()์ฐธ๊ณ :
transform์ CSV ์์ฑ๋ถ๋ ์๋ณธ์ ๋ช ์์ writerow๋ฅผ dict ๋ณํฉ์ผ๋ก ์ด์ง ์ค์๋ค. ๋์์ ๋์ผํ๋ค.
3.3 ๊ฐ ๋ถ๋ถ์ด ์ฌ๋ผ์ด๋์ ๋ฌด์์ธ์ง
| DAG ์ฝ๋ | ์ฌ๋ผ์ด๋/์ค๋ฌด์์ | ํต์ฌ ๊ฐ๋ |
|---|---|---|
events = [...] | Event Bus์์ ํ๋ฌ์จ ๋ก๊ทธ | ์ด๋ฒคํธ ์คํธ๋ฆผ |
extract_to_s3_like | S3 raw prefix์ ์ ์ฌ | object, prefix, ํํฐ์
(dt=) |
validate, data_quality_check | Data Quality ์ฒดํฌ | ํ์ ์ปฌ๋ผยทrow countยท์ค๋ณต |
transform | Spark / dbt transform | raw โ curated |
load_to_bq_like | BigQuery load job | schema, idempotent ์ ์ฌ |
์ฌ๊ธฐ์ ํนํ ์ค์ํ ๊ฐ๊ฐ์ idempotent(๋ฉฑ๋ฑ) ๋ค. ๊ฐ์ ๋ ์ง์ DAG๋ฅผ ๋ ๋ฒ ๋๋ ค๋ ๊ฒฐ๊ณผ๊ฐ ๋ง๊ฐ์ง๋ฉด ์ ๋๋ค. ๊ทธ๋์ INSERT OR REPLACE(= BigQuery์ ํํฐ์
insert_overwrite์ ๋์)๋ก ์งฐ๋ค. ์ด๊ฑด DBT + Airflow ํธ์์ Fact ๋ชจ๋ธ์ ๋ ์ง ํํฐ์
๋จ์๋ก ๋ฎ์ด์ฐ๋ ๊ฒ๊ณผ ๊ฐ์ ์๋ฆฌ๋ค.
3.4 Helm๊ณผ Terraform์ ์ด๋์?
์ด ์ค์ต์๋ Helm๋ Terraform๋ ๋ฑ์ฅํ์ง ์๋๋ค. ๋ก์ปฌ Docker๋ก ์ถฉ๋ถํ๊ธฐ ๋๋ฌธ์ด๋ค. ํ์ง๋ง ์ค์ ํ์ฌ์์ ์ด DAG๊ฐ ๋๋ ค๋ฉด ๊ทธ ์๋์ ๋ ์ธต์ด ๋ ์์ด์ผ ํ๋ค.
Terraform = ๋
ยท๊ฑด๋ฌผยท์ ๊ธฐยท์๋ ๋ง๋ค๊ธฐ (EKS ํด๋ฌ์คํฐ, S3, BigQuery dataset, IAM)
Helm = ๊ทธ ๊ฑด๋ฌผ์ Airflow๋ฅผ ์
์ ์ํค๊ธฐ (webserver/scheduler/worker ๋ฐฐํฌยท์ค์ )
Airflow = ์
์ ํ ์๋น์ค๊ฐ ๋งค์ผ ๋ฐ์ดํฐ ์์
์ ๋๋ฆฌ๊ธฐ- Docker Compose = ๋ก์ปฌ์์ Airflow ๋์ฐ๊ธฐ / Helm Chart = ์ฟ ๋ฒ๋คํฐ์ค์์ Airflow ๋์ฐ๊ธฐ.
values.yaml์ด ๊ณง Airflow ์ค์ ๊ฐ ๋ชจ์์ด๊ณ , ๊ทธ ๊ฐ์ดtemplates/*.yaml์ ์ฃผ์ ๋ผ ์ต์ข K8s manifest๊ฐ ๋ง๋ค์ด์ง๋ค. - ์ฒ์์ ์ค์น๋ณด๋ค ๋ ๋๋ง ๊ฒฐ๊ณผ ๊ตฌ๊ฒฝ์ด ์ดํด๊ฐ ๋น ๋ฅด๋ค.
helm repo add apache-airflow https://airflow.apache.org
helm repo update
# Chart๊ฐ ์ค์ ๋ก ์ด๋ค K8s YAML๋ก ๋ณํ๋๋์ง ํ์ธ (ํด๋ฌ์คํฐ ์์ด๋ ๊ฐ๋ฅ)
helm template airflow apache-airflow/airflow --namespace airflow > rendered-airflow.yaml4. ์คํ ๊ฒฐ๊ณผ
๐ง ์๋๋ ๋ด์ผ ๋ก์ปฌ์์ DAG๋ฅผ ์ง์ Triggerํ ๋ค ์คํฌ๋ฆฐ์ท/๋ก๊ทธ๋ก ์ฑ์ด๋ค.
(1) DAG ๋ชฉ๋ก์์ mini_data_pipeline ์ผ๊ธฐ
(2) Graph ๋ทฐ โ ์ด๋ก๋ถ ๐ข
(3) data_quality_check ๋ก๊ทธ
๊ธฐ๋๊ฐ:
row_count=3
total_amount=20000(4) ๋งฅ์ ์ค์ ๋ก ์์ธ ๊ฒฐ๊ณผ๋ฌผ (~/airflow-study/data/)
data/s3_like/raw/event_bus/dt=.../events.jsonl # S3 raw ํ๋ด
data/s3_like/curated/event_bus/dt=.../events.csv # S3 curated ํ๋ด
data/warehouse/bq_like.db # BigQuery ํ๋ด (SQLite)(5) ์ผ๋ถ๋ฌ ์คํจ์์ผ ๋ณด๊ธฐ โ validate์์ ํ์ ์ปฌ๋ผ์ ํ๋ ๋นผ๋ฉด ์ด๋ task๊ฐ ์ ์คํจํ๋์ง ๋ก๊ทธ๋ก ์ถ์ ํ๋ ์ฐ์ต์ด ๋๋ค. (์ค๋ฌด์์ ์ ์ผ ์์ฃผ ํ๋ ์ผ)
5. ๋ฐฐ์ด ๊ฒ / ๋ค์ ๋จ๊ณ
์๊ฒ ํ ๋ฒ ๋๋ฆฌ๊ณ ๋๋ฉด, ๊ฐ ์กฐ๊ฐ์ ์ค์ ๋๊ตฌ๋ก ํ๋์ฉ ์นํํ๋ฉฐ ํ์ฅํ ์ ์๋ค.
๋ก์ปฌ ํด๋ โ ์ง์ง S3
SQLite โ ์ง์ง BigQuery
ํ์ด์ฌ task โ Spark job / SQL / provider operator
Docker Compose โ Helm on Kubernetes
์๋ ํด๋ ์์ฑ โ Terraform“๋ค ์์์ผ ์์"์ด ์๋๋ผ ์๊ฒ ๋๋ฆฌ๊ณ โ ๊ฐ ๋๊ตฌ๊ฐ ์ ํ์ํ์ง ๋ถ์ฌ๊ฐ๋ ์์. ์ฌ๋ผ์ด๋์ “100 of 300"์ด ๋งํ ๊ฒ ๊ฒฐ๊ตญ ์ด๊ฑฐ์๋ค.
์ฐธ๊ณ
- Airflow โ Docs: https://airflow.apache.org/docs/apache-airflow/stable/
- Airflow โ Core Concepts (DAGยทTaskยทOperatorยทXComยทExecutor): https://airflow.apache.org/docs/apache-airflow/stable/core-concepts/index.html
- Airflow โ Running with Docker Compose: https://airflow.apache.org/docs/apache-airflow/stable/howto/docker-compose/index.html
- Airflow โ TaskFlow API tutorial: https://airflow.apache.org/docs/apache-airflow/stable/tutorial/taskflow.html
- Airflow โ Helm Chart: https://airflow.apache.org/docs/helm-chart/stable/index.html
- Helm โ Docs / Charts: https://helm.sh/docs/topics/charts/
- Terraform โ Intro / Providers: https://developer.hashicorp.com/terraform/language/providers
- BigQuery โ Loading data: https://cloud.google.com/bigquery/docs/loading-data
- Amazon S3 โ Organizing objects using prefixes: https://docs.aws.amazon.com/AmazonS3/latest/userguide/using-prefixes.html
- Apache Spark โ Docs: https://spark.apache.org/docs/latest/
- Apache Flink CDC โ Docs: https://nightlies.apache.org/flink/flink-cdc-docs-stable/
- GCP Pub/Sub / Dataflow: https://cloud.google.com/pubsub/docs ยท https://cloud.google.com/dataflow/docs
- Kubernetes โ Docs: https://kubernetes.io/docs/home/
- Karpenter โ Docs: https://karpenter.sh/