LLM WikiAccess-protected knowledge portal
← 스터디 홈
5편 · 약 26분

데이터 파이프라인 테스트 전략: 단위·통합·E2E 테스트

데이터 파이프라인 테스트가 어려운 이유

소프트웨어 개발에서 테스트 피라미드는 오래된 원칙이다. 하지만 데이터 파이프라인에 그대로 적용하면 잘 맞지 않는 부분이 생긴다. 파이프라인의 주요 로직은 SQL과 외부 시스템 호출로 이루어져 있어서, 전통적인 단위 테스트 방식으로는 커버하기 어렵다. DAG 하나가 Airflow → Spark → S3 → dbt → DW를 통과하는데, 어디서 어떻게 테스트를 끊을지 기준이 없으면 테스트는 "있어 보이지만 아무것도 잡지 못하는" 상태가 된다.

실제로 데이터 파이프라인에서 발생하는 장애 유형은 이렇다.

  • 로직 버그: SQL 조인 방향이 틀리거나, 집계 분모가 잘못됨
  • 스키마 드리프트: 소스 테이블 컬럼이 바뀌었는데 파이프라인이 모름
  • 데이터 품질 위반: NULL이 들어오면 안 되는 곳에 NULL이 생김
  • 통합 실패: Airflow가 Spark 잡을 트리거했는데 파라미터가 맞지 않음
  • 환경 불일치: dev에서는 됐는데 prod에서 실패

이 장에서는 이 다섯 가지 유형을 커버하는 테스트 계층을 설계하는 방법을 다룬다.


데이터 파이프라인 테스팅 피라미드

E2E / 파이프라인 테스트 전체 DAG 실행, 실제 저장소 사용 느림 · 비쌈 · 적게 통합 테스트 컴포넌트 연결: Airflow Task → DB, Spark → S3 중간 속도 · 환경 필요 단위 테스트 변환 로직, SQL 단위, DAG 구조 검증, 스키마 드리프트 dbt unit test, pytest, DAG validation 빠름 · 격리 · 많이 많음 ▲ 적음 ▼ 테스트 수
데이터 파이프라인 테스팅 피라미드

단위 테스트를 최대한 넓게 깔고, 통합 테스트는 컴포넌트 경계에서 필요한 것만, 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=buildable

mock 입력을 사용하므로 실제 웨어하우스 없이도 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_at

dbt-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: 1000000

Airflow 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_list

Python 변환 로직 단위 테스트

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"] > 0

Spark 변환 통합 테스트

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() == 1

3단계: 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"] > 0

CI/CD 파이프라인에 테스트 계층 통합

PR 오픈
DAG 파싱
검증
pytest
단위 테스트
dbt build
(dev 환경)
~2-5분
merge → staging
통합 테스트
(staging 환경)
dbt test
(staging 전체)
smoke E2E
(핵심 경로)
~15-30분
prod 배포
배포 후
dbt source freshness
핵심 mart
row count 체크
SLO 알림
연동
~5분
데이터 파이프라인 CI 테스트 흐름
# .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=short

Great 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/