Snowflake MigrationClickHouse Workshops

Snowflake의 dbt

소스 Medallion 파이프라인의 구축 방식: 소스, 스테이징 뷰, 증분 MERGE 모델, 스냅샷, 테스트.

이 문서는 NYC Taxi Snowflake Migration Lab에서 dbt(data build tool)가 어떻게 사용되는지 — 무엇을 하는지, 각 구성 요소가 왜 존재하는지, 그리고 어떻게 이해해야 하는지 — 설명합니다.


dbt가 하는 일(그리고 하지 않는 일)

dbt는 이미 데이터베이스에 들어 있는 데이터를 변환합니다. 외부에서 데이터를 로드하거나, 파일을 옮기거나, 인프라를 관리하지는 않습니다. dbt의 역할은 원시 테이블을 받아서 여러분이 작성한 SQL을 실행해 깔끔하고 테스트된, 분석에 바로 쓸 수 있는 테이블로 만드는 것입니다.

SQL을 위한 빌드 시스템이라고 생각하면 됩니다. models/ 안의 각 .sql 파일은 Snowflake에서 테이블 또는 뷰가 되는 모델입니다. dbt는 CREATE OR REPLACE 보일러플레이트를 처리하고, 모델 간 의존성을 해석하며, 테스트를 실행합니다.


프로젝트 구조

dbt/nyc_taxi_dbt/
├── dbt_project.yml          # Project config: name, folder layout, materialization defaults
├── profiles.yml.example     # Connection config template (copy to ~/.dbt/profiles.yml)
├── packages.yml             # Third-party dbt packages
│
├── macros/
│   ├── generate_schema_name.sql   # Overrides dbt's default schema naming logic
│   └── generate_surrogate_key.sql # Wrapper for consistent surrogate key generation
│
└── models/
    ├── sources.yml          # Declares RAW.TRIPS_RAW as an external source
    │
    ├── staging/             # Layer 1: clean and rename raw columns
    │   ├── schema.yml       # Column-level tests for staging models
    │   ├── stg_trips.sql
    │   └── stg_taxi_zones.sql
    │
    ├── intermediate/        # Layer 2: joins and enrichment (no physical table)
    │   └── int_trips_enriched.sql
    │
    └── analytics/           # Layer 3: final tables consumed by dashboards
        ├── schema.yml
        ├── fact_trips.sql
        ├── agg_hourly_zone_trips.sql
        ├── dim_date.sql
        ├── dim_payment_type.sql
        ├── dim_taxi_zones.sql
        └── dim_vendor.sql

세 개의 계층 (Medallion 아키텍처)

계층 1 — 스테이징 (models/staging/)

목적: 도착한 그대로의 원시 데이터를 사용 가능한 형태로 만듭니다.

이 모델들은 STAGING 스키마에 뷰로 생성됩니다(스토리지 비용 없음 — 쿼리 시점에 실행됩니다). 각 스테이징 모델은 한 가지 일만 합니다:

모델소스하는 일
stg_tripsRAW.TRIPS_RAW컬럼명을 snake_case로 변경, duration_minutes 추가, VARIANT 타입 TRIP_METADATA 컬럼을 타입이 지정된 컬럼들로 평탄화
stg_taxi_zonesANALYTICS.DIM_TAXI_ZONES가벼운 정리, COALESCE 가드 추가, dbt 계보(lineage) 노드 제공

여기서 가장 중요한 작업은 TRIP_METADATA VARIANT 컬럼을 평탄화하는 것입니다. Snowflake의 콜론 경로(colon-path) 문법은 중첩된 JSON 필드를 추출합니다:

-- Snowflake: colon-path notation
TRIP_METADATA:driver.rating::FLOAT     AS driver_rating,
TRIP_METADATA:app.surge_multiplier::FLOAT AS surge_multiplier

이것이 마이그레이션 과제 중 하나입니다 — ClickHouse는 대신 JSONExtractFloat(TRIP_METADATA, 'driver', 'rating')를 사용합니다.

계층 2 — 중간 (models/intermediate/)

목적: 모든 조인을 한 곳에서 수행해 반복하지 않도록 합니다.

int_trips_enriched는 stg_trips를 모든 차원(존, 결제 유형, 벤더, 날짜)과 조인하여 트립당 하나의 넓고 완전히 비정규화된 행을 만듭니다. 이 모델은 ephemeral로 선언되어 있으며, 이는 dbt가 이 모델을 참조하는 모델의 SQL 안에 인라인으로 삽입한다는 뜻입니다 — Snowflake에는 물리적 테이블이나 뷰가 생성되지 않습니다.

-- dbt_project.yml
intermediate:
  +materialized: ephemeral   # compiled inline, no CREATE TABLE

