7.1 데이터 품질이 중요한 이유
잘못된 데이터는 잘못된 결정으로 이어집니다. 데이터 품질은 선택이 아니라 필수입니다.
데이터 품질의 6가지 차원
- 완전성: 모든 필수 데이터가 있는가?
- 정확성: 데이터가 실제를 반영하는가?
- 일관성: 시스템 간 일치하는가?
- 적시성: 데이터가 최신인가?
- 유효성: 비즈니스 규칙을 만족하는가?
- 유일성: 중복이 없는가?
7.2 데이터 검증
스키마 검증
from pydantic import BaseModel, Field, validator
class OrderRecord(BaseModel):
order_id: int = Field(..., gt=0)
user_id: int = Field(..., gt=0)
total: float = Field(..., ge=0)
status: str = Field(..., regex='^(pending|completed|cancelled)$')
created_at: datetime
@validator('created_at')
def not_future(cls, v):
if v > datetime.now():
raise ValueError('created_at cannot be in the future')
return v
NULL 검증
SELECT
COUNT(*) as total_records,
COUNT(user_id) as non_null_user_id,
COUNT(*) - COUNT(user_id) as null_user_id,
(COUNT(*) - COUNT(user_id))::float / COUNT(*) as null_rate
FROM orders;
-- null_rate > 0.01이면 알림
범위 검증
SELECT COUNT(*)
FROM orders
WHERE total < 0 OR total > 1000000;
-- 이상치 발견 시 Dead Letter Queue로
7.3 멱등성 (Idempotency)
같은 작업을 여러 번 실행해도 같은 결과를 보장:
멱등성 패턴 1: UPSERT
MERGE INTO dim_users target
USING staging.users source
ON target.user_id = source.user_id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT ...;
멱등성 패턴 2: 파티션 덮어쓰기
-- 특정 날짜 파티션만 덮어씀
INSERT OVERWRITE TABLE sales PARTITION (date='2025-12-26')
SELECT * FROM staging.sales
WHERE date = '2025-12-26';
멱등성 패턴 3: 트랜잭션 ID
-- 중복 방지
INSERT INTO processed_files (file_name, processed_at)
VALUES ('data_2025_12_26.csv', NOW())
ON CONFLICT (file_name) DO NOTHING;
7.4 에러 핸들링 전략
1. Fail Fast
첫 오류에서 즉시 중단:
def process_batch(records):
for record in records:
if not validate(record):
raise ValidationError(f"Invalid record: {record}")
transform(record)
2. Quarantine Bad Records
잘못된 레코드는 격리:
good_records = []
bad_records = []
for record in records:
try:
validated = validate(record)
good_records.append(validated)
except ValidationError as e:
bad_records.append({'record': record, 'error': str(e)})
# 좋은 레코드는 처리
load(good_records)
# 나쁜 레코드는 별도 저장
save_to_quarantine(bad_records)
3. Retry with Backoff
일시적 오류는 재시도:
@retry(
stop=stop_after_attempt(3),
wait=wait_exponential(multiplier=1, min=4, max=10)
)
def fetch_data_from_api():
response = requests.get(url)
response.raise_for_status()
return response.json()
7.5 정확히 한 번 처리 (Exactly-Once)
At-Most-Once
최대 한 번 처리 (손실 가능):
# 커밋 먼저, 처리 나중
consumer.commit()
process(message) # 여기서 실패하면 메시지 손실
At-Least-Once
최소 한 번 처리 (중복 가능):
# 처리 먼저, 커밋 나중
process(message)
consumer.commit() # 여기서 실패하면 재처리
Exactly-Once
정확히 한 번 처리:
BEGIN TRANSACTION;
-- 메시지 처리
INSERT INTO processed_events VALUES (...);
-- 오프셋 저장
INSERT INTO consumer_offsets VALUES (partition, offset);
COMMIT;
7.6 데이터 품질 프레임워크
Great Expectations
import great_expectations as ge
# Expectation Suite 생성
suite = ge.core.ExpectationSuite("orders_suite")
# 기대값 정의
suite.add_expectation(
ge.core.ExpectationConfiguration(
expectation_type="expect_column_values_to_be_between",
kwargs={"column": "total", "min_value": 0}
)
)
# 검증 실행
results = context.run_validation_operator(
"action_list_operator",
assets_to_validate=[batch],
run_id="validation_run"
)
if not results["success"]:
send_alert("Data quality check failed!")
Soda
# checks.yml
checks for orders:
- row_count > 0
- missing_count(user_id) = 0
- invalid_count(status) = 0:
valid values: ['pending', 'completed', 'cancelled']
- duplicate_count(order_id) = 0
- avg(total) between 10 and 1000
7.7 데이터 계약 (Data Contract)
생산자와 소비자 간의 명시적 계약:
contract:
name: orders_v1
owner: data-team@company.com
schema:
columns:
- name: order_id
type: INTEGER
constraints:
- NOT NULL
- UNIQUE
- name: total
type: DECIMAL(10,2)
constraints:
- >= 0
- name: status
type: STRING
allowed_values: ['pending', 'completed', 'cancelled']
sla:
freshness: 1 hour
completeness: 99.9%
availability: 99.95%
7.8 테스트 주도 데이터 (TDD for Data)
단위 테스트
def test_clean_email():
input_df = pd.DataFrame({'email': ['TEST@example.com']})
result = clean_data(input_df)
assert result['email'][0] == 'test@example.com'
def test_remove_duplicates():
input_df = pd.DataFrame({
'id': [1, 1, 2],
'value': [10, 10, 20]
})
result = remove_duplicates(input_df)
assert len(result) == 2
통합 테스트
def test_full_pipeline():
# 테스트 데이터 준비
setup_test_database()
# 파이프라인 실행
pipeline.run()
# 결과 검증
result = query_result_table()
assert len(result) == expected_count
assert result['total'].sum() == expected_total
7.9 참조 무결성
Foreign Key 검증
SELECT
o.order_id,
o.user_id
FROM orders o
LEFT JOIN users u ON o.user_id = u.user_id
WHERE u.user_id IS NULL;
-- 결과가 있으면 참조 무결성 위반
Python으로 참조 무결성 체크
def check_referential_integrity(orders_df, users_df):
"""주문의 모든 user_id가 users 테이블에 존재하는지 확인"""
orphaned_orders = orders_df[
~orders_df['user_id'].isin(users_df['user_id'])
]
if len(orphaned_orders) > 0:
raise ValueError(
f"Found {len(orphaned_orders)} orders with invalid user_id"
)
return True
7.10 데이터 리니지 (Data Lineage)
데이터의 출처와 변환 과정을 추적하는 것이 품질 보장에 중요합니다.
Apache Atlas로 리니지 추적
from pyatlas import Atlas
atlas = Atlas(
base_url="http://atlas-server:21000",
username="admin",
password="admin"
)
# 데이터셋 등록
atlas.create_entity({
"typeName": "DataSet",
"attributes": {
"name": "orders_raw",
"qualifiedName": "s3://bucket/raw/orders@cluster",
"description": "Raw orders from e-commerce system"
}
})
# 프로세스 등록 (변환)
atlas.create_entity({
"typeName": "Process",
"attributes": {
"name": "clean_orders",
"qualifiedName": "clean_orders_process@cluster",
"inputs": [{"typeName": "DataSet", "uniqueAttributes": {"qualifiedName": "s3://bucket/raw/orders@cluster"}}],
"outputs": [{"typeName": "DataSet", "uniqueAttributes": {"qualifiedName": "s3://bucket/clean/orders@cluster"}}]
}
})
dbt로 리니지 자동 생성
-- models/staging/stg_orders.sql
{{
config(
materialized='view',
tags=['staging']
)
}}
SELECT
order_id,
user_id,
total,
created_at
FROM {{ source('raw', 'orders') }}
WHERE deleted_at IS NULL
-- dbt는 자동으로 리니지 그래프 생성
7.11 데이터 검증 자동화
Great Expectations Checkpoint
import great_expectations as ge
context = ge.get_context()
# Checkpoint 설정
checkpoint_config = {
"name": "orders_validation",
"config_version": 1,
"class_name": "SimpleCheckpoint",
"validations": [
{
"batch_request": {
"datasource_name": "my_datasource",
"data_connector_name": "default_runtime_data_connector",
"data_asset_name": "orders"
},
"expectation_suite_name": "orders_suite"
}
],
"action_list": [
{
"name": "store_validation_result",
"action": {"class_name": "StoreValidationResultAction"}
},
{
"name": "send_slack_notification",
"action": {
"class_name": "SlackNotificationAction",
"slack_webhook": "https://hooks.slack.com/services/xxx",
"notify_on": "failure"
}
}
]
}
# Checkpoint 실행
checkpoint = context.add_checkpoint(**checkpoint_config)
result = checkpoint.run()
if not result["success"]:
raise ValueError("Data validation failed!")
7.12 Schema Evolution
시간이 지남에 따라 스키마가 변경될 때 품질을 유지하는 방법입니다.
스키마 버전 관리
# schema_v1.py
from pydantic import BaseModel
class OrderV1(BaseModel):
order_id: int
user_id: int
total: float
# schema_v2.py
class OrderV2(BaseModel):
order_id: int
user_id: int
total: float
currency: str = "USD" # 새 필드 (기본값 제공)
# 마이그레이션
def migrate_v1_to_v2(order_v1):
return OrderV2(
**order_v1.dict(),
currency="USD"
)
Avro 스키마 진화
{
"type": "record",
"name": "Order",
"namespace": "com.example",
"fields": [
{"name": "order_id", "type": "int"},
{"name": "user_id", "type": "int"},
{"name": "total", "type": "double"},
{
"name": "currency",
"type": "string",
"default": "USD"
}
]
}
7.13 백프레셔와 데이터 흐름 제어
Producer-Consumer 패턴
import queue
import threading
def producer(q, max_size=1000):
"""데이터 생성"""
while True:
data = fetch_data()
# 큐가 가득 차면 대기 (백프레셔)
q.put(data, block=True)
def consumer(q):
"""데이터 처리"""
while True:
data = q.get(block=True)
process_data(data)
q.task_done()
# 큐 크기로 백프레셔 제어
q = queue.Queue(maxsize=100)
producer_thread = threading.Thread(target=producer, args=(q,))
consumer_thread = threading.Thread(target=consumer, args=(q,))
producer_thread.start()
consumer_thread.start()
Kafka 백프레셔
from kafka import KafkaConsumer
consumer = KafkaConsumer(
'orders-topic',
bootstrap_servers=['localhost:9092'],
max_poll_records=100, # 한 번에 최대 100개만 가져오기
fetch_max_wait_ms=500, # 최대 500ms 대기
)
for message in consumer:
# 처리 속도가 느리면 자동으로 백프레셔 발생
slow_processing(message.value)
7.14 트랜잭션과 원자성
Database Transaction
import psycopg2
def transfer_with_transaction(from_account, to_account, amount):
"""트랜잭션을 사용한 계좌 이체"""
conn = psycopg2.connect(DATABASE_URL)
try:
cursor = conn.cursor()
# 출금
cursor.execute(
"UPDATE accounts SET balance = balance - %s WHERE account_id = %s",
(amount, from_account)
)
# 입금
cursor.execute(
"UPDATE accounts SET balance = balance + %s WHERE account_id = %s",
(amount, to_account)
)
# 모두 성공하면 커밋
conn.commit()
except Exception as e:
# 하나라도 실패하면 롤백
conn.rollback()
raise e
finally:
conn.close()
Two-Phase Commit
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker
engine1 = create_engine('postgresql://db1/mydb')
engine2 = create_engine('postgresql://db2/mydb')
Session1 = sessionmaker(bind=engine1, twophase=True)
Session2 = sessionmaker(bind=engine2, twophase=True)
session1 = Session1()
session2 = Session2()
try:
# Phase 1: Prepare
session1.add(Record1(...))
session2.add(Record2(...))
session1.prepare()
session2.prepare()
# Phase 2: Commit
session1.commit()
session2.commit()
except Exception as e:
session1.rollback()
session2.rollback()
raise e
7.15 재처리 전략
Safe Reprocessing
def safe_reprocess(start_date, end_date):
"""안전한 재처리"""
# 1. 백업 생성
create_backup_table("orders", "orders_backup")
try:
# 2. 기존 데이터 삭제
delete_range("orders", start_date, end_date)
# 3. 재처리
reprocess_range(start_date, end_date)
# 4. 검증
validate_data("orders", start_date, end_date)
# 5. 백업 삭제
drop_table("orders_backup")
except Exception as e:
# 복구
restore_from_backup("orders", "orders_backup")
raise e
점진적 재처리
def incremental_reprocess(start_date, end_date):
"""하루씩 재처리하여 영향 최소화"""
current_date = start_date
while current_date <= end_date:
try:
# 하루치 재처리
reprocess_day(current_date)
# 검증
validate_day(current_date)
print(f"Successfully reprocessed {current_date}")
except Exception as e:
print(f"Failed to reprocess {current_date}: {e}")
# 실패한 날짜는 건너뛰고 계속
log_failure(current_date, str(e))
current_date += timedelta(days=1)
7.16 데이터 품질 SLA
| 메트릭 | 목표 | 측정 방법 | 조치 |
|---|---|---|---|
| 완전성 | > 99% | NULL 비율 | 소스 데이터 점검 |
| 정확성 | > 99.9% | 검증 규칙 통과율 | 변환 로직 수정 |
| 적시성 | < 2시간 | 데이터 신선도 | 파이프라인 최적화 |
| 일관성 | 100% | 참조 무결성 | FK 제약 추가 |
| 유일성 | 100% | 중복 레코드 비율 | UPSERT 로직 수정 |
7.17 Data Quality Scorecard
def calculate_quality_score(df):
"""데이터 품질 점수 계산"""
scores = {}
# 완전성 (NULL 비율)
null_rate = df.isnull().sum().sum() / (len(df) * len(df.columns))
scores['completeness'] = 1 - null_rate
# 유일성 (중복 비율)
duplicate_rate = df.duplicated(subset=['id']).sum() / len(df)
scores['uniqueness'] = 1 - duplicate_rate
# 유효성 (범위 검증)
invalid_count = len(df[df['total'] < 0])
scores['validity'] = 1 - (invalid_count / len(df))
# 종합 점수
overall_score = sum(scores.values()) / len(scores)
return {
'scores': scores,
'overall': overall_score,
'grade': 'A' if overall_score > 0.95 else 'B' if overall_score > 0.90 else 'C'
}
# 사용
result = calculate_quality_score(orders_df)
print(f"Quality Score: {result['overall']:.2%} (Grade: {result['grade']})")
# Quality Score: 97.50% (Grade: A)
요약
이 장에서는 데이터 품질과 신뢰성을 보장하는 다양한 기법을 배웠습니다:
- 데이터 품질 차원: 완전성, 정확성, 일관성, 적시성, 유효성, 유일성
- 검증 기법: 스키마, NULL, 범위, 참조 무결성
- 멱등성: UPSERT, 파티션 덮어쓰기, 트랜잭션 ID
- 에러 핸들링: Fail Fast, Quarantine, Retry with Backoff
- Exactly-Once: 트랜잭션, 멱등성 보장
- 품질 프레임워크: Great Expectations, Soda
- 데이터 계약: 생산자-소비자 간 명시적 합의
- TDD for Data: 단위 테스트, 통합 테스트
- 리니지: 데이터 출처와 변환 추적
- 스키마 진화: 버전 관리, 하위 호환성
데이터 품질 체크리스트
- ✅ 스키마 검증 (타입, 필수 컬럼)
- ✅ NULL 체크 (완전성)
- ✅ 중복 체크 (유일성)
- ✅ 범위 검증 (유효성)
- ✅ 참조 무결성 (일관성)
- ✅ 비즈니스 룰 검증
- ✅ 데이터 신선도 확인 (적시성)
- ✅ 멱등성 보장 (재실행 가능)
- ✅ 트랜잭션 처리 (원자성)
- ✅ 에러 핸들링 (복원력)
- ✅ 리니지 추적 (투명성)
- ✅ 품질 점수 계산 (측정)
복습 문제
- 데이터 품질의 6가지 차원을 설명하고 각각의 측정 방법을 제시하세요.
- 멱등성이 중요한 이유와 구현 방법 3가지를 설명하세요.
- At-Most-Once, At-Least-Once, Exactly-Once의 차이점을 비교하세요.
- Dead Letter Queue 패턴의 장점과 구현 방법을 설명하세요.
- Great Expectations와 Soda의 차이점은 무엇인가요?
- 데이터 계약(Data Contract)이 필요한 이유를 설명하세요.
- 스키마 진화 시 하위 호환성을 유지하는 방법은?
- 데이터 품질 점수를 계산하는 공식을 제시하세요.
실습 과제
- Pydantic을 사용하여 스키마 검증을 구현하세요.
- 멱등한 UPSERT 파이프라인을 작성하세요.
- Great Expectations로 데이터 품질 검증 스위트를 만드세요.
- 데이터 품질 점수를 계산하고 Slack으로 알림을 보내는 스크립트를 작성하세요.