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.sql3 つのレイヤー (Medallion アーキテクチャ)
レイヤー 1 — Staging (models/staging/)
目的: 到着したそのままの生データを受け取り、使える形にする。
これらのモデルは STAGING スキーマにビューとして作られます (ストレージコストなし — クエリ実行時に評価されます)。各 staging モデルは 1 つの仕事だけを担います:
| モデル | ソース | 何をするか |
|---|---|---|
stg_trips | RAW.TRIPS_RAW | カラム名を snake_case にリネームし、duration_minutes を追加し、VARIANT 型の TRIP_METADATA カラムを型付きカラムへフラット化 |
stg_taxi_zones | ANALYTICS.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 TABLEephemeral は、中間結果が 1 つの下流モデルからのみ必要で、ストレージやクエリのコンパイルオーバーヘッドを払いたくない場合に使います。
レイヤー 3 — Analytics (models/analytics/)
目的: ダッシュボードにそのまま使える最終テーブル。
これらは ANALYTICS スキーマに作られます。2 種類あります:
静的なディメンションテーブル — 小さく、dbt run ごとに全件再ロードされます:
| モデル | 行数 | 備考 |
|---|---|---|
dim_date | 約 7,670 | 2009–2029 年の日付スパイン。会計四半期と米国連邦祝日を含む |
dim_payment_type | 6 | シードデータからのパススルー |
dim_vendor | 3 | シードデータからのパススルー |
dim_taxi_zones | 265 | stg_taxi_zones 経由のパススルー |
インクリメンタルなファクト/集計テーブル — 大きく、実行ごとに MERGE で更新されます:
| モデル | 行数 | 備考 |
|---|---|---|
fact_trips | 50M | トリップごとに 1 行、完全非正規化 |
agg_hourly_zone_trips | 約 9M | zone ごとの時間別カウントを事前集計 |
マテリアライゼーション
マテリアライゼーションは、あるモデルに対して dbt が Snowflake 上に何を作るかを決めます。
| マテリアライゼーション | 作られる Snowflake オブジェクト | 使いどころ |
|---|---|---|
view | CREATE VIEW | 安価。常に最新データを反映。staging で使用 |
table | CREATE 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-refreshMERGE 戦略 (主要な移行課題)
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: 1000not_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 depsdbt_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 deps | packages.yml からパッケージをインストール |
dbt run | すべてのモデルをビルド (可能な場合はインクリメンタル) |
dbt run --full-refresh | すべてのインクリメンタルモデルをゼロから再構築 |
dbt run -s fact_trips | fact_trips とその依存のみをビルド |
dbt test | すべてのスキーマテストとカスタムテストを実行 |
dbt build | dbt 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 taskdbt はパイプラインの中間に位置します。生のテーブルが存在してデータが入るまでは実行できません。この順序は setup.sh スクリプトが面倒を見ます。