중간 결과가 하나의 다운스트림 모델에만 필요하고 스토리지나 쿼리 컴파일 오버헤드를 지불하고 싶지 않을 때 ephemeral을 사용하세요.

계층 3 — 분석 (models/analytics/)

목적: 최종적인, 대시보드에 바로 쓸 수 있는 테이블.

이 모델들은 ANALYTICS 스키마에 생성됩니다. 두 종류가 있습니다:

정적 차원 테이블 — 작고, dbt run마다 전체 재적재됩니다:

모델행 수비고
dim_date~7,6702009–2029 날짜 스파인, 회계 분기와 미국 연방 공휴일 포함
dim_payment_type6시드 데이터 그대로 전달
dim_vendor3시드 데이터 그대로 전달
dim_taxi_zones265stg_taxi_zones를 통해 전달

증분 팩트/집계 테이블 — 크고, 실행마다 MERGE로 갱신됩니다:

모델행 수비고
fact_trips50M트립당 한 행, 완전히 비정규화
agg_hourly_zone_trips~9M존별 시간당 사전 집계 카운트

Materialization

materialization은 특정 모델에 대해 dbt가 Snowflake에 무엇을 생성할지 제어합니다.

Materialization생성되는 Snowflake 객체사용 시점
viewCREATE VIEW저렴함; 항상 최신 데이터 반영; 스테이징에 사용
tableCREATE TABLE AS SELECT실행마다 전체 재구축; 작은 차원 테이블에 사용
incremental기존 테이블에 MERGE INTO큰 테이블; 새 행만 처리
ephemeral(객체 없음 — CTE로 인라인)하나의 다운스트림 모델이 공유하는 중간 로직

두 개의 증분 모델은 서로 다른 증분 전략을 보여줍니다:

fact_trips — 마지막 실행 이후의 새 트립을 처리합니다:

{% if is_incremental() %}
  WHERE pickup_at > (SELECT MAX(pickup_at) FROM {{ this }})
{% endif %}

agg_hourly_zone_trips — 늦게 도착하는 데이터를 잡기 위해 롤링 2시간 윈도우를 재집계합니다:

{% if is_incremental() %}
  WHERE pickup_at >= DATEADD('hour', -2, CURRENT_TIMESTAMP())
{% endif %}

첫 실행(빈 테이블)에서는 is_incremental()이 false를 반환하고 전체 데이터셋이 처리됩니다. 이후 실행에서는 새 데이터만 처리됩니다. 스키마가 변경되어 처음부터 재구축해야 한다면 다음을 실행하세요:

dbt run --full-refresh

MERGE 전략 (핵심 마이그레이션 과제)

incremental_strategy = 'merge'일 때 dbt는 Snowflake MERGE INTO 문을 생성합니다:

MERGE INTO ANALYTICS.FACT_TRIPS AS target
USING (SELECT ...) AS source
ON target.trip_id = source.trip_id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT ...;

이것은 이 랩에서 문서화된 가장 중요한 마이그레이션 과제 중 하나입니다. ClickHouse에는 MERGE 문이 없습니다. ClickHouse에서의 등가 방식은 ReplacingMergeTree 테이블 엔진을 사용하고 쿼리에 FINAL을 추가하는 것, 또는 명시적인 삽입/삭제 의미가 필요하면 CollapsingMergeTree를 사용하는 것입니다.


스키마 이름 지정: generate_schema_name 매크로

dbt의 기본 동작은 profiles.yml의 타깃 스키마와 dbt_project.yml의 커스텀 스키마를 연결합니다:

target schema = STAGING  +  custom schema = ANALYTICS  →  STAGING_ANALYTICS  (wrong)

이 프로젝트는 macros/generate_schema_name.sql의 커스텀 매크로로 그 동작을 재정의합니다:

{% macro generate_schema_name(custom_schema_name, node) -%}
  {%- if custom_schema_name is none -%}
    {{ target.schema | upper }}      -- no custom schema → use target schema
  {%- else -%}
    {{ custom_schema_name | upper }}  -- custom schema → use it directly
  {%- endif -%}
{%- endmacro %}

결과: +schema: ANALYTICS를 가진 모델은 STAGING_ANALYTICS가 아니라 ANALYTICS에 생성됩니다.

이 매크로는 하나의 dbt 프로젝트에 여러 스키마가 있고 타깃 스키마 이름이 앞에 붙는 것을 원하지 않을 때 항상 필요합니다.


연결과 자격 증명 (profiles.yml)

dbt는 ~/.dbt/profiles.yml에 정의된 프로필을 사용해 Snowflake에 연결합니다(git에 절대 커밋하지 않습니다). dbt_project.yml의 프로필 이름은 일치해야 합니다:

# dbt_project.yml
profile: 'nyc_taxi'

