3.1 데이터 변환이란?
데이터 변환(Transformation)은 원시 데이터를 분석 가능한 형태로 만드는 과정입니다. 다이아몬드 원석을 보석으로 다듬는 것과 같습니다.
변환의 목표
데이터를 정확하고, 일관되며, 유용한 형태로 만드는 것
3.2 변환의 5가지 유형
1. 정제 (Cleaning)
잘못되거나 불완전한 데이터 수정:
-- NULL 값 처리
UPDATE users SET country = 'Unknown' WHERE country IS NULL;
-- 중복 제거
WITH ranked AS (
SELECT *, ROW_NUMBER() OVER (PARTITION BY email ORDER BY created_at DESC) as rn
FROM users
)
DELETE FROM users WHERE id IN (
SELECT id FROM ranked WHERE rn > 1
);
-- 이상치 제거
DELETE FROM orders WHERE total < 0 OR total > 1000000;
2. 정규화 (Normalization)
일관된 형식으로 변환:
-- 전화번호 정규화
UPDATE users
SET phone = REGEXP_REPLACE(phone, '[^0-9]', '', 'g');
-- 이메일 소문자화
UPDATE users
SET email = LOWER(email);
-- 날짜 형식 통일
UPDATE events
SET event_date = TO_TIMESTAMP(event_date_str, 'YYYY-MM-DD HH24:MI:SS');
3. 보강 (Enrichment)
외부 데이터로 보완:
-- 사용자 정보 추가
SELECT
o.*,
u.name,
u.email,
u.segment
FROM orders o
LEFT JOIN users u ON o.user_id = u.id;
-- 지역 정보 추가
SELECT
s.*,
g.city,
g.country,
g.latitude,
g.longitude
FROM sales s
LEFT JOIN geo_lookup g ON s.postal_code = g.postal_code;
4. 집계 (Aggregation)
요약 통계 계산:
-- 일별 매출 집계
SELECT
DATE(order_date) as date,
COUNT(*) as order_count,
SUM(total) as revenue,
AVG(total) as avg_order_value
FROM orders
GROUP BY DATE(order_date);
-- 사용자별 활동 집계
SELECT
user_id,
COUNT(*) as login_count,
MIN(login_at) as first_login,
MAX(login_at) as last_login,
COUNT(DISTINCT DATE(login_at)) as active_days
FROM logins
GROUP BY user_id;
5. 필터링 (Filtering)
필요한 데이터만 선택:
-- 활성 사용자만
SELECT * FROM users
WHERE last_login_at > NOW() - INTERVAL '30 days'
AND status = 'active';
-- 유효한 주문만
SELECT * FROM orders
WHERE total > 0
AND status IN ('completed', 'shipped')
AND created_at > '2025-01-01';
3.3 SQL 기반 변환
dbt (data build tool)
현대 데이터 변환의 표준:
-- models/mart/fct_daily_sales.sql
{{
config(
materialized='incremental',
unique_key='date'
)
}}
SELECT
DATE(order_date) as date,
COUNT(DISTINCT user_id) as customers,
COUNT(*) as orders,
SUM(total) as revenue,
AVG(total) as avg_order_value
FROM {{ ref('stg_orders') }}
{% if is_incremental() %}
WHERE order_date > (SELECT MAX(date) FROM {{ this }})
{% endif %}
GROUP BY DATE(order_date)
dbt의 장점:
- SQL로 모든 변환 정의
- 버전 관리 (Git)
- 테스트 내장
- 문서 자동 생성
- 증분 업데이트 지원
3.4 코드 기반 변환
Pandas (Python)
import pandas as pd
# 데이터 읽기
df = pd.read_csv('raw_data.csv')
# 정제
df = df.dropna(subset=['user_id']) # NULL 제거
df = df.drop_duplicates(subset=['email']) # 중복 제거
# 정규화
df['email'] = df['email'].str.lower()
df['phone'] = df['phone'].str.replace(r'[^0-9]', '', regex=True)
# 파생 컬럼
df['age'] = (pd.Timestamp.now() - pd.to_datetime(df['birth_date'])).dt.days // 365
df['is_vip'] = df['total_spent'] > 10000
# 집계
daily_sales = df.groupby(df['order_date'].dt.date).agg({
'order_id': 'count',
'total': ['sum', 'mean']
})
# 저장
df.to_parquet('cleaned_data.parquet')
PySpark (대용량 데이터)
from pyspark.sql import functions as F
# 데이터 읽기
df = spark.read.parquet('s3://bucket/raw/')
# 정제 및 변환
df_clean = (df
.filter(F.col('user_id').isNotNull())
.dropDuplicates(['email'])
.withColumn('email', F.lower(F.col('email')))
.withColumn('age',
(F.datediff(F.current_date(), F.col('birth_date')) / 365).cast('int')
)
.withColumn('is_vip', F.col('total_spent') > 10000)
)
# 집계
daily_sales = (df_clean
.groupBy(F.to_date('order_date').alias('date'))
.agg(
F.count('*').alias('order_count'),
F.sum('total').alias('revenue'),
F.avg('total').alias('avg_order_value')
)
)
# 저장
daily_sales.write.mode('overwrite').parquet('s3://bucket/mart/daily_sales')
3.5 복잡한 변환 예제
Customer Lifetime Value (CLV) 계산
WITH user_orders AS (
SELECT
user_id,
MIN(order_date) as first_order,
MAX(order_date) as last_order,
COUNT(*) as order_count,
SUM(total) as total_spent,
AVG(total) as avg_order_value
FROM orders
GROUP BY user_id
),
user_metrics AS (
SELECT
user_id,
order_count,
total_spent,
avg_order_value,
DATEDIFF(day, first_order, last_order) as customer_days,
CASE
WHEN DATEDIFF(day, first_order, last_order) = 0 THEN 0
ELSE order_count::float / (DATEDIFF(day, first_order, last_order) / 30.0)
END as orders_per_month
FROM user_orders
)
SELECT
user_id,
order_count,
total_spent,
avg_order_value,
orders_per_month,
-- CLV 예측 (단순화)
avg_order_value * orders_per_month * 12 * 3 as predicted_clv_3yr
FROM user_metrics
세션화 (Sessionization)
이벤트를 세션으로 그룹화:
WITH events_with_gaps AS (
SELECT
user_id,
event_time,
CASE
WHEN DATEDIFF(minute,
LAG(event_time) OVER (PARTITION BY user_id ORDER BY event_time),
event_time
) > 30 THEN 1
ELSE 0
END as is_new_session
FROM events
),
sessions AS (
SELECT
user_id,
event_time,
SUM(is_new_session) OVER (
PARTITION BY user_id
ORDER BY event_time
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
) as session_id
FROM events_with_gaps
)
SELECT
user_id,
session_id,
MIN(event_time) as session_start,
MAX(event_time) as session_end,
COUNT(*) as event_count,
DATEDIFF(second, MIN(event_time), MAX(event_time)) as session_duration_sec
FROM sessions
GROUP BY user_id, session_id
3.6 변환 패턴
Slowly Changing Dimensions (SCD)
시간에 따라 변하는 차원 데이터 처리:
Type 1: 덮어쓰기
-- 이전 값은 유지 안 함
UPDATE dim_users SET city = 'Seoul' WHERE user_id = 123;
Type 2: 히스토리 유지
-- 이전 레코드 종료
UPDATE dim_users
SET valid_to = CURRENT_DATE
WHERE user_id = 123 AND valid_to IS NULL;
-- 새 레코드 삽입
INSERT INTO dim_users (user_id, city, valid_from, valid_to)
VALUES (123, 'Seoul', CURRENT_DATE, NULL);
Type 3: 이전 값만 유지
UPDATE dim_users
SET
city = 'Seoul',
previous_city = city
WHERE user_id = 123;
Star Schema 변환
-- Fact Table
CREATE TABLE fct_sales AS
SELECT
o.order_id,
o.user_id,
o.product_id,
o.store_id,
o.date_id,
o.quantity,
o.unit_price,
o.total
FROM orders o;
-- Dimension Tables
CREATE TABLE dim_users AS
SELECT DISTINCT
user_id,
name,
email,
segment
FROM users;
CREATE TABLE dim_products AS
SELECT DISTINCT
product_id,
product_name,
category,
brand
FROM products;
3.7 변환 최적화
1. 파티셔닝 활용
-- 날짜로 파티션
CREATE TABLE orders_partitioned (
order_id INT,
user_id INT,
total DECIMAL,
order_date DATE
)
PARTITION BY RANGE (order_date) (
PARTITION p_2025_01 VALUES LESS THAN ('2025-02-01'),
PARTITION p_2025_02 VALUES LESS THAN ('2025-03-01'),
...
);
2. 증분 처리
-- 새로운 데이터만 처리
INSERT INTO orders_mart
SELECT * FROM transform_orders(
SELECT * FROM orders_raw
WHERE updated_at > (SELECT MAX(updated_at) FROM orders_mart)
);
3. 병렬 처리
from multiprocessing import Pool
def transform_partition(partition_id):
df = read_partition(partition_id)
df_transformed = transform(df)
write_partition(df_transformed, partition_id)
# 8개 프로세스로 병렬 처리
with Pool(8) as p:
p.map(transform_partition, range(100))
3.8 변환 테스트
dbt 테스트
-- models/schema.yml
version: 2
models:
- name: fct_daily_sales
description: "Daily sales metrics"
columns:
- name: date
tests:
- unique
- not_null
- name: revenue
tests:
- not_null
- dbt_utils.accepted_range:
min_value: 0
inclusive: true
- name: orders
tests:
- not_null
- dbt_utils.expression_is_true:
expression: ">= customers"
Python 테스트
import pytest
def test_transform_removes_nulls():
input_df = pd.DataFrame({
'user_id': [1, None, 3],
'value': [10, 20, 30]
})
result = transform(input_df)
assert result['user_id'].isna().sum() == 0
def test_transform_normalizes_email():
input_df = pd.DataFrame({
'email': ['TEST@EXAMPLE.COM', 'user@DOMAIN.com']
})
result = transform(input_df)
assert (result['email'] == result['email'].str.lower()).all()
변환의 황금률
- 멱등성: 여러 번 실행해도 같은 결과
- 테스트: 모든 변환 로직을 테스트
- 문서화: 왜 이렇게 변환하는지 설명
- 모니터링: 데이터 품질 지표 추적
3.9 다음 장 미리보기
다음 장에서는 변환된 데이터를 목적지에 적재하는 방법을 배웁니다:
- 다양한 저장소 옵션 (Warehouse, Lake, Mart)
- Batch vs Stream 적재
- Upsert와 Merge 전략
- 파티셔닝과 인덱싱
- 성능 최적화