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(๊ทธ ๊ธธ์ด ์‹ค์ œ๋กœ ๋Œ์•„๊ฐ€๋Š” ๋ฐ”๋‹ฅ) ์œผ๋กœ ๋‚˜๋‰˜๊ณ , ๊ทธ ์•ˆ์— ๋‹ค์„ฏ ์ถ•์ด ์žˆ์—ˆ๋‹ค.

text
Pipelines โ”€โ”ฌโ”€ Orchestration : Airflow
           โ”œโ”€ Batch         : Spark, RDBMS, NoSQL
           โ””โ”€ Streaming     : Flink, Kafka, Pub/Sub, Dataflow

Infra โ”€โ”€โ”€โ”€โ”€โ”ฌโ”€ Kubernetes, Karpenter
           โ””โ”€ DPaaS         : DT Platform + @

๊ฐ ์ถ•์€ ์•ž์„  ๋…ธํŠธ๋“ค๊ณผ ์ด๋ ‡๊ฒŒ ์—ฐ๊ฒฐ๋œ๋‹ค.

์Šฌ๋ผ์ด๋“œ ์ถ•ํ•ต์‹ฌ ๋„๊ตฌํ•œ ์ค„ ์ •์˜์ด์–ด์ง€๋Š” ๋…ธํŠธ
OrchestrationAirflow์–ธ์ œยท์–ด๋–ค ์ˆœ์„œ๋กœ ์ž‘์—…์„ ์‹คํ–‰ํ• ์ง€ ๊ด€๋ฆฌDBT + Airflow ํŽธ
BatchSpark / RDBMS / NoSQL์Œ“์ธ ๋ฐ์ดํ„ฐ๋ฅผ ์ฃผ๊ธฐ์ ์œผ๋กœ ํ•œ ๋ฒˆ์— ์ฒ˜๋ฆฌDT Platform ํŽธ
StreamingFlink / Kafka / Pub/Sub / Dataflow๋ฐ์ดํ„ฐ๊ฐ€ ์ƒ๊ธฐ๋Š” ์ฆ‰์‹œ ํ˜๋ ค๋ณด๋‚ด๋ฉฐ ์ฒ˜๋ฆฌMongoDB CDC ํŽธ
InfraKubernetes / Karpenter์œ„ ์ž‘์—…๋“ค์ด ์‹ค์ œ๋กœ ์‹คํ–‰๋˜๋Š” ํด๋Ÿฌ์Šคํ„ฐโ€”
DPaaSDT Platform์œ„ ์ „๋ถ€๋ฅผ ๊ตฌ์„ฑ์›์ด ์‰ฝ๊ฒŒ ์“ฐ๊ฒŒ ๋งŒ๋“  ๋‚ด๋ถ€ ํ”Œ๋žซํผDT Platform ํŽธ

ํ•œ ์ค„๋กœ ๋‹ค์‹œ ๋ฌถ์œผ๋ฉด:

text
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(๋น„๊ด€๊ณ„ํ˜•).
  • 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 โ€” ๋กœ์ปฌ์— ๋ฏธ๋‹ˆ ํŒŒ์ดํ”„๋ผ์ธ ์„ธ์šฐ๊ธฐ

ํ•ต์‹ฌ ์•„์ด๋””์–ด๋Š” ์ง„์งœ ํด๋ผ์šฐ๋“œ๋Š” ์•ˆ ์“ฐ๋˜ ๊ตฌ์กฐ๋Š” ํšŒ์‚ฌ์™€ ๋˜‘๊ฐ™์ด ๊ฐ€์ ธ๊ฐ€๋Š” ๊ฒƒ์ด๋‹ค. ์ด๋ ‡๊ฒŒ ์น˜ํ™˜ํ–ˆ๋‹ค.

text
๋กœ์ปฌ ํด๋” (ํŒŒ์ผ ์“ฐ๊ธฐ)     โ†’  S3 (raw / curated)
SQLite ํ…Œ์ด๋ธ”             โ†’  BigQuery (fact_event)
Airflow DAG              โ†’  ์‹ค์ œ ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ด์…˜ (๊ทธ๋Œ€๋กœ)
ํŒŒ์ด์ฌ ๋ฆฌ์ŠคํŠธ๋กœ ๋งŒ๋“  ์ด๋ฒคํŠธ  โ†’  Event Bus์—์„œ ํ˜๋Ÿฌ์˜จ ๋กœ๊ทธ

ํŒŒ์ดํ”„๋ผ์ธ ๋‹จ๊ณ„๋Š” ์‹ค์ œ ๋ฐฐ์น˜ ํŒŒ์ดํ”„๋ผ์ธ๊ณผ ๊ฐ™์€ ๊ณจ๊ฒฉ์ด๋‹ค.

text
extract_to_s3_like โ†’ validate โ†’ transform โ†’ load_to_bq_like โ†’ data_quality_check

3.1 ํ™˜๊ฒฝ: Docker Compose๋กœ Airflow ๋„์šฐ๊ธฐ

Airflow ๊ณต์‹ ๋ฌธ์„œ์˜ Docker Compose quick start๋ฅผ ์‚ฌ์šฉํ•œ๋‹ค. (ํ”„๋กœ๋•์…˜์šฉ์ด ์•„๋‹ˆ๋ผ ๋กœ์ปฌ ํ•™์Šต์šฉ)

bash
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์„ ํ•œ ์ค„ ์ถ”๊ฐ€ํ•œ๋‹ค.

yaml
# 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์˜ ์ธ์ž๋กœ ๋„˜์–ด๊ฐ€๋Š” ๊ฒƒ์ด ๊ณง ์˜์กด์„ฑ์ด ๋œ๋‹ค.