# ~/.dbt/profiles.yml
nyc_taxi:
  target: dev
  outputs:
    dev:
      type: snowflake
      account: "{{ env_var('SNOWFLAKE_ORG') }}-{{ env_var('SNOWFLAKE_ACCOUNT') }}"
      role: DBT_ROLE
      database: NYC_TAXI_DB
      warehouse: TRANSFORM_WH
      schema: STAGING       # ← this is the "target schema" / default schema
      threads: 4

핵심 포인트:

  • schema: STAGING은 기본 스키마입니다. +schema: 재정의가 없는 모델은 여기에 생성됩니다.
  • role: DBT_ROLE은 dbt가 필요한 권한만 가진, Terraform이 생성한 최소 권한 역할입니다.
  • threads: 4는 dbt가 몇 개의 모델을 병렬로 빌드할지 제어합니다.
  • 자격 증명은 환경 변수에서 오며, 설정 실행 전에 .env에서 로드됩니다.

테스트

dbt 테스트에는 두 가지 형태가 있습니다:

스키마 테스트 (schema.yml에 선언)

- name: trip_id
  tests:
    - not_null
    - unique
- name: total_amount_usd
  tests:
    - dbt_expectations.expect_column_values_to_be_between:
        min_value: 0
        max_value: 1000

not_null과 unique는 내장 테스트입니다. dbt_expectations 테스트는 packages.yml에 선언된 calogica/dbt_expectations 패키지에서 옵니다.

커스텀 SQL 테스트 (tests/)

-- tests/assert_revenue_positive.sql
-- A passing test returns 0 rows
SELECT trip_id, total_amount_usd
FROM {{ ref('fact_trips') }}
WHERE total_amount_usd < 0

커스텀 테스트는 그냥 SQL 쿼리입니다. dbt가 실행하고 행이 하나라도 반환되면 실패로 처리합니다.

모든 테스트를 실행하려면:

dbt test

서드파티 패키지 (packages.yml)

packages:
  - package: dbt-labs/dbt_utils
    version: [">=1.0.0", "<2.0.0"]
  - package: calogica/dbt_expectations
    version: [">=0.10.0", "<1.0.0"]

처음 사용하기 전에 설치하세요:

dbt deps

dbt_utils는 dim_date.sql에서 사용하는 date_spine 생성기를 제공합니다. dbt_expectations는 내장 not_null/unique를 넘어서는 범위/분포 테스트를 제공합니다.


의존성 그래프

dbt는 {{ ref() }} 호출을 따라가며 자동으로 올바른 순서로 모델을 빌드합니다:

RAW.TRIPS_RAW (source — not managed by dbt)
    └── stg_trips (view)
            └── int_trips_enriched (ephemeral)
                    ├── fact_trips (incremental table)
                    └── agg_hourly_zone_trips (incremental table)

ANALYTICS.DIM_TAXI_ZONES (seeded by SQL script)
    └── stg_taxi_zones (view)
            ├── int_trips_enriched
            └── dim_taxi_zones (table)

dbt_utils.date_spine
    └── dim_date (table)

{{ ref('stg_trips') }}는 한 모델이 다른 모델에 대한 의존성을 선언하는 방법입니다. {{ source('raw', 'TRIPS_RAW') }}는 외부 테이블(sources.yml에 정의됨)에 대한 의존성을 선언합니다.


자주 쓰는 명령

명령하는 일
dbt depspackages.yml의 패키지 설치
dbt run모든 모델 빌드(가능한 경우 증분)
dbt run --full-refresh모든 증분 모델을 처음부터 재구축
dbt run -s fact_tripsfact_trips와 그 의존성만 빌드
dbt test모든 스키마 테스트와 커스텀 테스트 실행
dbt builddbt run + dbt test를 함께 실행
dbt compile실행하지 않고 SQL만 생성(디버깅에 유용)
dbt docs generate && dbt docs serve계보 그래프를 빌드하고 브라우저에서 탐색

이 프로젝트에서는 fact_trips가 비어 있으면(첫 실행 또는 tear-down 이후) setup.sh가 dbt run --full-refresh를 자동으로 실행합니다.


dbt가 전체 설정에서 차지하는 위치

terraform apply          → creates warehouses, database, schemas, roles
scripts/01_create_tables.sql → creates raw tables, seeds dimension data
scripts/02_seed_data.sql → loads 50M synthetic trip rows
dbt deps && dbt build    → transforms raw data into analytics-ready tables
scripts/03_create_streams_tasks.sql → creates CDC stream and scheduled task

dbt는 파이프라인의 중간에 위치합니다. 원시 테이블이 존재하고 데이터가 들어 있어야 실행할 수 있습니다. setup.sh 스크립트가 이 순서를 처리합니다.

이 페이지의 내용

KO