Snowflake MigrationClickHouse Workshops

Snowflake 上の dbt

ソースとなる Medallion パイプラインの構成: sources、staging ビュー、インクリメンタルな MERGE モデル、snapshot、テスト。

このドキュメントでは、NYC Taxi Snowflake Migration Lab で dbt (data build tool) がどのように使われているか — 何をするのか、各パーツがなぜ存在するのか、どう考えればよいのか — を説明します。


dbt がやること(とやらないこと)

dbt は、すでにデータベースに入っているデータを変換します。外部からデータをロードしたり、ファイルを移動したり、インフラを管理したりはしません。その役割は、生のテーブルを受け取り、あなたが書いた SQL を実行することで、クリーンでテスト済みの分析可能なテーブルに変えることです。

SQL 向けのビルドシステムだと考えてください。models/ にある各 .sql ファイルが 1 つのモデルであり、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

3 つのレイヤー (Medallion アーキテクチャ)

レイヤー 1 — Staging (models/staging/)

目的: 到着したそのままの生データを受け取り、使える形にする。

これらのモデルは STAGING スキーマにビューとして作られます (ストレージコストなし — クエリ実行時に評価されます)。各 staging モデルは 1 つの仕事だけを担います:

モデルソース何をするか
stg_tripsRAW.TRIPS_RAWカラム名を snake_case にリネームし、duration_minutes を追加し、VARIANT 型の TRIP_METADATA カラムを型付きカラムへフラット化
stg_taxi_zonesANALYTICS.DIM_TAXI_ZONES軽いクリーンアップ、COALESCE によるガードの追加、dbt のリネージノードの提供

ここで最も重要な処理は、TRIP_METADATA VARIANT カラムのフラット化です。Snowflake のコロンパス構文はネストされた JSON フィールドを抽出します:

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

これは移行課題の 1 つです — ClickHouse では代わりに JSONExtractFloat(TRIP_METADATA, 'driver', 'rating') を使います。

レイヤー 2 — Intermediate (models/intermediate/)

目的: すべての結合を 1 か所で行い、繰り返さずに済むようにする。

int_trips_enriched は stg_trips をすべてのディメンション (zone、payment type、vendor、date) と結合し、トリップごとに横に広い完全非正規化の行を生成します。これは ephemeral として宣言されており、dbt はその SQL を参照元のモデルにインライン展開します — Snowflake 上に物理テーブルもビューも作られません。

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

ephemeral は、中間結果が 1 つの下流モデルからのみ必要で、ストレージやクエリのコンパイルオーバーヘッドを払いたくない場合に使います。

レイヤー 3 — Analytics (models/analytics/)

目的: ダッシュボードにそのまま使える最終テーブル。

これらは ANALYTICS スキーマに作られます。2 種類あります:

静的なディメンションテーブル — 小さく、dbt run ごとに全件再ロードされます:

モデル行数備考
dim_date約 7,6702009–2029 年の日付スパイン。会計四半期と米国連邦祝日を含む
dim_payment_type6シードデータからのパススルー
dim_vendor3シードデータからのパススルー
dim_taxi_zones265stg_taxi_zones 経由のパススルー

インクリメンタルなファクト/集計テーブル — 大きく、実行ごとに MERGE で更新されます:

モデル行数備考
fact_trips50Mトリップごとに 1 行、完全非正規化
agg_hourly_zone_trips約 9Mzone ごとの時間別カウントを事前集計

マテリアライゼーション

マテリアライゼーションは、あるモデルに対して dbt が Snowflake 上に何を作るかを決めます。

マテリアライゼーション作られる Snowflake オブジェクト使いどころ
viewCREATE VIEW安価。常に最新データを反映。staging で使用
tableCREATE TABLE AS SELECT実行ごとに全件再構築。小さなディメンションで使用
incremental既存テーブルへの MERGE INTO大きなテーブル。新しい行のみを処理
ephemeral(オブジェクトなし — CTE としてインライン展開)1 つの下流モデルが共有する中間ロジック

2 つのインクリメンタルモデルは、それぞれ異なるインクリメンタル戦略を示しています:

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 ...;

これはこのラボで扱う最も重要な移行課題の 1 つです。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 に作られます。

このマクロは、1 つの 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 は Terraform が作成する最小権限のロールで、dbt に必要な権限のみを持ちます。
  • threads: 4 は dbt が並列にビルドするモデル数を制御します。
  • 認証情報は環境変数から取得され、セットアップ実行前に .env から読み込まれます。

テスト

dbt のテストには 2 つの形式があります:

スキーマテスト (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 はそれを実行し、行が 1 つでも返れば失敗とします。

すべてのテストを実行するには:

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 が空の場合 (初回実行時、またはテアダウン後)、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 スクリプトが面倒を見ます。

このページの内容

JA