데이터 파이프라인 테스트가 어려운 이유
소프트웨어 개발에서 테스트 피라미드는 오래된 원칙이다. 하지만 데이터 파이프라인에 그대로 적용하면 잘 맞지 않는 부분이 생긴다. 파이프라인의 주요 로직은 SQL과 외부 시스템 호출로 이루어져 있어서, 전통적인 단위 테스트 방식으로는 커버하기 어렵다. DAG 하나가 Airflow → Spark → S3 → dbt → DW를 통과하는데, 어디서 어떻게 테스트를 끊을지 기준이 없으면 테스트는 "있어 보이지만 아무것도 잡지 못하는" 상태가 된다.
실제로 데이터 파이프라인에서 발생하는 장애 유형은 이렇다.
- 로직 버그: SQL 조인 방향이 틀리거나, 집계 분모가 잘못됨
- 스키마 드리프트: 소스 테이블 컬럼이 바뀌었는데 파이프라인이 모름
- 데이터 품질 위반: NULL이 들어오면 안 되는 곳에 NULL이 생김
- 통합 실패: Airflow가 Spark 잡을 트리거했는데 파라미터가 맞지 않음
- 환경 불일치: dev에서는 됐는데 prod에서 실패
이 장에서는 이 다섯 가지 유형을 커버하는 테스트 계층을 설계하는 방법을 다룬다.
데이터 파이프라인 테스팅 피라미드
단위 테스트를 최대한 넓게 깔고, 통합 테스트는 컴포넌트 경계에서 필요한 것만, E2E는 가장 중요한 경로만 커버한다.
1단계: 단위 테스트
dbt 단위 테스트 (dbt 1.8+)
dbt 1.8부터 공식 단위 테스트 기능이 추가됐다. SQL 변환 로직을 실제 데이터 없이 mock 입력으로 검증한다.
# models/orders/schema.yml
unit_tests:
- name: test_order_total_with_discount
model: orders_mart
given:
- input: ref('stg_orders')
rows:
- {order_id: 1, subtotal: 10000, discount: 1000}
- {order_id: 2, subtotal: 5000, discount: 0}
expect:
rows:
- {order_id: 1, total_amount: 9000}
- {order_id: 2, total_amount: 5000}# 단위 테스트만 실행 (DW 연결 없이도 가능)
dbt test --select orders_mart --indirect-selection=buildablemock 입력을 사용하므로 실제 웨어하우스 없이도 CI에서 빠르게 실행된다.
dbt 데이터 테스트
빌트인 테스트로 데이터 품질 기대치를 선언한다.
# models/schema.yml
models:
- name: orders_mart
columns:
- name: order_id
tests:
- unique
- not_null
- name: status
tests:
- accepted_values:
values: ['pending', 'completed', 'cancelled']
- name: customer_id
tests:
- relationships:
to: ref('customers')
field: id
sources:
- name: raw
tables:
- name: orders
freshness:
warn_after: {count: 1, period: hour}
error_after: {count: 6, period: hour}
loaded_at_field: _loaded_atdbt-expectations 패키지를 추가하면 Great Expectations 스타일의 더 풍부한 assertion이 가능하다.
# dbt-expectations 활용
- name: amount
tests:
- dbt_expectations.expect_column_values_to_be_between:
min_value: 0
max_value: 10000000
- dbt_expectations.expect_column_quantile_values_to_be_between:
quantile: 0.99
min_value: 0
max_value: 1000000Airflow DAG 구조 단위 테스트
DAG 파일이 파싱 가능한지, 태스크 의존성이 올바른지, 설정 값이 기준을 충족하는지 검사한다.
# tests/test_dag_validation.py
import pytest
from airflow.models import DagBag
@pytest.fixture(scope="module")
def dagbag():
return DagBag(dag_folder="dags/", include_examples=False)
def test_no_import_errors(dagbag):
assert len(dagbag.import_errors) == 0, \
f"DAG import errors: {dagbag.import_errors}"
def test_all_dags_have_tags(dagbag):
for dag_id, dag in dagbag.dags.items():
assert dag.tags, f"DAG '{dag_id}' has no tags"
def test_all_dags_have_catchup_false(dagbag):
for dag_id, dag in dagbag.dags.items():
assert dag.catchup is False, \
f"DAG '{dag_id}' has catchup=True — backfill 폭발 위험"
def test_task_dependencies(dagbag):
dag = dagbag.get_dag("orders_daily_pipeline")
assert dag is not None
# extract → transform → load 순서 확인
extract_task = dag.get_task("extract_orders")
transform_task = dag.get_task("transform_orders")
assert transform_task in extract_task.downstream_listPython 변환 로직 단위 테스트
Airflow PythonOperator나 Spark UDF에 담긴 변환 로직은 pytest로 단위 테스트한다.
# src/transformations/orders.py
def compute_order_total(subtotal: float, discount: float, tax_rate: float) -> float:
if subtotal < 0:
raise ValueError("subtotal must be non-negative")
discounted = subtotal - discount
return round(discounted * (1 + tax_rate), 2)
# tests/test_orders.py
import pytest
from src.transformations.orders import compute_order_total
def test_basic_total():
assert compute_order_total(10000, 1000, 0.1) == 9900.0
def test_zero_discount():
assert compute_order_total(5000, 0, 0.1) == 5500.0
def test_negative_subtotal_raises():
with pytest.raises(ValueError):
compute_order_total(-100, 0, 0.1)2단계: 통합 테스트
통합 테스트는 컴포넌트 간 연결을 검증한다. 실제 DB나 오브젝트 스토리지가 필요하지만, 프로덕션 대신 테스트 환경을 쓴다.
Airflow Task 통합 테스트
dag.test() 메서드는 Executor 없이 로컬에서 태스크를 실행한다. 외부 연결이 필요하지만 Airflow 메타데이터 DB는 불필요하다.
# tests/integration/test_dag_run.py
import pytest
from datetime import datetime
from airflow.models import DagBag
def test_extract_task_runs(tmp_path):
"""extract_orders 태스크가 S3에서 데이터를 가져오는지 검증"""
dagbag = DagBag(dag_folder="dags/")
dag = dagbag.get_dag("orders_daily_pipeline")
# dag.test()는 모든 태스크를 순서대로 실행
dag.test(
execution_date=datetime(2026, 7, 1),
run_conf={"env": "test", "s3_bucket": "my-test-bucket"}
)
# 테스트 S3 버킷에서 결과 파일 존재 여부 확인
import boto3
s3 = boto3.client("s3")
resp = s3.list_objects_v2(
Bucket="my-test-bucket",
Prefix="orders/2026/07/01/"
)
assert resp["KeyCount"] > 0Spark 변환 통합 테스트
PySpark 잡은 로컬 SparkSession을 사용해 통합 테스트한다. 실제 클러스터 없이도 HDFS API를 흉내 낼 수 있다.
# tests/integration/test_spark_transform.py
import pytest
from pyspark.sql import SparkSession
from src.spark_jobs.transform_orders import transform_orders
@pytest.fixture(scope="module")
def spark():
return SparkSession.builder \
.master("local[2]") \
.appName("test") \
.config("spark.sql.shuffle.partitions", "2") \
.getOrCreate()
def test_transform_orders_deduplication(spark):
raw_data = [
(1, "2026-07-01", 10000),
(1, "2026-07-01", 10000), # 중복 행
(2, "2026-07-01", 5000),
]
df = spark.createDataFrame(raw_data, ["order_id", "date", "amount"])
result = transform_orders(df)
assert result.count() == 2 # 중복 제거 확인
assert result.filter("order_id = 1").count() == 13단계: E2E 파이프라인 테스트
E2E 테스트는 전체 파이프라인을 샘플 데이터로 실행하고 최종 출력을 검증한다. 비용이 크므로 주요 경로(Happy Path)와 핵심 장애 시나리오만 커버한다.
# tests/e2e/test_orders_pipeline_e2e.py
import pytest
from datetime import datetime, timezone
# CI 환경에서만 실행 (로컬은 건너뜀)
@pytest.mark.skipif(
not os.getenv("CI_E2E_ENABLED"),
reason="E2E tests require CI environment"
)
def test_orders_pipeline_end_to_end():
"""
실제 파이프라인 실행:
1. 테스트 소스 DB에 샘플 주문 데이터 삽입
2. DAG trigger
3. DW에서 mart 테이블 결과 검증
"""
# 1. 소스 데이터 준비
source_db = get_test_db_connection()
source_db.execute("""
INSERT INTO orders (id, customer_id, amount, status, created_at)
VALUES (9999, 100, 50000, 'completed', NOW())
""")
# 2. DAG 트리거
airflow_client = AirflowAPIClient(base_url=os.getenv("AIRFLOW_URL"))
run_id = airflow_client.trigger_dag(
dag_id="orders_daily_pipeline",
conf={"execution_date": "2026-07-08", "env": "e2e-test"}
)
# 3. 완료 대기 (최대 10분)
wait_for_dag_completion(airflow_client, run_id, timeout_seconds=600)
# 4. DW 결과 검증
dw = get_dw_connection()
result = dw.execute("""
SELECT total_amount, status
FROM orders_mart
WHERE order_id = 9999
""").fetchone()
assert result is not None, "주문 9999가 mart에 존재해야 함"
assert result["status"] == "completed"
assert result["total_amount"] > 0CI/CD 파이프라인에 테스트 계층 통합
검증
단위 테스트
(dev 환경)
(staging 환경)
(staging 전체)
(핵심 경로)
dbt source freshness
row count 체크
연동
# .github/workflows/pipeline-test.yml
name: Pipeline CI
on:
pull_request:
branches: [main, staging]
paths:
- 'dags/**'
- 'models/**'
- 'src/**'
- 'tests/**'
jobs:
unit-test:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: Python 설치
uses: actions/setup-python@v5
with:
python-version: "3.11"
- name: 의존성 설치
run: pip install -r requirements-dev.txt
- name: DAG 파싱 검증
run: |
python -m pytest tests/test_dag_validation.py -v
- name: 단위 테스트
run: |
python -m pytest tests/unit/ -v --tb=short --cov=src --cov-report=xml
- name: dbt 컴파일 & 단위 테스트
run: |
dbt deps
dbt compile --profiles-dir ./ci-profiles
dbt test --select tag:unit --profiles-dir ./ci-profiles
integration-test:
runs-on: ubuntu-latest
needs: unit-test
if: github.base_ref == 'main'
services:
postgres:
image: postgres:15
env:
POSTGRES_PASSWORD: testpass
POSTGRES_DB: testdb
options: >-
--health-cmd pg_isready
--health-interval 10s
--health-retries 5
steps:
- uses: actions/checkout@v4
- name: 통합 테스트
env:
DB_URL: postgresql://postgres:testpass@localhost/testdb
run: |
python -m pytest tests/integration/ -v --tb=shortGreat Expectations: 선언적 데이터 품질 게이트
dbt 테스트만으로 부족할 때 Great Expectations(GX)를 추가하면 더 풍부한 통계 기반 검증이 가능하다.
# great_expectations/checkpoints/orders_mart_checkpoint.yaml
# 체크포인트 실행
import great_expectations as gx
context = gx.get_context()
results = context.run_checkpoint(checkpoint_name="orders_mart_checkpoint")
if not results["success"]:
raise RuntimeError(f"Data quality gate failed: {results}")# Expectation Suite 예시 (코드 기반 정의)
suite = context.add_expectation_suite("orders_mart")
# 행 수
suite.add_expectation(
gx.expectations.ExpectTableRowCountToBeBetween(min_value=1000, max_value=10_000_000)
)
# NULL 금지
suite.add_expectation(
gx.expectations.ExpectColumnValuesToNotBeNull(column="order_id")
)
# 분포 기반 체크 (99th percentile)
suite.add_expectation(
gx.expectations.ExpectColumnQuantileValuesToBeBetween(
column="amount",
quantile_ranges={
"quantiles": [0.25, 0.75, 0.99],
"value_ranges": [[0, 50000], [0, 200000], [0, 5000000]]
}
)
)GX는 dbt 테스트보다 복잡한 분포 검사, 교차 테이블 검증, 외부 데이터소스 비교가 필요할 때 쓴다.
테스트 전략 결정 기준
| 상황 | 선택 |
|---|---|
| SQL 변환 로직의 정확성 검증 | dbt unit test |
| 데이터 품질 (NULL, 범위, 유니크) | dbt built-in test + dbt-expectations |
| DAG 구조가 의도대로인지 | pytest + DagBag |
| Airflow Task의 외부 API 호출 | pytest + mock (unittest.mock) |
| Spark 변환 로직 | pytest + 로컬 SparkSession |
| 컴포넌트 간 연결 (Airflow → DB) | dag.test() + 테스트 환경 |
| 전체 파이프라인 행복 경로 | E2E (CI에서 선택적 실행) |
| 통계 분포 / 교차 검증 | Great Expectations |
흔한 실수와 대응
실수 1: 모든 테스트를 E2E로 작성 느리고 깨지기 쉬운 테스트만 남는다. 변환 로직은 단위 테스트로 격리한다.
실수 2: Mock을 너무 많이 써서 실제 연결을 전혀 테스트 안 함 통합 테스트 계층이 없으면 "로컬에선 됐는데 prod에서 실패"가 반복된다.
실수 3: dbt 테스트를 prod 데이터에서만 실행 prod에서 처음 실패를 발견하면 이미 다운스트림 대시보드가 망가져 있다. CI에서 dev/staging 데이터로 먼저 실행해야 한다.
실수 4: 테스트가 외부 상태에 의존 테스트 전후 데이터 정리(teardown)를 빠뜨리면 테스트 간 오염이 생긴다. pytest fixture의 yield + 정리 코드 패턴을 쓴다.
@pytest.fixture
def clean_test_data(db):
# setup
db.execute("INSERT INTO test_orders VALUES (9999, ...)")
yield
# teardown — 테스트 후 항상 정리
db.execute("DELETE FROM test_orders WHERE id = 9999")References
- Astronomer Docs, "Test Airflow DAGs": https://www.astronomer.io/docs/learn/testing-airflow
- Astronomer Blog, "Testing Your Apache Airflow DAGs": https://www.astronomer.io/blog/testing-your-apache-airflow-dags/
- dbt Labs, "dbt Unit Tests": https://docs.getdbt.com/docs/build/unit-tests
- Datacoves, "dbt Testing: A Complete Guide to Data Tests, Unit Tests, and Testing Packages": https://datacoves.com/post/dbt-test-options
- Datafold Blog, "Testing data pipelines: The Modern Data Stack challenge": https://www.datafold.com/blog/testing-data-pipelines/
- Atlan, "Testing Data Pipelines: A Complete Guide for 2026": https://atlan.com/testing-data-pipelines/
- Airbyte Blog, "Testing Data Pipelines with dbt-expectations: A Beginner's Guide": https://airbyte.com/blog/testing-with-dbt-expectations
- Medium (Chandu Kavar), "Testing in Airflow Part 2 — Integration Tests and End-To-End Pipeline Tests": https://medium.com/@chandukavar/testing-in-airflow-part-2-integration-tests-and-end-to-end-pipeline-tests-af0555cd1a82
- Great Expectations Documentation: https://docs.greatexpectations.io/