4.1 데이터 계보란?
Data Lineage(데이터 계보)는 데이터의 출처부터 최종 목적지까지의 전체 여정을 추적하고 시각화하는 것입니다. 데이터가 어디서 왔고, 어떻게 변환되었으며, 어디로 흘러가는지를 명확하게 보여줍니다.
- 신뢰성: "이 대시보드의 숫자는 정확한가?" 질문에 답할 수 있음
- 영향도 분석: 테이블 변경 시 영향받는 모든 다운스트림 파악
- 규제 준수: GDPR, CCPA 등 데이터 흐름 추적 요구사항 충족
- 디버깅: 데이터 문제 발생 시 원인 추적
- 최적화: 불필요한 ETL 파이프라인 식별
4.2 계보의 3가지 레벨
1) 데이터셋 레벨 계보 (Coarse-grained)
테이블/파일 단위의 의존성을 추적합니다.
raw_customers → stg_customers → dim_customers → customer_report
2) 컬럼 레벨 계보 (Fine-grained)
개별 컬럼의 변환과 흐름을 추적합니다.
orders.total_amount → customer_metrics.total_revenue
orders.discount → customer_metrics.total_revenue (SUM aggregation)
3) 필드 레벨 계보 (Ultra-fine-grained)
개별 값의 변환 로직까지 추적합니다.
customer.email → LOWER(customer.email) → HASH(email) → anonymized_email
4.3 계보 그래프 구조
데이터 계보는 방향성 비순환 그래프(DAG, Directed Acyclic Graph)로 표현됩니다.
그래프 구성 요소
- 노드 (Node): 데이터셋 (테이블, 뷰, 파일)
- 엣지 (Edge): 데이터 흐름 (변환, 복사, 이동)
- 프로세스 (Process): ETL 작업, SQL 쿼리, 스크립트
{
"nodes": [
{"id": "src_orders", "type": "TABLE", "platform": "PostgreSQL"},
{"id": "stg_orders", "type": "TABLE", "platform": "Snowflake"},
{"id": "fct_orders", "type": "TABLE", "platform": "Snowflake"}
],
"edges": [
{"from": "src_orders", "to": "stg_orders", "process": "etl_ingest"},
{"from": "stg_orders", "to": "fct_orders", "process": "dbt_transform"}
],
"processes": [
{"id": "etl_ingest", "type": "SPARK_JOB", "schedule": "0 2 * * *"},
{"id": "dbt_transform", "type": "DBT_MODEL", "schedule": "0 3 * * *"}
]
}
4.4 계보 자동 추적 방법
1) SQL 파싱
SQL 쿼리를 분석하여 테이블 간 의존성을 추출합니다.
-- SQL 쿼리
CREATE TABLE customer_summary AS
SELECT
c.customer_id,
c.email,
COUNT(o.order_id) as order_count,
SUM(o.total_amount) as total_spent
FROM customers c
LEFT JOIN orders o ON c.customer_id = o.customer_id
GROUP BY c.customer_id, c.email;
-- 추출된 계보
Upstream: customers, orders
Downstream: customer_summary
Columns: customers.customer_id, customers.email,
orders.order_id, orders.total_amount
2) 실행 로그 분석
Spark, Airflow 등의 실행 로그에서 데이터 흐름을 추출합니다.
3) API 후킹
데이터베이스 드라이버, Spark API에 후크를 설치하여 실시간 추적합니다.
4) 메타스토어 통합
Hive Metastore, AWS Glue 등에서 계보 정보를 가져옵니다.
4.5 주요 플랫폼별 계보 추적
| 플랫폼 | 계보 추적 방법 | 레벨 |
|---|---|---|
| Apache Spark | Listener API, Query Plans | Dataset, Column |
| dbt | manifest.json 파일 | Dataset, Column |
| Airflow | Lineage Backend | Dataset |
| Snowflake | ACCESS_HISTORY view | Dataset, Column |
| BigQuery | Audit Logs, INFORMATION_SCHEMA | Dataset |
dbt 계보 추출 예시
import json
# manifest.json 파일 읽기
with open('target/manifest.json') as f:
manifest = json.load(f)
# 모델 계보 추출
for node_id, node in manifest['nodes'].items():
if node['resource_type'] == 'model':
print(f"Model: {node['name']}")
print(f" Depends on: {node['depends_on']['nodes']}")
print(f" Columns: {list(node['columns'].keys())}")
4.6 영향도 분석 (Impact Analysis)
특정 데이터셋을 변경할 때, 영향을 받는 모든 다운스트림을 자동으로 파악합니다.
영향받는 다운스트림:
1. stg_customers.email (직접 의존)
2. dim_customers.email_domain (파생 컬럼)
3. customer_segments (세그먼트 계산에 사용)
4. marketing_campaigns (타겟팅에 사용)
5. customer_analytics_dashboard (Tableau 대시보드)
총 영향: 5개 데이터셋, 12개 리포트, 3개 대시보드
영향도 분석 알고리즘
def analyze_impact(dataset_id, change_type):
"""
주어진 데이터셋 변경의 영향도를 분석합니다.
"""
visited = set()
impact = []
def dfs_downstream(node_id, depth=0):
if node_id in visited:
return
visited.add(node_id)
# 현재 노드의 다운스트림 찾기
downstream = get_downstream_nodes(node_id)
for ds_node in downstream:
impact.append({
'node': ds_node,
'depth': depth + 1,
'impact_type': calculate_impact_type(change_type, ds_node)
})
dfs_downstream(ds_node['id'], depth + 1)
dfs_downstream(dataset_id)
return impact
4.7 계보 시각화
계보 그래프를 사용자가 이해하기 쉽게 시각화하는 것이 중요합니다.
시각화 라이브러리
- D3.js: 커스터마이즈 가능한 웹 기반 그래프
- Cytoscape.js: 대규모 그래프 시각화
- Mermaid: 텍스트 기반 다이어그램
- Graphviz: DOT 언어 기반 레이아웃
- 계층적 레이아웃 사용 (왼쪽 → 오른쪽 또는 위 → 아래)
- 플랫폼별로 색상 코딩
- 중요한 노드 강조 표시
- 확대/축소, 필터링 기능 제공
- 노드 클릭 시 상세 정보 팝업
4.8 컬럼 레벨 계보
테이블 레벨보다 더 세밀한 컬럼 단위 계보를 추적합니다.
-- 원본 SQL
SELECT
LOWER(TRIM(email)) as clean_email,
CASE
WHEN lifetime_value > 10000 THEN 'VIP'
WHEN lifetime_value > 1000 THEN 'Regular'
ELSE 'New'
END as customer_segment
FROM customers;
-- 추출된 컬럼 계보
clean_email:
- source: customers.email
- transformations: [TRIM, LOWER]
customer_segment:
- source: customers.lifetime_value
- transformations: [CASE WHEN segmentation]
- logic: "VIP: >10000, Regular: >1000, New: else"
4.9 계보 저장소
그래프 데이터베이스
계보는 본질적으로 그래프 구조이므로 그래프 데이터베이스가 적합합니다.
Neo4j 예시
// 노드 생성
CREATE (c:Dataset {
name: 'customers',
platform: 'PostgreSQL',
type: 'TABLE'
})
CREATE (s:Dataset {
name: 'customer_summary',
platform: 'Snowflake',
type: 'TABLE'
})
// 관계 생성
CREATE (c)-[:FEEDS_INTO {
process: 'etl_pipeline',
transformation: 'aggregation'
}]->(s)
// 다운스트림 조회
MATCH (d:Dataset {name: 'customers'})-[:FEEDS_INTO*]->(downstream)
RETURN downstream.name, downstream.platform;
4.10 실전 활용 시나리오
대시보드에서 이상한 수치가 발견되었습니다. 계보를 역추적하여 문제의 근원을 찾습니다.
대시보드 → BI 레이어 → DW 테이블 → ETL → 소스 DB
역추적 결과: 소스 DB에서 NULL 값이 증가함 → ETL이 잘못 처리
고객이 개인 정보 삭제를 요청했습니다. 계보를 통해 모든 복사본 위치를 파악합니다.
고객 데이터가 저장된 위치:
1. prod_db.customers (원본)
2. analytics_db.dim_customers (복사본)
3. marketing_db.customer_segments (파생)
4. logs.user_activities (로그)
5. s3://backup/customers/ (백업)
→ 모든 위치에서 해당 고객 데이터 삭제 필요
4.11 다음 장 예고
다음 장에서는 Business Glossary를 다룹니다. 조직의 공통 언어를 정의하고 관리하는 방법을 배우게 됩니다.