python
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_likeS3 raw prefix์— ์ ์žฌobject, prefix, ํŒŒํ‹ฐ์…˜(dt=)
validate, data_quality_checkData Quality ์ฒดํฌํ•„์ˆ˜ ์ปฌ๋Ÿผยทrow countยท์ค‘๋ณต
transformSpark / dbt transformraw โ†’ curated
load_to_bq_likeBigQuery load jobschema, idempotent ์ ์žฌ

์—ฌ๊ธฐ์„œ ํŠนํžˆ ์ค‘์š”ํ•œ ๊ฐ๊ฐ์€ idempotent(๋ฉฑ๋“ฑ) ๋‹ค. ๊ฐ™์€ ๋‚ ์งœ์˜ DAG๋ฅผ ๋‘ ๋ฒˆ ๋Œ๋ ค๋„ ๊ฒฐ๊ณผ๊ฐ€ ๋ง๊ฐ€์ง€๋ฉด ์•ˆ ๋œ๋‹ค. ๊ทธ๋ž˜์„œ INSERT OR REPLACE(= BigQuery์˜ ํŒŒํ‹ฐ์…˜ insert_overwrite์— ๋Œ€์‘)๋กœ ์งฐ๋‹ค. ์ด๊ฑด DBT + Airflow ํŽธ์—์„œ Fact ๋ชจ๋ธ์„ ๋‚ ์งœ ํŒŒํ‹ฐ์…˜ ๋‹จ์œ„๋กœ ๋ฎ์–ด์“ฐ๋˜ ๊ฒƒ๊ณผ ๊ฐ™์€ ์›๋ฆฌ๋‹ค.

3.4 Helm๊ณผ Terraform์€ ์–ด๋””์—?

์ด ์‹ค์Šต์—๋Š” Helm๋„ Terraform๋„ ๋“ฑ์žฅํ•˜์ง€ ์•Š๋Š”๋‹ค. ๋กœ์ปฌ Docker๋กœ ์ถฉ๋ถ„ํ•˜๊ธฐ ๋•Œ๋ฌธ์ด๋‹ค. ํ•˜์ง€๋งŒ ์‹ค์ œ ํšŒ์‚ฌ์—์„œ ์ด DAG๊ฐ€ ๋Œ๋ ค๋ฉด ๊ทธ ์•„๋ž˜์— ๋‘ ์ธต์ด ๋” ์žˆ์–ด์•ผ ํ•œ๋‹ค.

text
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๊ฐ€ ๋งŒ๋“ค์–ด์ง„๋‹ค.
  • ์ฒ˜์Œ์—” ์„ค์น˜๋ณด๋‹ค ๋ Œ๋”๋ง ๊ฒฐ๊ณผ ๊ตฌ๊ฒฝ์ด ์ดํ•ด๊ฐ€ ๋น ๋ฅด๋‹ค.
bash
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.yaml

4. ์‹คํ—˜ ๊ฒฐ๊ณผ

๐Ÿšง ์•„๋ž˜๋Š” ๋‚ด์ผ ๋กœ์ปฌ์—์„œ DAG๋ฅผ ์ง์ ‘ Triggerํ•œ ๋’ค ์Šคํฌ๋ฆฐ์ƒท/๋กœ๊ทธ๋กœ ์ฑ„์šด๋‹ค.

(1) DAG ๋ชฉ๋ก์—์„œ mini_data_pipeline ์ผœ๊ธฐ

(2) Graph ๋ทฐ โ€” ์ดˆ๋ก๋ถˆ ๐ŸŸข

(3) data_quality_check ๋กœ๊ทธ

๊ธฐ๋Œ€๊ฐ’:

text
row_count=3
total_amount=20000

(4) ๋งฅ์— ์‹ค์ œ๋กœ ์Œ“์ธ ๊ฒฐ๊ณผ๋ฌผ (~/airflow-study/data/)

text
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. ๋ฐฐ์šด ๊ฒƒ / ๋‹ค์Œ ๋‹จ๊ณ„

์ž‘๊ฒŒ ํ•œ ๋ฒˆ ๋Œ๋ฆฌ๊ณ  ๋‚˜๋ฉด, ๊ฐ ์กฐ๊ฐ์„ ์‹ค์ œ ๋„๊ตฌ๋กœ ํ•˜๋‚˜์”ฉ ์น˜ํ™˜ํ•˜๋ฉฐ ํ™•์žฅํ•  ์ˆ˜ ์žˆ๋‹ค.

text
๋กœ์ปฌ ํด๋”        โ†’  ์ง„์งœ S3
SQLite          โ†’  ์ง„์งœ BigQuery
ํŒŒ์ด์ฌ task      โ†’  Spark job / SQL / provider operator
Docker Compose  โ†’  Helm on Kubernetes
์ˆ˜๋™ ํด๋” ์ƒ์„ฑ    โ†’  Terraform

“๋‹ค ์•Œ์•„์•ผ ์‹œ์ž‘"์ด ์•„๋‹ˆ๋ผ ์ž‘๊ฒŒ ๋Œ๋ฆฌ๊ณ  โ†’ ๊ฐ ๋„๊ตฌ๊ฐ€ ์™œ ํ•„์š”ํ•œ์ง€ ๋ถ™์—ฌ๊ฐ€๋Š” ์ˆœ์„œ. ์Šฌ๋ผ์ด๋“œ์˜ “100 of 300"์ด ๋งํ•œ ๊ฒŒ ๊ฒฐ๊ตญ ์ด๊ฑฐ์˜€๋‹ค.

์ฐธ๊ณ 

Discussion