8.1 Kafka + Spark Streaming 실시간 파이프라인
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
spark = SparkSession.builder.appName("RealTimePipeline").getOrCreate()
# Kafka에서 스트림 읽기
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "user-events") \
.load()
# JSON 파싱
events = df.selectExpr("CAST(value AS STRING) as json") \
.select(from_json("json", schema).alias("data")) \
.select("data.*")
# 변환
processed = events \
.filter(col("event_type") == "purchase") \
.withColumn("hour", date_format("timestamp", "yyyy-MM-dd-HH")) \
.groupBy("hour", "product_id") \
.agg(
count("*").alias("purchases"),
sum("amount").alias("revenue")
)
# 출력
query = processed \
.writeStream \
.outputMode("complete") \
.format("console") \
.start()
query.awaitTermination()
8.2 Change Data Capture (CDC) 파이프라인
from debezium import DebeziumEngine
# Debezium 설정
config = {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "localhost",
"database.port": "3306",
"database.user": "root",
"database.password": "password",
"database.server.name": "myapp"
}
def handle_change(event):
operation = event['op'] # 'c'=create, 'u'=update, 'd'=delete
table = event['source']['table']
data = event['after'] # 변경 후 데이터
if operation == 'c':
insert_to_warehouse(table, data)
elif operation == 'u':
update_warehouse(table, data)
elif operation == 'd':
mark_deleted(table, event['before']['id'])
# 실행
engine = DebeziumEngine(config)
engine.run(handle_change)
8.3 Lambda Architecture
# Batch Layer (정확성)
def batch_processing():
df = spark.read.parquet("s3://lake/raw/")
daily_stats = df \
.groupBy("date", "user_id") \
.agg(sum("amount").alias("total"))
daily_stats.write.mode("overwrite").parquet("s3://lake/batch/")
# Speed Layer (실시간성)
def stream_processing():
stream = spark.readStream.format("kafka").load()
near_realtime = stream \
.groupBy(window("timestamp", "5 minutes"), "user_id") \
.agg(sum("amount").alias("total"))
near_realtime.writeStream.format("delta").start()
# Serving Layer (결합)
def serve_query(user_id):
batch_result = query_batch_layer(user_id)
stream_result = query_speed_layer(user_id)
return merge(batch_result, stream_result)
8.4 Data Mesh 패턴
도메인 중심의 분산 데이터 아키텍처:
# 각 도메인이 자신의 파이프라인 소유
class OrdersDomain:
def provide_orders_data(self):
"""주문 데이터 제공 (Data Product)"""
return self.pipeline.run()
def pipeline(self):
extract = self.extract_from_oltp()
clean = self.clean_orders(extract)
enrich = self.enrich_with_user_data(clean)
return self.publish_to_data_catalog(enrich)
class UsersDomain:
def provide_users_data(self):
"""사용자 데이터 제공"""
return self.pipeline.run()
8.5 CI/CD for Data Pipelines
GitHub Actions
# .github/workflows/pipeline-test.yml
name: Test Data Pipeline
on: [push, pull_request]
jobs:
test:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v2
- name: Set up Python
uses: actions/setup-python@v2
with:
python-version: '3.9'
- name: Install dependencies
run: pip install -r requirements.txt
- name: Run unit tests
run: pytest tests/
- name: Run dbt tests
run: dbt test
- name: Validate schemas
run: python validate_schemas.py
8.6 비용 최적화 전략
1. 파티셔닝 최적화
-- 날짜 파티션으로 스캔 범위 축소
SELECT *
FROM orders
WHERE date BETWEEN '2025-12-01' AND '2025-12-31';
-- ✅ 1개월 파티션만 스캔
-- ❌ 전체 테이블 스캔 방지
2. 압축 포맷 사용
# Parquet: 컬럼 기반 + 압축
df.write.format("parquet").save("s3://bucket/data/")
# 비용 절감:
# CSV: 10GB → Parquet: 2GB (80% 감소)
3. 증분 처리
-- 전체 재처리 ❌
SELECT * FROM large_table;
-- 증분 처리 ✅
SELECT * FROM large_table
WHERE updated_at > (SELECT MAX(updated_at) FROM processed);
8.7 스케일링 전략
수직 스케일링
- 더 큰 인스턴스 사용
- 더 많은 메모리, CPU
- 한계: 단일 노드 최대 용량
수평 스케일링
- 더 많은 워커 노드 추가
- Spark, Flink 클러스터
- 무한 확장 가능
8.8 실전 체크리스트
Production 배포 전
- ✅ 모든 테스트 통과
- ✅ 에러 핸들링 구현
- ✅ 로깅 추가
- ✅ 모니터링 설정
- ✅ 알림 구성
- ✅ 멱등성 보장
- ✅ 백업 전략
- ✅ 롤백 계획
- ✅ 문서 작성
- ✅ On-call 준비
8.9 베스트 프랙티스 요약
1. 설계 원칙
- 멱등성: 재실행 가능하게
- 격리성: 각 단계를 독립적으로
- 관찰성: 모든 것을 로깅
- 복원력: 실패해도 복구
2. 코드 품질
- 버전 관리 (Git)
- 코드 리뷰
- 자동화된 테스트
- CI/CD 파이프라인
3. 운영
- 모니터링 및 알림
- 문서화
- On-call 로테이션
- Post-mortem 분석
8.10 다음 단계
이제 여러분은 데이터 파이프라인의 모든 것을 배웠습니다:
- ✅ Extract: 데이터 수집
- ✅ Transform: 데이터 변환
- ✅ Load: 데이터 적재
- ✅ Orchestration: 워크플로우 관리
- ✅ Monitoring: 관찰성
- ✅ Quality: 데이터 품질
- ✅ Advanced: 고급 패턴
8.9 Kappa 아키텍처
Lambda의 복잡성을 제거한 스트리밍 전용 아키텍처:
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
spark = SparkSession.builder.appName("KappaArchitecture").getOrCreate()
# 단일 스트리밍 레이어로 모든 처리
stream = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "events") \
.load()
# 실시간 집계
real_time_stats = stream \
.groupBy(window("timestamp", "1 hour"), "user_id") \
.agg(
count("*").alias("event_count"),
sum("amount").alias("total_amount")
)
# Delta Lake에 저장 (시계열 형태)
real_time_stats.writeStream \
.format("delta") \
.outputMode("append") \
.option("checkpointLocation", "/checkpoints/stats") \
.start("/data/hourly_stats")
# 과거 데이터 재처리가 필요하면 같은 코드 재사용
historical_stats = spark.read \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("startingOffsets", "earliest") \
.option("endingOffsets", "latest") \
.load() \
.groupBy(window("timestamp", "1 hour"), "user_id") \
.agg(
count("*").alias("event_count"),
sum("amount").alias("total_amount")
)
historical_stats.write.format("delta").mode("overwrite").save("/data/hourly_stats")
Lambda vs Kappa 비교
| 항목 | Lambda | Kappa |
|---|---|---|
| 레이어 | Batch + Speed | Stream만 |
| 복잡도 | 높음 | 낮음 |
| 코드 중복 | 있음 | 없음 |
| 레이턴시 | 배치 레이어 느림 | 모두 빠름 |
| 재처리 | 배치로 자동 | 수동 재처리 |
| 적용 사례 | 대규모 배치 + 실시간 | 순수 스트리밍 |
8.10 Feature Store
머신러닝을 위한 피처 저장소는 데이터 파이프라인의 고급 활용 사례입니다.
Feast Feature Store
# feature_repo/features.py
from feast import Entity, Feature, FeatureView, Field
from feast.types import Float32, Int64
from datetime import timedelta
# Entity 정의
user = Entity(
name="user_id",
join_keys=["user_id"],
description="User entity"
)
# Feature View 정의
user_features = FeatureView(
name="user_features",
entities=[user],
ttl=timedelta(days=1),
schema=[
Field(name="total_purchases", dtype=Int64),
Field(name="avg_order_value", dtype=Float32),
Field(name="days_since_last_purchase", dtype=Int64),
],
online=True,
source=BatchSource(
path="s3://bucket/user_features/",
timestamp_field="event_timestamp"
)
)
# 피처 저장소 초기화
# feast init feature_repo
# feast apply
# 피처 수집 (배치)
from feast import FeatureStore
store = FeatureStore(repo_path="feature_repo")
# 훈련용 피처 가져오기
training_df = store.get_historical_features(
entity_df=entity_df,
features=[
"user_features:total_purchases",
"user_features:avg_order_value",
"user_features:days_since_last_purchase"
]
).to_df()
# 온라인 서빙
features = store.get_online_features(
features=[
"user_features:total_purchases",
"user_features:avg_order_value"
],
entity_rows=[{"user_id": 123}]
).to_dict()
피처 파이프라인
from airflow import DAG
from airflow.operators.python import PythonOperator
def compute_user_features():
"""사용자 피처 계산"""
df = spark.sql("""
SELECT
user_id,
COUNT(*) as total_purchases,
AVG(total) as avg_order_value,
DATEDIFF(CURRENT_DATE(), MAX(order_date)) as days_since_last_purchase,
CURRENT_TIMESTAMP() as event_timestamp
FROM orders
WHERE order_date >= DATE_SUB(CURRENT_DATE(), 90)
GROUP BY user_id
""")
# Feature Store에 저장
df.write.format("parquet").mode("overwrite").save("s3://bucket/user_features/")
with DAG('feature_pipeline', schedule_interval='@daily') as dag:
compute_features = PythonOperator(
task_id='compute_user_features',
python_callable=compute_user_features
)
8.11 CI/CD for Data Pipelines
완전한 GitHub Actions 워크플로우
# .github/workflows/pipeline-ci-cd.yml
name: Data Pipeline CI/CD
on:
push:
branches: [main, develop]
pull_request:
branches: [main]
env:
PYTHON_VERSION: '3.9'
jobs:
lint:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v3
- name: Set up Python
uses: actions/setup-python@v4
with:
python-version: ${{ env.PYTHON_VERSION }}
- name: Install dependencies
run: |
pip install flake8 black isort
- name: Run linters
run: |
flake8 pipelines/ --max-line-length=100
black --check pipelines/
isort --check-only pipelines/
test:
runs-on: ubuntu-latest
services:
postgres:
image: postgres:14
env:
POSTGRES_PASSWORD: postgres
options: >-
--health-cmd pg_isready
--health-interval 10s
--health-timeout 5s
--health-retries 5
ports:
- 5432:5432
steps:
- uses: actions/checkout@v3
- name: Set up Python
uses: actions/setup-python@v4
with:
python-version: ${{ env.PYTHON_VERSION }}
- name: Install dependencies
run: |
pip install -r requirements.txt
pip install pytest pytest-cov
- name: Run unit tests
env:
DATABASE_URL: postgresql://postgres:postgres@localhost:5432/test
run: |
pytest tests/unit --cov=pipelines --cov-report=xml
- name: Run integration tests
env:
DATABASE_URL: postgresql://postgres:postgres@localhost:5432/test
run: |
pytest tests/integration
- name: Upload coverage
uses: codecov/codecov-action@v3
with:
file: ./coverage.xml
dbt-test:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v3
- name: Set up Python
uses: actions/setup-python@v4
with:
python-version: ${{ env.PYTHON_VERSION }}
- name: Install dbt
run: pip install dbt-core dbt-postgres
- name: dbt deps
run: dbt deps
- name: dbt compile
run: dbt compile
- name: dbt test
run: dbt test
deploy-dev:
runs-on: ubuntu-latest
needs: [lint, test, dbt-test]
if: github.ref == 'refs/heads/develop'
steps:
- uses: actions/checkout@v3
- name: Deploy to Development
run: |
# Airflow DAGs 배포
aws s3 sync dags/ s3://airflow-dev-dags/
# dbt 모델 배포
dbt run --target dev
deploy-prod:
runs-on: ubuntu-latest
needs: [lint, test, dbt-test]
if: github.ref == 'refs/heads/main'
environment:
name: production
url: https://airflow.company.com
steps:
- uses: actions/checkout@v3
- name: Deploy to Production
run: |
# Airflow DAGs 배포
aws s3 sync dags/ s3://airflow-prod-dags/
# dbt 모델 배포
dbt run --target prod
# Slack 알림
curl -X POST ${{ secrets.SLACK_WEBHOOK }} \
-H 'Content-Type: application/json' \
-d '{"text":"Pipeline deployed to production successfully!"}'
8.12 Data Mesh 구현
도메인별 데이터 제품
# domains/orders/pipeline.py
class OrdersDataProduct:
"""주문 도메인 데이터 제품"""
def __init__(self):
self.domain = "orders"
self.owner = "orders-team@company.com"
self.sla = {
"freshness": timedelta(hours=1),
"availability": 0.999,
"quality": 0.99
}
def extract(self):
"""소스에서 주문 데이터 추출"""
return spark.read.jdbc(
url="jdbc:mysql://orders-db:3306/orders",
table="orders",
properties={"user": "reader", "password": "xxx"}
)
def transform(self, df):
"""비즈니스 로직 적용"""
return df \
.filter(col("status") != "cancelled") \
.withColumn("revenue", col("quantity") * col("price")) \
.withColumn("order_hour", hour("order_timestamp"))
def publish(self, df):
"""데이터 제품 발행"""
# 1. Delta Lake에 저장
df.write.format("delta") \
.mode("overwrite") \
.save("s3://data-products/orders/")
# 2. 메타데이터 등록
self.register_metadata({
"domain": self.domain,
"table": "orders",
"schema": df.schema.json(),
"row_count": df.count(),
"updated_at": datetime.now().isoformat()
})
# 3. 데이터 계약 검증
self.validate_contract(df)
def validate_contract(self, df):
"""데이터 계약 검증"""
contract = {
"required_columns": ["order_id", "user_id", "total", "order_timestamp"],
"constraints": {
"order_id": "unique, not null",
"total": ">= 0"
}
}
# 검증 로직
for col in contract["required_columns"]:
if col not in df.columns:
raise ValueError(f"Missing required column: {col}")
def run(self):
"""전체 파이프라인 실행"""
raw = self.extract()
transformed = self.transform(raw)
self.publish(transformed)
return transformed
# 도메인 팀이 독립적으로 실행
orders_product = OrdersDataProduct()
orders_product.run()
8.13 Medallion 아키텍처
Bronze, Silver, Gold 레이어로 데이터 품질을 점진적으로 개선:
Bronze Layer (Raw)
# 원본 데이터 그대로 저장
bronze = spark.read.json("s3://raw/orders/")
bronze.write.format("delta") \
.mode("append") \
.save("s3://lake/bronze/orders/")
Silver Layer (Cleaned)
# 클렌징 및 표준화
silver = spark.read.format("delta").load("s3://lake/bronze/orders/") \
.filter(col("order_id").isNotNull()) \
.dropDuplicates(["order_id"]) \
.withColumn("order_date", to_date("order_timestamp")) \
.withColumn("total", col("total").cast("decimal(10,2)"))
silver.write.format("delta") \
.mode("overwrite") \
.partitionBy("order_date") \
.save("s3://lake/silver/orders/")
Gold Layer (Business)
# 비즈니스 로직 적용
gold = spark.read.format("delta").load("s3://lake/silver/orders/") \
.join(users, "user_id") \
.join(products, "product_id") \
.groupBy("order_date", "category") \
.agg(
sum("total").alias("revenue"),
count("order_id").alias("order_count"),
avg("total").alias("avg_order_value")
)
gold.write.format("delta") \
.mode("overwrite") \
.save("s3://lake/gold/daily_sales_by_category/")
| 레이어 | 목적 | 특징 | 사용자 |
|---|---|---|---|
| Bronze | 원본 보관 | Raw, 중복 가능, 스키마 느슨 | 데이터 엔지니어 |
| Silver | 정제 데이터 | 중복 제거, 타입 변환, 유효성 검증 | 데이터 분석가 |
| Gold | 비즈니스 로직 | 집계, 조인, 비즈니스 룰 적용 | 비즈니스 사용자, BI |
8.14 실시간 Feature Engineering
from pyspark.sql.streaming import StreamingQuery
# 스트림에서 피처 계산
stream = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "user-events") \
.load()
# 윈도우 집계로 피처 생성
features = stream \
.groupBy(
col("user_id"),
window("timestamp", "30 minutes", "5 minutes")
) \
.agg(
count("event_type").alias("event_count_30m"),
countDistinct("page_url").alias("unique_pages_30m"),
avg("session_duration").alias("avg_session_30m")
)
# Feature Store에 실시간 쓰기
query = features.writeStream \
.foreachBatch(lambda df, epoch: df.write \
.format("redis") \
.option("table", "user_features") \
.option("key.column", "user_id") \
.mode("append") \
.save()
) \
.start()
8.15 성능 벤치마킹
벤치마크 프레임워크
import time
import psutil
import pandas as pd
class PipelineBenchmark:
def __init__(self, name):
self.name = name
self.results = []
def benchmark(self, func, *args, **kwargs):
"""함수 성능 측정"""
# 시작 시간 및 리소스
start_time = time.time()
start_cpu = psutil.cpu_percent()
start_mem = psutil.virtual_memory().percent
# 실행
result = func(*args, **kwargs)
# 종료 시간 및 리소스
end_time = time.time()
end_cpu = psutil.cpu_percent()
end_mem = psutil.virtual_memory().percent
# 결과 저장
self.results.append({
'function': func.__name__,
'duration_sec': end_time - start_time,
'cpu_delta': end_cpu - start_cpu,
'mem_delta': end_mem - start_mem
})
return result
def report(self):
"""벤치마크 리포트"""
df = pd.DataFrame(self.results)
print(f"\n=== {self.name} Benchmark Report ===")
print(df.to_string(index=False))
return df
# 사용 예시
benchmark = PipelineBenchmark("ETL Pipeline")
raw_data = benchmark.benchmark(extract_data)
cleaned = benchmark.benchmark(transform_data, raw_data)
benchmark.benchmark(load_data, cleaned)
benchmark.report()
8.16 모범 사례 종합
1. 아키텍처 원칙
- 느슨한 결합: 각 컴포넌트 독립적으로 개발/배포
- 높은 응집도: 관련 기능을 함께 그룹화
- 확장성: 수평 확장 가능하게 설계
- 복원력: 장애 시 자동 복구
2. 코드 품질
- 타입 힌트: Python 타입 어노테이션 사용
- 문서화: Docstring, README, 아키텍처 다이어그램
- 테스트: 80% 이상 커버리지 목표
- 코드 리뷰: 최소 1명 이상 승인 필요
3. 운영
- 모니터링: 모든 파이프라인 메트릭 수집
- 알림: 심각도별 알림 채널 분리
- 로깅: 구조화된 로그, 추적 ID 포함
- 백업: 주요 데이터 자동 백업
4. 보안
- 인증: IAM, OAuth2 사용
- 암호화: 전송/저장 시 암호화
- 비밀 관리: AWS Secrets Manager, Vault
- 감사: 모든 데이터 접근 로깅
요약
이 장에서는 데이터 파이프라인의 고급 패턴과 실전 사례를 배웠습니다:
- Lambda 아키텍처: 배치 + 스트림 레이어
- Kappa 아키텍처: 스트리밍 전용, 낮은 복잡도
- Feature Store: ML 피처 중앙 관리
- CI/CD: 자동화된 테스트와 배포
- Data Mesh: 도메인 중심 데이터 제품
- Medallion: Bronze-Silver-Gold 레이어
- CDC: 변경 데이터 캡처
- 성능 최적화: 파티셔닝, 압축, 증분 처리
- 비용 최적화: 리소스 효율적 사용
- 실전 체크리스트: 프로덕션 배포 준비
프로덕션 체크리스트
- ✅ 모든 테스트 통과 (단위, 통합, E2E)
- ✅ 에러 핸들링 구현 (재시도, 알림)
- ✅ 구조화된 로깅 추가
- ✅ 모니터링 설정 (메트릭, 대시보드)
- ✅ 알림 구성 (Slack, PagerDuty)
- ✅ 멱등성 보장
- ✅ 백업 전략 (자동 백업, 복구 테스트)
- ✅ 롤백 계획 (버전 관리, 배포 스크립트)
- ✅ 문서 작성 (README, 아키텍처, Runbook)
- ✅ On-call 준비 (연락처, 에스컬레이션)
- ✅ 성능 테스트 (부하 테스트, 벤치마크)
- ✅ 보안 검토 (인증, 암호화, 접근 제어)
복습 문제
- Lambda 아키텍처와 Kappa 아키텍처의 차이점과 각각의 장단점을 비교하세요.
- Feature Store가 필요한 이유와 주요 구성 요소를 설명하세요.
- Data Mesh의 핵심 원칙 4가지를 설명하세요.
- Medallion 아키텍처의 Bronze, Silver, Gold 레이어 각각의 역할은?
- CDC(Change Data Capture)의 동작 원리와 활용 사례를 설명하세요.
- 데이터 파이프라인 CI/CD에서 반드시 포함되어야 할 단계는?
- 비용 최적화를 위한 3가지 이상의 전략을 제시하세요.
- 프로덕션 배포 전 검증해야 할 항목들을 나열하세요.
실습 과제
- Kafka + Spark Streaming으로 실시간 파이프라인을 구현하세요.
- Lambda 아키텍처 또는 Kappa 아키텍처로 전체 시스템을 설계하세요.
- GitHub Actions로 CI/CD 파이프라인을 구축하세요.
- Medallion 아키텍처로 Bronze-Silver-Gold 레이어를 구현하세요.
마치며
축하합니다! 이제 여러분은 데이터 파이프라인의 모든 것을 배웠습니다:
- ✅ 1장: 데이터 파이프라인 개요
- ✅ 2장: Extract - 데이터 수집
- ✅ 3장: Transform - 데이터 변환
- ✅ 4장: Load - 데이터 로딩
- ✅ 5장: Orchestration - 워크플로우 관리
- ✅ 6장: Monitoring - 관찰성
- ✅ 7장: Quality - 데이터 품질
- ✅ 8장: Advanced - 고급 패턴
弘益人間 (홍익人間) - 널리 인간을 이롭게 하라
데이터 파이프라인은 단순한 코드가 아닙니다. 올바른 데이터를 올바른 시간에 올바른 사람에게 전달함으로써, 우리는 더 나은 결정을 내리고, 더 나은 제품을 만들며, 궁극적으로 인류를 이롭게 할 수 있습니다.
여러분이 만드는 파이프라인이 세상을 더 나은 곳으로 만들기를 바랍니다.
- WIA (World Certification Industry Association)
다음 단계
- 실습 프로젝트: 실제 데이터로 엔드투엔드 파이프라인 구축
- 오픈소스 기여: Airflow, Spark 등 프로젝트에 참여
- 커뮤니티: Data Engineering 커뮤니티 가입
- 인증: AWS Certified Data Analytics, Google Professional Data Engineer
- 심화 학습: Distributed Systems, Stream Processing
Happy Data Engineering!